This is an automated email from the ASF dual-hosted git repository.
yuxiqian pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
The following commit(s) were added to refs/heads/master by this push:
new 55cccc909 [FLINK-40223][runtime] Support additional Flink SQL string
functions in YAML Transform (#4488)
55cccc909 is described below
commit 55cccc90922f215d9d7e21cbda3d1203c5597e20
Author: haruki <[email protected]>
AuthorDate: Fri Aug 7 17:36:08 2026 +0800
[FLINK-40223][runtime] Support additional Flink SQL string functions in
YAML Transform (#4488)
---
docs/content.zh/docs/core-concept/transform.md | 19 ++
docs/content/docs/core-concept/transform.md | 19 ++
.../src/test/resources/specs/string.yaml | 128 ++++++++
.../runtime/functions/impl/StringFunctions.java | 365 +++++++++++++++++++++
.../flink/cdc/runtime/parser/JaninoCompiler.java | 1 +
.../parser/metadata/TransformSqlOperatorTable.java | 186 +++++++++++
.../functions/impl/StringFunctionsTest.java | 105 ++++++
.../transform/PostTransformOperatorTest.java | 139 ++++++++
.../cdc/runtime/parser/TransformParserTest.java | 121 +++++++
9 files changed, 1083 insertions(+)
diff --git a/docs/content.zh/docs/core-concept/transform.md
b/docs/content.zh/docs/core-concept/transform.md
index a563c040b..726769254 100644
--- a/docs/content.zh/docs/core-concept/transform.md
+++ b/docs/content.zh/docs/core-concept/transform.md
@@ -177,6 +177,9 @@ Flink CDC 使用 [Calcite](https://calcite.apache.org/) 来解析表达式并且
| UPPER(string) | upper(string)
| 返回大写形式的字符串。
|
| LOWER(string) | lower(string)
| 返回小写形式的字符串。
|
| TRIM(string1) | trim('BOTH',string1)
| 返回去除两端空格的字符串。
|
+| LTRIM(string[, trimString]) | ltrim(string[,
trimString]) | 返回去除开头 trimString 字符后的字符串,默认去除空格。
|
+| RTRIM(string[, trimString]) | rtrim(string[,
trimString]) | 返回去除末尾 trimString 字符后的字符串,默认去除空格。
|
+| BTRIM(string[, trimString]) | btrim(string[,
trimString]) | 返回去除开头和末尾 trimString 字符后的字符串,默认去除空格。
|
| REGEXP_REPLACE(string1, string2, string3) | regexpReplace(string1,
string2, string3) | 返回将 STRING1 中所有匹配正则表达式 STRING2 的子串替换为 STRING3
后的字符串。例如,'foobar'.regexpReplace('oo\|ar', '') 返回 "fb"。 |
| REGEXP_EXTRACT(string, regex[, extractIndex]) | regexpExtract(string,
regex[, extractIndex]) | 返回正则表达式组 extractIndex 捕获的子串。extractIndex 默认为 0,0
表示完整匹配。输入为 NULL、无匹配、正则表达式非法或组索引非法时返回 NULL。 |
| REGEXP_EXTRACT_ALL(string, regex[, extractIndex]) | regexpExtractAll(string,
regex[, extractIndex]) | 返回组 extractIndex 捕获的所有子串组成的
ARRAY<STRING>。extractIndex 默认为 1,0 表示完整匹配。无匹配时返回空数组,输入为
NULL、正则表达式非法或组索引非法时返回 NULL。 |
@@ -185,7 +188,23 @@ Flink CDC 使用 [Calcite](https://calcite.apache.org/)
来解析表达式并且
| REGEXP_SUBSTR(string, regex) | regexpSubstr(string,
regex) | 返回第一个匹配正则表达式的子串。输入为 NULL、无匹配或正则表达式非法时返回 NULL。 |
| SUBSTR(string, integer1[, integer2]) |
substr(string,integer1,integer2) | 返回 STRING 从位置 integer1 开始、长度为
integer2(默认到末尾)的子串。
|
| SUBSTRING(string FROM integer1 [ FOR integer2 ]) |
substring(string,integer1,integer2) | 返回 STRING 从位置 integer1 开始、长度为
integer2(默认到末尾)的子串。
|
+| OVERLAY(string1 PLACING string2 FROM integer1 [FOR integer2]) |
overlay(string1, string2, integer1[, integer2]) | 从位置 integer1 开始,用 STRING2 替换
STRING1 的子串,替换长度默认为 STRING2 的长度,支持字符字符串和二进制字符串。
|
+| POSITION(string1 IN string2 [ FROM integer ]) |
position(string1, string2[, integer]) | 返回 STRING1 在 STRING2
中第一次出现的位置,可指定从 integer 开始查找。起始位置为 1,未找到时返回 0,支持字符字符串和二进制字符串。
|
+| LOCATE(string1, string2[, integer]) | locate(string1,
string2[, integer]) | 返回 STRING1 在 STRING2 中第一次出现的位置,可指定从 integer
开始查找。起始位置为 1,未找到时返回 0。 |
+| INSTR(string1, string2) | instr(string1,
string2) | 返回 STRING2 在 STRING1 中第一次出现的位置。起始位置为 1,未找到时返回
0。 |
| CONCAT(string1, string2,…) | concat(string1,
string2,…) | 返回连接 string1、string2、… 后的字符串。例如,CONCAT('AA', 'BB',
'CC') 返回 'AABBCC'。 |
+| CONCAT_WS(separator, string1, string2,...) |
concatWs(separator, string1, string2,...) | 使用分隔符连接 string1、string2、...
后返回字符串。NULL 字符串参数会被跳过。
|
+| LPAD(string1, integer, string2) | lpad(string1,
integer, string2) | 返回使用 STRING2 左填充 STRING1 至 integer
个字符后的字符串。如果 STRING1 更长,则截断到 integer 个字符。 |
+| RPAD(string1, integer, string2) | rpad(string1,
integer, string2) | 返回使用 STRING2 右填充 STRING1 至 integer
个字符后的字符串。如果 STRING1 更长,则截断到 integer 个字符。 |
+| REPLACE(string1, string2, string3) |
replace(string1, string2, string3) | 返回将 STRING1 中所有 STRING2 替换为
STRING3 后的字符串。
|
+| REPEAT(string, integer) | repeat(string,
integer) | 返回将 STRING 重复 integer 次后的字符串。
|
+| LEFT(string, integer) | left(string,
integer) | 返回 STRING 最左侧 integer 个字符。
|
+| RIGHT(string, integer) | right(string,
integer) | 返回 STRING 最右侧 integer 个字符。
|
+| STARTSWITH(string1, string2) |
startswith(string1, string2) | 返回 STRING1 是否以 STRING2
开头,支持字符字符串和二进制字符串。
|
+| ENDSWITH(string1, string2) |
endswith(string1, string2) | 返回 STRING1 是否以 STRING2
结尾,支持字符字符串和二进制字符串。
|
+| TO_BASE64(string \| binary) | toBase64(string
\| binary) | 将字符或二进制字符串编码为 base64 字符串。
|
+| FROM_BASE64(string) |
fromBase64(string) | 将 base64 字符串按 UTF-8
解码,并假定解码后的字节是有效的 UTF-8。对于任意二进制内容,请使用 FROM_BASE64_BINARY。 |
+| FROM_BASE64_BINARY(string) |
fromBase64Binary(string) | 将 base64 字符串解码为
VARBINARY,并保留原始字节。
|
## 时间函数
diff --git a/docs/content/docs/core-concept/transform.md
b/docs/content/docs/core-concept/transform.md
index e0a8451dd..f60f4d7c7 100644
--- a/docs/content/docs/core-concept/transform.md
+++ b/docs/content/docs/core-concept/transform.md
@@ -178,6 +178,9 @@ Logical functions follow SQL three-valued logic for
nullable BOOLEAN values. `AN
| UPPER(string) | upper(string)
| Returns string in uppercase.
|
| LOWER(string) | lower(string)
| Returns string in lowercase.
|
| TRIM(string1) | trim('BOTH',string1)
| Returns a string that removes whitespaces at both sides.
|
+| LTRIM(string[, trimString]) | ltrim(string[,
trimString]) | Returns a string with leading characters in
trimString removed. Whitespace is removed by default.
|
+| RTRIM(string[, trimString]) | rtrim(string[,
trimString]) | Returns a string with trailing characters
in trimString removed. Whitespace is removed by default.
|
+| BTRIM(string[, trimString]) | btrim(string[,
trimString]) | Returns a string with leading and trailing
characters in trimString removed. Whitespace is removed by default.
|
| REGEXP_REPLACE(string1, string2, string3) | regexpReplace(string1,
string2, string3) | Returns a string from STRING1 with all the substrings that
match a regular expression STRING2 consecutively being replaced with STRING3.
E.g., 'foobar'.regexpReplace('oo\|ar', '') returns "fb". |
| REGEXP_EXTRACT(string, regex[, extractIndex]) | regexpExtract(string,
regex[, extractIndex]) | Returns the substring captured by the regular
expression group extractIndex. extractIndex defaults to 0, where 0 means the
whole match. Returns NULL for NULL input, no match, invalid regex, or invalid
group index. |
| REGEXP_EXTRACT_ALL(string, regex[, extractIndex]) | regexpExtractAll(string,
regex[, extractIndex]) | Returns ARRAY<STRING> with all substrings
captured by group extractIndex. extractIndex defaults to 1, and 0 means the
whole match. Returns an empty array if there is no match, and NULL for NULL
input, invalid regex, or invalid group index. |
@@ -186,7 +189,23 @@ Logical functions follow SQL three-valued logic for
nullable BOOLEAN values. `AN
| REGEXP_SUBSTR(string, regex) | regexpSubstr(string,
regex) | Returns the first substring that matches regex. Returns
NULL for NULL input, no match, or invalid regex. |
| SUBSTR(string, integer1[, integer2]) |
substr(string,integer1,integer2) | Returns a substring of STRING
starting from position integer1 with length integer2 (to the end by default).
|
| SUBSTRING(string FROM integer1 [ FOR integer2 ]) |
substring(string,integer1,integer2) | Returns a substring of STRING
starting from position integer1 with length integer2 (to the end by default).
|
+| OVERLAY(string1 PLACING string2 FROM integer1 [FOR integer2]) |
overlay(string1, string2, integer1[, integer2]) | Replaces a substring of
STRING1 with STRING2 from position integer1. The replaced length defaults to
the length of STRING2. Character and binary strings are supported.
|
+| POSITION(string1 IN string2 [ FROM integer ]) |
position(string1, string2[, integer]) | Returns the position of the
first occurrence of STRING1 in STRING2, optionally starting from integer. The
first position is 1. Returns 0 if not found. Character and binary strings are
supported. |
+| LOCATE(string1, string2[, integer]) |
locate(string1, string2[, integer]) | Returns the position of the
first occurrence of STRING1 in STRING2, optionally starting from integer. The
first position is 1. Returns 0 if not found.
|
+| INSTR(string1, string2) | instr(string1,
string2) | Returns the position of the first
occurrence of STRING2 in STRING1. The first position is 1. Returns 0 if not
found.
|
| CONCAT(string1, string2,…) | concat(string1,
string2,…) | Returns a string that concatenates string1, string2,
…. E.g., CONCAT('AA', 'BB', 'CC') returns 'AABBCC'.
|
+| CONCAT_WS(separator, string1, string2,...) |
concatWs(separator, string1, string2,...) | Returns a string that
concatenates string1, string2, ... with a separator. Null string arguments are
skipped.
|
+| LPAD(string1, integer, string2) | lpad(string1,
integer, string2) | Returns STRING1 left-padded with STRING2
to a length of integer characters. If STRING1 is longer, it is shortened to
integer characters. |
+| RPAD(string1, integer, string2) | rpad(string1,
integer, string2) | Returns STRING1 right-padded with STRING2
to a length of integer characters. If STRING1 is longer, it is shortened to
integer characters. |
+| REPLACE(string1, string2, string3) |
replace(string1, string2, string3) | Returns STRING1 with all
occurrences of STRING2 replaced by STRING3.
|
+| REPEAT(string, integer) | repeat(string,
integer) | Returns a string that repeats STRING
integer times.
|
+| LEFT(string, integer) | left(string,
integer) | Returns the leftmost integer characters
from STRING.
|
+| RIGHT(string, integer) | right(string,
integer) | Returns the rightmost integer characters
from STRING.
|
+| STARTSWITH(string1, string2) |
startswith(string1, string2) | Returns whether STRING1
starts with STRING2. Character and binary strings are supported.
|
+| ENDSWITH(string1, string2) |
endswith(string1, string2) | Returns whether STRING1 ends
with STRING2. Character and binary strings are supported.
|
+| TO_BASE64(string \| binary) |
toBase64(string \| binary) | Encodes a character or
binary string to a base64 string.
|
+| FROM_BASE64(string) |
fromBase64(string) | Decodes a base64 string as
UTF-8 and assumes the decoded bytes are valid UTF-8. Use FROM_BASE64_BINARY for
arbitrary binary content.
|
+| FROM_BASE64_BINARY(string) |
fromBase64Binary(string) | Decodes a base64 string to
VARBINARY and preserves the original bytes.
|
## Temporal Functions
diff --git a/flink-cdc-composer/src/test/resources/specs/string.yaml
b/flink-cdc-composer/src/test/resources/specs/string.yaml
index 9069005f4..10c43c5fd 100644
--- a/flink-cdc-composer/src/test/resources/specs/string.yaml
+++ b/flink-cdc-composer/src/test/resources/specs/string.yaml
@@ -147,3 +147,131 @@
DataChangeEvent{tableId=foo.bar.baz, before=[-1, 天地玄黄宇宙洪荒, {{天地玄黄宇宙洪荒}}],
after=[], op=DELETE, meta=()}
DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, {{null}}],
op=INSERT, meta=()}
DataChangeEvent{tableId=foo.bar.baz, before=[0, null, {{null}}], after=[],
op=DELETE, meta=()}
+- do: Additional String Search Functions
+ projection: |-
+ id_
+ string_
+ POSITION('Z' IN string_) AS position_
+ POSITION('i' IN string_ FROM 15) AS position_from_
+ LOCATE('Z', string_) AS locate_
+ LOCATE('i', string_, 16) AS locate_from_
+ INSTR(string_, 'Z') AS instr_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`string_` STRING,`position_` INT,`position_from_`
INT,`locate_` INT,`locate_from_` INT,`instr_` INT}, primaryKeys=id_, options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, From A to Z is
Lie, 11, 17, 11, 17, 11], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, From A to Z is Lie, 11,
17, 11, 17, 11], after=[-1, 天地玄黄宇宙洪荒, 0, 0, 0, 0, 0], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, 天地玄黄宇宙洪荒, 0, 0, 0, 0, 0],
after=[], op=DELETE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, null,
null, null, null, null], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, null, null, null, null,
null, null], after=[], op=DELETE, meta=()}
+- do: Additional String Trim Functions
+ projection: |-
+ id_
+ LTRIM(' CDC') AS ltrim_default_
+ LTRIM('xxCDC', 'x') AS ltrim_chars_
+ RTRIM('CDC ') AS rtrim_default_
+ RTRIM('CDCyy', 'y') AS rtrim_chars_
+ BTRIM(' CDC ') AS btrim_default_
+ BTRIM('xyCDCxy', 'xy') AS btrim_chars_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`ltrim_default_` STRING,`ltrim_chars_`
STRING,`rtrim_default_` STRING,`rtrim_chars_` STRING,`btrim_default_`
STRING,`btrim_chars_` STRING}, primaryKeys=id_, options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, CDC, CDC, CDC,
CDC, CDC, CDC], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, CDC, CDC, CDC, CDC, CDC,
CDC], after=[-1, CDC, CDC, CDC, CDC, CDC, CDC], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, CDC, CDC, CDC, CDC, CDC,
CDC], after=[], op=DELETE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, CDC, CDC, CDC,
CDC, CDC, CDC], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, CDC, CDC, CDC, CDC, CDC,
CDC], after=[], op=DELETE, meta=()}
+- do: Additional String Join Functions
+ projection: |-
+ id_
+ string_
+ CONCAT_WS('|', string_, NULL, 'CDC') AS joined_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`string_` STRING,`joined_` STRING}, primaryKeys=id_,
options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, From A to Z is
Lie, From A to Z is Lie|CDC], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, From A to Z is Lie, From A
to Z is Lie|CDC], after=[-1, 天地玄黄宇宙洪荒, 天地玄黄宇宙洪荒|CDC], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, 天地玄黄宇宙洪荒, 天地玄黄宇宙洪荒|CDC],
after=[], op=DELETE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, CDC],
op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, null, CDC], after=[],
op=DELETE, meta=()}
+- do: Additional String Pad Replace Repeat Functions
+ projection: |-
+ id_
+ string_
+ LPAD('CDC', 5, '0') AS lpad_
+ LPAD('CDCLONG', 3, '0') AS lpad_truncated_
+ RPAD('CDC', 5, '0') AS rpad_
+ RPAD('CDCLONG', 3, '0') AS rpad_truncated_
+ REPLACE(string_, ' ', '_') AS replaced_
+ REPEAT('CDC', 2) AS repeated_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`string_` STRING,`lpad_` STRING,`lpad_truncated_`
STRING,`rpad_` STRING,`rpad_truncated_` STRING,`replaced_` STRING,`repeated_`
STRING}, primaryKeys=id_, options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, From A to Z is
Lie, 00CDC, CDC, CDC00, CDC, From_A_to_Z_is_Lie, CDCCDC], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, From A to Z is Lie, 00CDC,
CDC, CDC00, CDC, From_A_to_Z_is_Lie, CDCCDC], after=[-1, 天地玄黄宇宙洪荒, 00CDC, CDC,
CDC00, CDC, 天地玄黄宇宙洪荒, CDCCDC], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, 天地玄黄宇宙洪荒, 00CDC, CDC,
CDC00, CDC, 天地玄黄宇宙洪荒, CDCCDC], after=[], op=DELETE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, 00CDC,
CDC, CDC00, CDC, null, CDCCDC], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, null, 00CDC, CDC, CDC00,
CDC, null, CDCCDC], after=[], op=DELETE, meta=()}
+- do: Additional String Slice Functions
+ projection: |-
+ id_
+ string_
+ OVERLAY(string_ PLACING 'CDC' FROM 1) AS overlay_default_
+ OVERLAY(string_ PLACING 'CDC' FROM 1 FOR 4) AS overlay_with_length_
+ LEFT(string_, 2) AS left_
+ RIGHT(string_, 2) AS right_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`string_` STRING,`overlay_default_`
STRING,`overlay_with_length_` STRING,`left_` STRING,`right_` STRING},
primaryKeys=id_, options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, From A to Z is
Lie, CDCm A to Z is Lie, CDC A to Z is Lie, Fr, ie], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, From A to Z is Lie, CDCm A
to Z is Lie, CDC A to Z is Lie, Fr, ie], after=[-1, 天地玄黄宇宙洪荒, CDC黄宇宙洪荒,
CDC宇宙洪荒, 天地, 洪荒], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, 天地玄黄宇宙洪荒, CDC黄宇宙洪荒,
CDC宇宙洪荒, 天地, 洪荒], after=[], op=DELETE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, null,
null, null, null], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, null, null, null, null,
null], after=[], op=DELETE, meta=()}
+- do: Additional String Predicate Functions
+ projection: |-
+ id_
+ string_
+ STARTSWITH(string_, 'From') AS starts_
+ ENDSWITH(string_, 'Lie') AS ends_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`string_` STRING,`starts_` BOOLEAN,`ends_` BOOLEAN},
primaryKeys=id_, options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, From A to Z is
Lie, true, true], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, From A to Z is Lie, true,
true], after=[-1, 天地玄黄宇宙洪荒, false, false], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, 天地玄黄宇宙洪荒, false, false],
after=[], op=DELETE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, null,
null], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, null, null, null],
after=[], op=DELETE, meta=()}
+- do: Additional Base64 Functions
+ projection: |-
+ id_
+ string_
+ TO_BASE64(string_) AS encoded_string_
+ FROM_BASE64(TO_BASE64(string_)) AS decoded_string_
+ TO_BASE64(varbinary_) AS encoded_binary_
+ FROM_BASE64_BINARY(TO_BASE64(varbinary_)) AS decoded_binary_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`string_` STRING,`encoded_string_` STRING,`decoded_string_`
STRING,`encoded_binary_` STRING,`decoded_binary_` VARBINARY(65536)},
primaryKeys=id_, options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, From A to Z is
Lie, RnJvbSBBIHRvIFogaXMgTGll, From A to Z is Lie, ZG9sb3Igc2l0IGFtZXQ=,
ZG9sb3Igc2l0IGFtZXQ=], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, From A to Z is Lie,
RnJvbSBBIHRvIFogaXMgTGll, From A to Z is Lie, ZG9sb3Igc2l0IGFtZXQ=,
ZG9sb3Igc2l0IGFtZXQ=], after=[-1, 天地玄黄宇宙洪荒, 5aSp5Zyw546E6buE5a6H5a6Z5rSq6I2S,
天地玄黄宇宙洪荒, 5YWt5LiD5YWr5Lmd5Y2B, 5YWt5LiD5YWr5Lmd5Y2B], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, 天地玄黄宇宙洪荒,
5aSp5Zyw546E6buE5a6H5a6Z5rSq6I2S, 天地玄黄宇宙洪荒, 5YWt5LiD5YWr5Lmd5Y2B,
5YWt5LiD5YWr5Lmd5Y2B], after=[], op=DELETE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, null,
null, null, null], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, null, null, null, null,
null], after=[], op=DELETE, meta=()}
+- do: Additional Binary String Functions
+ projection: |-
+ id_
+ TO_BASE64(OVERLAY(binary_ PLACING varbinary_ FROM 2)) AS
overlay_binary_default_
+ TO_BASE64(OVERLAY(binary_ PLACING varbinary_ FROM 2 FOR 3)) AS
overlay_binary_with_length_
+ POSITION(binary_ IN binary_) AS position_binary_
+ POSITION(binary_ IN binary_ FROM 2) AS position_binary_from_
+ STARTSWITH(binary_, binary_) AS starts_binary_
+ ENDSWITH(binary_, binary_) AS ends_binary_
+ primary-key: id_
+ expect: |-
+ CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT
NULL 'Identifier',`overlay_binary_default_`
STRING,`overlay_binary_with_length_` STRING,`position_binary_`
INT,`position_binary_from_` INT,`starts_binary_` BOOLEAN,`ends_binary_`
BOOLEAN}, primaryKeys=id_, options=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1,
TGRvbG9yIHNpdCBhbWV0, TGRvbG9yIHNpdCBhbWV0bSBpcHN1bQ==, 1, 0, true, true],
op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[1, TGRvbG9yIHNpdCBhbWV0,
TGRvbG9yIHNpdCBhbWV0bSBpcHN1bQ==, 1, 0, true, true], after=[-1,
5OWFreS4g+WFq+S5neWNgQ==, 5OWFreS4g+WFq+S5neWNgbqM5LiJ5Zub5LqU, 1, 0, true,
true], op=UPDATE, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[-1, 5OWFreS4g+WFq+S5neWNgQ==,
5OWFreS4g+WFq+S5neWNgbqM5LiJ5Zub5LqU, 1, 0, true, true], after=[], op=DELETE,
meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, null,
null, null, null, null], op=INSERT, meta=()}
+ DataChangeEvent{tableId=foo.bar.baz, before=[0, null, null, null, null,
null, null], after=[], op=DELETE, meta=()}
diff --git
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctions.java
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctions.java
index fca030674..c665b44df 100644
---
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctions.java
+++
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctions.java
@@ -25,7 +25,9 @@ import org.apache.calcite.runtime.SqlFunctions;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
+import java.util.Base64;
import java.util.List;
import java.util.UUID;
import java.util.regex.Matcher;
@@ -138,6 +140,25 @@ public class StringFunctions {
return String.join("", str);
}
+ public static String concatWs(String separator, String... str) {
+ if (separator == null || str == null) {
+ return null;
+ }
+ StringBuilder builder = new StringBuilder();
+ boolean first = true;
+ for (String s : str) {
+ if (s == null) {
+ continue;
+ }
+ if (!first) {
+ builder.append(separator);
+ }
+ builder.append(s);
+ first = false;
+ }
+ return builder.toString();
+ }
+
public static boolean like(String str, String regex) {
return Pattern.compile(regex).matcher(str).find();
}
@@ -248,6 +269,256 @@ public class StringFunctions {
return str.toLowerCase();
}
+ public static String overlay(String str, String replacement, Number start)
{
+ if (replacement == null) {
+ return null;
+ }
+ return overlay(str, replacement, start, replacement.length());
+ }
+
+ public static String overlay(String str, String replacement, Number start,
Number length) {
+ if (str == null || replacement == null || start == null || length ==
null) {
+ return null;
+ }
+ int startPosition = start.intValue();
+ int len = length.intValue();
+ if (startPosition <= 0 || startPosition > str.length()) {
+ return str;
+ }
+
+ StringBuilder builder = new StringBuilder();
+ builder.append(str, 0, startPosition - 1);
+ builder.append(replacement);
+ if (startPosition + len <= str.length() && len > 0) {
+ builder.append(str.substring(startPosition - 1 + len));
+ }
+ return builder.toString();
+ }
+
+ public static byte[] overlay(byte[] bytes, byte[] replacement, Number
start) {
+ if (replacement == null) {
+ return null;
+ }
+ return overlay(bytes, replacement, start, replacement.length);
+ }
+
+ public static byte[] overlay(byte[] bytes, byte[] replacement, Number
start, Number length) {
+ if (bytes == null || replacement == null || start == null || length ==
null) {
+ return null;
+ }
+ int startPosition = start.intValue();
+ int len = length.intValue();
+ if (startPosition <= 0 || startPosition > bytes.length) {
+ return bytes;
+ }
+
+ int prefixLength = startPosition - 1;
+ int suffixStart = prefixLength + Math.max(len, 0);
+ int suffixLength = len > 0 && suffixStart < bytes.length ?
bytes.length - suffixStart : 0;
+ byte[] result = new byte[prefixLength + replacement.length +
suffixLength];
+ System.arraycopy(bytes, 0, result, 0, prefixLength);
+ System.arraycopy(replacement, 0, result, prefixLength,
replacement.length);
+ if (suffixLength > 0) {
+ System.arraycopy(
+ bytes, suffixStart, result, prefixLength +
replacement.length, suffixLength);
+ }
+ return result;
+ }
+
+ public static Integer position(String seek, String str) {
+ return position(seek, str, 1);
+ }
+
+ public static Integer position(String seek, String str, Number from) {
+ if (seek == null || str == null || from == null) {
+ return null;
+ }
+ if (seek.isEmpty()) {
+ return 1;
+ }
+ int fromCodePoint = Math.max(from.intValue() - 1, 0);
+ int codePointLength = str.codePointCount(0, str.length());
+ if (fromCodePoint > codePointLength) {
+ return 0;
+ }
+ int fromIndex = str.offsetByCodePoints(0, fromCodePoint);
+ int index = str.indexOf(seek, fromIndex);
+ return index < 0 ? 0 : str.codePointCount(0, index) + 1;
+ }
+
+ public static Integer position(byte[] seek, byte[] bytes) {
+ return position(seek, bytes, 1);
+ }
+
+ public static Integer position(byte[] seek, byte[] bytes, Number from) {
+ if (seek == null || bytes == null || from == null) {
+ return null;
+ }
+ if (seek.length == 0) {
+ return 1;
+ }
+ int fromIndex = Math.max(from.intValue() - 1, 0);
+ if (fromIndex > bytes.length) {
+ return 0;
+ }
+ int index = indexOf(bytes, seek, fromIndex);
+ return index < 0 ? 0 : index + 1;
+ }
+
+ public static Integer instr(String str, String subString) {
+ if (str == null || subString == null) {
+ return null;
+ }
+ int index = str.indexOf(subString);
+ return index < 0 ? 0 : str.codePointCount(0, index) + 1;
+ }
+
+ public static Integer locate(String seek, String str) {
+ return position(seek, str);
+ }
+
+ public static Integer locate(String seek, String str, Number from) {
+ return position(seek, str, from);
+ }
+
+ public static String ltrim(String str) {
+ return ltrim(str, " ");
+ }
+
+ public static String ltrim(String str, String trimStr) {
+ return trim(str, trimStr, true, false);
+ }
+
+ public static String rtrim(String str) {
+ return rtrim(str, " ");
+ }
+
+ public static String rtrim(String str, String trimStr) {
+ return trim(str, trimStr, false, true);
+ }
+
+ public static String btrim(String str) {
+ return btrim(str, " ");
+ }
+
+ public static String btrim(String str, String trimStr) {
+ return trim(str, trimStr, true, true);
+ }
+
+ public static String lpad(String base, Number len, String pad) {
+ return pad(base, len, pad, true);
+ }
+
+ public static String rpad(String base, Number len, String pad) {
+ return pad(base, len, pad, false);
+ }
+
+ public static String replace(String str, String oldStr, String
replacement) {
+ if (str == null || oldStr == null || replacement == null) {
+ return null;
+ }
+ return str.replace(oldStr, replacement);
+ }
+
+ public static String repeat(String str, Number repeat) {
+ if (str == null || repeat == null) {
+ return null;
+ }
+ int count = repeat.intValue();
+ if (count <= 0) {
+ return "";
+ }
+ return str.repeat(count);
+ }
+
+ public static String left(String str, Number length) {
+ if (str == null || length == null) {
+ return null;
+ }
+ int len = length.intValue();
+ if (len <= 0) {
+ return "";
+ }
+ int codePointLength = str.codePointCount(0, str.length());
+ if (len >= codePointLength) {
+ return str;
+ }
+ return str.substring(0, str.offsetByCodePoints(0, len));
+ }
+
+ public static String right(String str, Number length) {
+ if (str == null || length == null) {
+ return null;
+ }
+ int len = length.intValue();
+ if (len <= 0) {
+ return "";
+ }
+ int codePointLength = str.codePointCount(0, str.length());
+ if (len >= codePointLength) {
+ return str;
+ }
+ return str.substring(str.offsetByCodePoints(0, codePointLength - len));
+ }
+
+ public static Boolean startswith(String str, String prefix) {
+ if (str == null || prefix == null) {
+ return null;
+ }
+ return str.startsWith(prefix);
+ }
+
+ public static Boolean startswith(byte[] bytes, byte[] prefix) {
+ if (bytes == null || prefix == null) {
+ return null;
+ }
+ return matchesAt(bytes, prefix, 0);
+ }
+
+ public static Boolean endswith(String str, String suffix) {
+ if (str == null || suffix == null) {
+ return null;
+ }
+ return str.endsWith(suffix);
+ }
+
+ public static Boolean endswith(byte[] bytes, byte[] suffix) {
+ if (bytes == null || suffix == null) {
+ return null;
+ }
+ return matchesAt(bytes, suffix, bytes.length - suffix.length);
+ }
+
+ public static String toBase64(String str) {
+ if (str == null) {
+ return null;
+ }
+ return
Base64.getEncoder().encodeToString(str.getBytes(StandardCharsets.UTF_8));
+ }
+
+ public static String toBase64(byte[] bytes) {
+ if (bytes == null) {
+ return null;
+ }
+ return Base64.getEncoder().encodeToString(bytes);
+ }
+
+ public static String fromBase64(String str) {
+ if (str == null) {
+ return null;
+ }
+ return new String(
+
Base64.getDecoder().decode(str.getBytes(StandardCharsets.UTF_8)),
+ StandardCharsets.UTF_8);
+ }
+
+ public static byte[] fromBase64Binary(String str) {
+ if (str == null) {
+ return null;
+ }
+ return
Base64.getDecoder().decode(str.getBytes(StandardCharsets.UTF_8));
+ }
+
public static String uuid() {
return UUID.randomUUID().toString();
}
@@ -299,4 +570,98 @@ public class StringFunctions {
return null;
}
}
+
+ private static String trim(
+ String str, String trimStr, boolean trimLeading, boolean
trimTrailing) {
+ if (str == null || trimStr == null) {
+ return null;
+ }
+ if (trimStr.isEmpty()) {
+ return str;
+ }
+ int begin = 0;
+ int end = str.length();
+ while (trimLeading && begin < end) {
+ int codePoint = str.codePointAt(begin);
+ if (trimStr.indexOf(codePoint) < 0) {
+ break;
+ }
+ begin += Character.charCount(codePoint);
+ }
+ while (trimTrailing && begin < end) {
+ int codePoint = str.codePointBefore(end);
+ if (trimStr.indexOf(codePoint) < 0) {
+ break;
+ }
+ end -= Character.charCount(codePoint);
+ }
+ return str.substring(begin, end);
+ }
+
+ private static boolean matchesAt(byte[] bytes, byte[] target, int offset) {
+ if (offset < 0 || offset + target.length > bytes.length) {
+ return false;
+ }
+ for (int i = 0; i < target.length; i++) {
+ if (bytes[offset + i] != target[i]) {
+ return false;
+ }
+ }
+ return true;
+ }
+
+ private static int indexOf(byte[] bytes, byte[] target, int fromIndex) {
+ int lastIndex = bytes.length - target.length;
+ for (int i = fromIndex; i <= lastIndex; i++) {
+ if (matchesAt(bytes, target, i)) {
+ return i;
+ }
+ }
+ return -1;
+ }
+
+ private static String pad(String base, Number length, String pad, boolean
leftPad) {
+ if (base == null || length == null || pad == null) {
+ return null;
+ }
+ int len = length.intValue();
+ if (len < 0 || pad.isEmpty()) {
+ return null;
+ } else if (len == 0) {
+ return "";
+ }
+
+ char[] data = new char[len];
+ char[] baseChars = base.toCharArray();
+ char[] padChars = pad.toCharArray();
+
+ if (leftPad) {
+ int pos = Math.max(len - base.length(), 0);
+ for (int i = 0; i < pos; i += pad.length()) {
+ for (int j = 0; j < pad.length() && j < pos - i; j++) {
+ data[i + j] = padChars[j];
+ }
+ }
+ int i = 0;
+ while (pos + i < len && i < base.length()) {
+ data[pos + i] = baseChars[i];
+ i++;
+ }
+ } else {
+ int pos = 0;
+ while (pos < base.length() && pos < len) {
+ data[pos] = baseChars[pos];
+ pos++;
+ }
+ while (pos < len) {
+ int i = 0;
+ while (i < pad.length() && i < len - pos) {
+ data[pos + i] = padChars[i];
+ i++;
+ }
+ pos += pad.length();
+ }
+ }
+ return new String(data);
+ }
}
diff --git
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/JaninoCompiler.java
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/JaninoCompiler.java
index 539a2cd55..53e81a5dc 100644
---
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/JaninoCompiler.java
+++
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/JaninoCompiler.java
@@ -301,6 +301,7 @@ public class JaninoCompiler {
case NOT_IN:
case LIKE:
case SIMILAR:
+ case POSITION:
case CEIL:
case FLOOR:
case TRIM:
diff --git
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSqlOperatorTable.java
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSqlOperatorTable.java
index a278799af..4850815cc 100644
---
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSqlOperatorTable.java
+++
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSqlOperatorTable.java
@@ -20,6 +20,7 @@ package org.apache.flink.cdc.runtime.parser.metadata;
import org.apache.flink.cdc.runtime.functions.BuiltInScalarFunction;
import org.apache.flink.cdc.runtime.functions.BuiltInTimestampFunction;
+import org.apache.calcite.rel.type.RelDataTypeSystem;
import org.apache.calcite.sql.SqlBinaryOperator;
import org.apache.calcite.sql.SqlFunction;
import org.apache.calcite.sql.SqlFunctionCategory;
@@ -174,6 +175,38 @@ public class TransformSqlOperatorTable extends
ReflectiveSqlOperatorTable {
public static final SqlFunction UPPER = SqlStdOperatorTable.UPPER;
public static final SqlFunction LOWER = SqlStdOperatorTable.LOWER;
public static final SqlFunction TRIM = SqlStdOperatorTable.TRIM;
+ public static final SqlFunction LTRIM =
+ new SqlFunction(
+ "LTRIM",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.ARG0_NULLABLE_VARYING,
+ null,
+ OperandTypes.or(
+ OperandTypes.family(SqlTypeFamily.CHARACTER),
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.CHARACTER)),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction RTRIM =
+ new SqlFunction(
+ "RTRIM",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.ARG0_NULLABLE_VARYING,
+ null,
+ OperandTypes.or(
+ OperandTypes.family(SqlTypeFamily.CHARACTER),
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.CHARACTER)),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction BTRIM =
+ new SqlFunction(
+ "BTRIM",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.or(
+ OperandTypes.family(SqlTypeFamily.CHARACTER),
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.CHARACTER)),
+ SqlFunctionCategory.STRING);
public static final SqlFunction REGEXP_REPLACE =
new SqlFunction(
"REGEXP_REPLACE",
@@ -253,6 +286,159 @@ public class TransformSqlOperatorTable extends
ReflectiveSqlOperatorTable {
SqlTypeFamily.INTEGER)),
SqlFunctionCategory.STRING);
public static final SqlFunction SUBSTRING = SqlStdOperatorTable.SUBSTRING;
+ public static final SqlFunction OVERLAY = SqlStdOperatorTable.OVERLAY;
+ public static final SqlFunction POSITION = SqlStdOperatorTable.POSITION;
+ public static final SqlFunction LOCATE =
+ new SqlFunction(
+ "LOCATE",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.INTEGER_NULLABLE,
+ null,
+ OperandTypes.or(
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.CHARACTER),
+ OperandTypes.family(
+ SqlTypeFamily.CHARACTER,
+ SqlTypeFamily.CHARACTER,
+ SqlTypeFamily.INTEGER)),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction INSTR =
+ new SqlFunction(
+ "INSTR",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.INTEGER_NULLABLE,
+ null,
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.CHARACTER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction CONCAT_WS =
+ new SqlFunction(
+ "CONCAT_WS",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.repeat(SqlOperandCountRanges.from(2),
OperandTypes.CHARACTER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction LPAD =
+ new SqlFunction(
+ "LPAD",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.family(
+ SqlTypeFamily.CHARACTER,
+ SqlTypeFamily.INTEGER,
+ SqlTypeFamily.CHARACTER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction RPAD =
+ new SqlFunction(
+ "RPAD",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.family(
+ SqlTypeFamily.CHARACTER,
+ SqlTypeFamily.INTEGER,
+ SqlTypeFamily.CHARACTER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction REPLACE =
+ new SqlFunction(
+ "REPLACE",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.family(
+ SqlTypeFamily.CHARACTER,
+ SqlTypeFamily.CHARACTER,
+ SqlTypeFamily.CHARACTER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction REPEAT =
+ new SqlFunction(
+ "REPEAT",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.INTEGER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction LEFT =
+ new SqlFunction(
+ "LEFT",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.ARG0_NULLABLE_VARYING,
+ null,
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.INTEGER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction RIGHT =
+ new SqlFunction(
+ "RIGHT",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.ARG0_NULLABLE_VARYING,
+ null,
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.INTEGER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction STARTSWITH =
+ new SqlFunction(
+ "STARTSWITH",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.BOOLEAN_NULLABLE,
+ null,
+ OperandTypes.or(
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.CHARACTER),
+ OperandTypes.family(SqlTypeFamily.BINARY,
SqlTypeFamily.BINARY)),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction ENDSWITH =
+ new SqlFunction(
+ "ENDSWITH",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.BOOLEAN_NULLABLE,
+ null,
+ OperandTypes.or(
+ OperandTypes.family(SqlTypeFamily.CHARACTER,
SqlTypeFamily.CHARACTER),
+ OperandTypes.family(SqlTypeFamily.BINARY,
SqlTypeFamily.BINARY)),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction TO_BASE64 =
+ new SqlFunction(
+ "TO_BASE64",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.or(
+ OperandTypes.family(SqlTypeFamily.CHARACTER),
+ OperandTypes.family(SqlTypeFamily.BINARY)),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction FROM_BASE64 =
+ new SqlFunction(
+ "FROM_BASE64",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(SqlTypeName.VARCHAR),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.family(SqlTypeFamily.CHARACTER),
+ SqlFunctionCategory.STRING);
+ public static final SqlFunction FROM_BASE64_BINARY =
+ new SqlFunction(
+ "FROM_BASE64_BINARY",
+ SqlKind.OTHER_FUNCTION,
+ ReturnTypes.cascade(
+ ReturnTypes.explicit(
+ SqlTypeName.VARBINARY,
+ RelDataTypeSystem.DEFAULT.getMaxPrecision(
+ SqlTypeName.VARBINARY)),
+ SqlTypeTransforms.TO_NULLABLE),
+ null,
+ OperandTypes.family(SqlTypeFamily.CHARACTER),
+ SqlFunctionCategory.STRING);
// ------------------
// Temporal Functions
diff --git
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctionsTest.java
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctionsTest.java
index 4e55db85b..4994dc0a1 100644
---
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctionsTest.java
+++
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/StringFunctionsTest.java
@@ -24,6 +24,111 @@ import static org.assertj.core.api.Assertions.assertThat;
/** Unit tests for {@link StringFunctions}. */
class StringFunctionsTest {
+ @Test
+ void testStringFunctions() {
+ assertThat(StringFunctions.overlay("abcdef", "ZZ",
3)).isEqualTo("abZZef");
+ assertThat(StringFunctions.overlay("abcdef", "ZZ", 3,
2)).isEqualTo("abZZef");
+ assertThat(StringFunctions.overlay("abcdef", "ZZ", 0,
2)).isEqualTo("abcdef");
+ assertThat(StringFunctions.overlay(new byte[] {1, 2, 3, 4}, new byte[]
{8, 9}, 2))
+ .containsExactly((byte) 1, (byte) 8, (byte) 9, (byte) 4);
+ assertThat(StringFunctions.overlay(new byte[] {1, 2, 3, 4}, new byte[]
{8, 9}, 2, 2))
+ .containsExactly((byte) 1, (byte) 8, (byte) 9, (byte) 4);
+ assertThat(StringFunctions.position("cd", "abcdef")).isEqualTo(3);
+ assertThat(StringFunctions.position("xy", "abcdef")).isZero();
+ assertThat(StringFunctions.position("", "abcdef")).isEqualTo(1);
+ assertThat(StringFunctions.position("cd", "abcdef", 4)).isZero();
+ assertThat(StringFunctions.position(new byte[] {2, 3}, new byte[] {1,
2, 3, 4}))
+ .isEqualTo(2);
+ assertThat(StringFunctions.position(new byte[] {2, 3}, new byte[] {1,
2, 3, 4}, 3))
+ .isZero();
+ assertThat(StringFunctions.locate("cd", "abcdef")).isEqualTo(3);
+ assertThat(StringFunctions.locate("", "abcdef", 4)).isEqualTo(1);
+ assertThat(StringFunctions.instr("abcabc", "bc")).isEqualTo(2);
+ assertThat(StringFunctions.ltrim(" abc ")).isEqualTo("abc ");
+ assertThat(StringFunctions.rtrim(" abc ")).isEqualTo(" abc");
+ assertThat(StringFunctions.btrim(" abc ")).isEqualTo("abc");
+ assertThat(StringFunctions.btrim("xyabcxy", "xy")).isEqualTo("abc");
+ assertThat(StringFunctions.concatWs(",", "a", null,
"b")).isEqualTo("a,b");
+ assertThat(StringFunctions.concatWs(",", null, null)).isEmpty();
+ assertThat(StringFunctions.concatWs("~", "AA", null, "BB", "", "CC"))
+ .isEqualTo("AA~BB~~CC");
+ assertThat(StringFunctions.lpad("hi", 5, "?")).isEqualTo("???hi");
+ assertThat(StringFunctions.rpad("hi", 5, "?")).isEqualTo("hi???");
+ assertThat(StringFunctions.lpad("hello", 2, "?")).isEqualTo("he");
+ assertThat(StringFunctions.lpad("hi", -1, "?")).isNull();
+ assertThat(StringFunctions.lpad("hi", 5, "")).isNull();
+ assertThat(StringFunctions.replace("hello", "l",
"x")).isEqualTo("hexxo");
+ assertThat(StringFunctions.repeat("ab", 3)).isEqualTo("ababab");
+ assertThat(StringFunctions.repeat("ab", 0)).isEmpty();
+ assertThat(StringFunctions.left("abcdef", 2)).isEqualTo("ab");
+ assertThat(StringFunctions.right("abcdef", 2)).isEqualTo("ef");
+ assertThat(StringFunctions.left("abcdef", 0)).isEmpty();
+ assertThat(StringFunctions.right("abcdef", -1)).isEmpty();
+ assertThat(StringFunctions.startswith("abcdef", "abc")).isTrue();
+ assertThat(StringFunctions.endswith("abcdef", "def")).isTrue();
+ assertThat(StringFunctions.startswith(new byte[] {1, 2, 3}, new byte[]
{1, 2})).isTrue();
+ assertThat(StringFunctions.startswith(new byte[] {1, 2, 3}, new byte[]
{2, 3})).isFalse();
+ assertThat(StringFunctions.endswith(new byte[] {1, 2, 3}, new byte[]
{2, 3})).isTrue();
+ assertThat(StringFunctions.endswith(new byte[] {1, 2, 3}, new byte[]
{1, 2})).isFalse();
+ assertThat(StringFunctions.toBase64("hello")).isEqualTo("aGVsbG8=");
+ assertThat(StringFunctions.toBase64(new byte[] {(byte) 0xC3,
0x28})).isEqualTo("wyg=");
+ assertThat(StringFunctions.fromBase64("aGVsbG8=")).isEqualTo("hello");
+ }
+
+ @Test
+ void testFromBase64ReplacesMalformedUtf8() {
+ assertThat(StringFunctions.fromBase64("wyg=")).isEqualTo("\uFFFD(");
+
assertThat(StringFunctions.toBase64(StringFunctions.fromBase64("wyg=")))
+ .isEqualTo("77+9KA==");
+
assertThat(StringFunctions.fromBase64Binary("wyg=")).containsExactly((byte)
0xC3, 0x28);
+ }
+
+ @Test
+ void testStringFunctionsReturnNullOnNullInput() {
+ assertThat(StringFunctions.concatWs(null, "a")).isNull();
+ assertThat(StringFunctions.overlay(null, "x", 1)).isNull();
+ assertThat(StringFunctions.overlay((byte[]) null, new byte[] {1},
1)).isNull();
+ assertThat(StringFunctions.position(null, "abc")).isNull();
+ assertThat(StringFunctions.position((byte[]) null, new byte[]
{1})).isNull();
+ assertThat(StringFunctions.locate("a", null)).isNull();
+ assertThat(StringFunctions.instr("abc", null)).isNull();
+ assertThat(StringFunctions.ltrim(null)).isNull();
+ assertThat(StringFunctions.rtrim(null)).isNull();
+ assertThat(StringFunctions.btrim(null)).isNull();
+ assertThat(StringFunctions.lpad(null, 1, "?")).isNull();
+ assertThat(StringFunctions.rpad("a", null, "?")).isNull();
+ assertThat(StringFunctions.replace("a", null, "b")).isNull();
+ assertThat(StringFunctions.repeat(null, 1)).isNull();
+ assertThat(StringFunctions.left(null, 1)).isNull();
+ assertThat(StringFunctions.right("a", null)).isNull();
+ assertThat(StringFunctions.startswith(null, "a")).isNull();
+ assertThat(StringFunctions.endswith("a", null)).isNull();
+ assertThat(StringFunctions.startswith((byte[]) null, new byte[]
{1})).isNull();
+ assertThat(StringFunctions.endswith(new byte[] {1}, (byte[])
null)).isNull();
+ assertThat(StringFunctions.toBase64((String) null)).isNull();
+ assertThat(StringFunctions.toBase64((byte[]) null)).isNull();
+ assertThat(StringFunctions.fromBase64(null)).isNull();
+ assertThat(StringFunctions.fromBase64Binary(null)).isNull();
+ }
+
+ @Test
+ void testStringFunctionsCountUnicodeCodePoints() {
+ String emoji = "\uD83D\uDE00";
+ String sameHighSurrogate = "\uD83D\uDE01";
+ String sameLowSurrogate = "\uD87D\uDE00";
+
+ assertThat(StringFunctions.position("x", emoji + "x")).isEqualTo(2);
+ assertThat(StringFunctions.locate("x", emoji + "x")).isEqualTo(2);
+ assertThat(StringFunctions.instr(emoji + "x", "x")).isEqualTo(2);
+ assertThat(StringFunctions.left(emoji + "x", 1)).isEqualTo(emoji);
+ assertThat(StringFunctions.right("x" + emoji, 1)).isEqualTo(emoji);
+ assertThat(StringFunctions.ltrim(sameHighSurrogate + "x", emoji))
+ .isEqualTo(sameHighSurrogate + "x");
+ assertThat(StringFunctions.rtrim("x" + sameLowSurrogate, emoji))
+ .isEqualTo("x" + sameLowSurrogate);
+ assertThat(StringFunctions.btrim(emoji + "x" + emoji,
emoji)).isEqualTo("x");
+ }
+
@Test
void testLegacyLikeUsesJavaRegex() {
assertThat(StringFunctions.like("Alice", "A.*")).isTrue();
diff --git
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
index fa2277a99..ba837efc0 100644
---
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
+++
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
@@ -2912,6 +2912,36 @@ class PostTransformOperatorTest {
testExpressionConditionTransform("concat('123', 'abc') = '123abc'");
testExpressionConditionTransform("upper('abc') = 'ABC'");
testExpressionConditionTransform("lower('ABC') = 'abc'");
+ testExpressionConditionTransform("OVERLAY('abcdef' PLACING 'ZZ' FROM
3) = 'abZZef'");
+ testExpressionConditionTransform("OVERLAY('abcdef' PLACING 'ZZ' FROM 3
FOR 2) = 'abZZef'");
+ testExpressionConditionTransform("POSITION('cd' IN 'abcdef') = 3");
+ testExpressionConditionTransform("POSITION('x' IN '\uD83D\uDE00x') =
2");
+ testExpressionConditionTransform("LOCATE('cd', 'abcdef') = 3");
+ testExpressionConditionTransform("LOCATE('cd', 'abcdef', 4) = 0");
+ testExpressionConditionTransform("LOCATE('', 'abcdef', 4) = 1");
+ testExpressionConditionTransform("INSTR('abcabc', 'bc') = 2");
+ testExpressionConditionTransform("LTRIM(' abc ') = 'abc '");
+ testExpressionConditionTransform("RTRIM(' abc ') = ' abc'");
+ testExpressionConditionTransform("BTRIM(' abc ') = 'abc'");
+ testExpressionConditionTransform(
+ "LTRIM('\uD83D\uDE01x', '\uD83D\uDE00') = '\uD83D\uDE01x'");
+ testExpressionConditionTransform(
+ "RTRIM('x\uD87D\uDE00', '\uD83D\uDE00') = 'x\uD87D\uDE00'");
+ testExpressionConditionTransform("BTRIM('xyabcxy', 'xy') = 'abc'");
+ testExpressionConditionTransform("CONCAT_WS(',', 'a', null, 'b') =
'a,b'");
+ testExpressionConditionTransform("LPAD('hi', 5, '?') = '???hi'");
+ testExpressionConditionTransform("RPAD('hi', 5, '?') = 'hi???'");
+ testExpressionConditionTransform("REPLACE('hello', 'l', 'x') =
'hexxo'");
+ testExpressionConditionTransform("REPEAT('ab', 3) = 'ababab'");
+ testExpressionConditionTransform("LEFT('abcdef', 2) = 'ab'");
+ testExpressionConditionTransform("RIGHT('abcdef', 2) = 'ef'");
+ testExpressionConditionTransform("LEFT('\uD83D\uDE00x', 1) =
'\uD83D\uDE00'");
+ testExpressionConditionTransform("RIGHT('x\uD83D\uDE00', 1) =
'\uD83D\uDE00'");
+ testExpressionConditionTransform("STARTSWITH('abcdef', 'abc')");
+ testExpressionConditionTransform("ENDSWITH('abcdef', 'def')");
+ testExpressionConditionTransform("TO_BASE64('hello') = 'aGVsbG8='");
+ testExpressionConditionTransform("FROM_BASE64('aGVsbG8=') = 'hello'");
+ testExpressionConditionTransform("TO_BASE64(FROM_BASE64('wyg=')) =
'77+9KA=='");
testExpressionConditionTransform("SUBSTR('ABC', -1) = 'C'");
testExpressionConditionTransform("SUBSTR('ABC', -2, 2) = 'BC'");
testExpressionConditionTransform("SUBSTR('ABC', 0) = 'ABC'");
@@ -2974,6 +3004,115 @@ class PostTransformOperatorTest {
testExpressionConditionTransform("cast(null as TIMESTAMP(3)) is null");
}
+ @Test
+ void testBinaryStartsWithAndEndsWithTransform() throws Exception {
+ TableId tableId = TableId.tableId("binary_string_functions");
+ Schema schema =
+ Schema.newBuilder()
+ .physicalColumn("value_", DataTypes.VARBINARY(16))
+ .physicalColumn("prefix_", DataTypes.VARBINARY(4))
+ .physicalColumn("suffix_", DataTypes.VARBINARY(4))
+ .build();
+ PostTransformOperator transform =
+ PostTransformOperator.newBuilder()
+ .addTransform(
+ tableId.identifier(),
+ null,
+ "STARTSWITH(value_, prefix_) AND
ENDSWITH(value_, suffix_)")
+ .addTimezone("UTC")
+ .build();
+ RegularEventOperatorTestHarness<PostTransformOperator, Event>
testHarness =
+ RegularEventOperatorTestHarness.with(transform, 1);
+ testHarness.open();
+
+ CreateTableEvent createTableEvent = new CreateTableEvent(tableId,
schema);
+ BinaryRecordDataGenerator recordDataGenerator =
+ new BinaryRecordDataGenerator(((RowType)
schema.toRowDataType()));
+ DataChangeEvent insertEvent =
+ DataChangeEvent.insertEvent(
+ tableId,
+ recordDataGenerator.generate(
+ new Object[] {
+ new byte[] {1, 2, 3, 4},
+ new byte[] {1, 2},
+ new byte[] {3, 4}
+ }));
+
+ transform.processElement(new StreamRecord<>(createTableEvent));
+ Assertions.assertThat(testHarness.getOutputRecords().poll())
+ .isEqualTo(new StreamRecord<>(createTableEvent));
+ transform.processElement(new StreamRecord<>(insertEvent));
+ Assertions.assertThat(testHarness.getOutputRecords().poll())
+ .isEqualTo(new StreamRecord<>(insertEvent));
+ testHarness.close();
+ }
+
+ @Test
+ void testBinaryOverlayPositionAndBase64Transform() throws Exception {
+ TableId tableId = TableId.tableId("binary_overlay_position_base64");
+ Schema schema =
+ Schema.newBuilder()
+ .physicalColumn("value_", DataTypes.VARBINARY(16))
+ .physicalColumn("replacement_", DataTypes.VARBINARY(4))
+ .physicalColumn("needle_", DataTypes.VARBINARY(4))
+ .physicalColumn("encoded_", DataTypes.STRING())
+ .build();
+ Schema expectedSchema =
+ Schema.newBuilder()
+ .physicalColumn("overlaid_", DataTypes.VARBINARY(20))
+ .physicalColumn("position_", DataTypes.INT())
+ .physicalColumn("encoded_value_", DataTypes.STRING())
+ .physicalColumn("decoded_", DataTypes.VARBINARY(65536))
+ .build();
+ PostTransformOperator transform =
+ PostTransformOperator.newBuilder()
+ .addTransform(
+ tableId.identifier(),
+ "OVERLAY(value_ PLACING replacement_ FROM 2
FOR 2) AS overlaid_, "
+ + "POSITION(needle_ IN value_) AS
position_, "
+ + "TO_BASE64(value_) AS
encoded_value_, "
+ + "FROM_BASE64_BINARY(encoded_) AS
decoded_",
+ null)
+ .addTimezone("UTC")
+ .build();
+ RegularEventOperatorTestHarness<PostTransformOperator, Event>
testHarness =
+ RegularEventOperatorTestHarness.with(transform, 1);
+ testHarness.open();
+
+ BinaryRecordDataGenerator recordDataGenerator =
+ new BinaryRecordDataGenerator(((RowType)
schema.toRowDataType()));
+ BinaryRecordDataGenerator expectedRecordDataGenerator =
+ new BinaryRecordDataGenerator(((RowType)
expectedSchema.toRowDataType()));
+ DataChangeEvent insertEvent =
+ DataChangeEvent.insertEvent(
+ tableId,
+ recordDataGenerator.generate(
+ new Object[] {
+ new byte[] {1, 2, 3, 4},
+ new byte[] {8, 9},
+ new byte[] {2, 3},
+ new BinaryStringData("wyg=")
+ }));
+ DataChangeEvent expectedInsertEvent =
+ DataChangeEvent.insertEvent(
+ tableId,
+ expectedRecordDataGenerator.generate(
+ new Object[] {
+ new byte[] {1, 8, 9, 4},
+ 2,
+ new BinaryStringData("AQIDBA=="),
+ new byte[] {(byte) 0xC3, 0x28}
+ }));
+
+ transform.processElement(new StreamRecord<>(new
CreateTableEvent(tableId, schema)));
+ Assertions.assertThat(testHarness.getOutputRecords().poll())
+ .isEqualTo(new StreamRecord<>(new CreateTableEvent(tableId,
expectedSchema)));
+ transform.processElement(new StreamRecord<>(insertEvent));
+ Assertions.assertThat(testHarness.getOutputRecords().poll())
+ .isEqualTo(new StreamRecord<>(expectedInsertEvent));
+ testHarness.close();
+ }
+
private void testExpressionConditionTransform(String expression) throws
Exception {
PostTransformOperator transform =
PostTransformOperator.newBuilder()
diff --git
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
index 281ea4721..c46525cd7 100644
---
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
+++
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
@@ -33,6 +33,7 @@ import org.apache.calcite.jdbc.CalciteSchema;
import org.apache.calcite.prepare.CalciteCatalogReader;
import org.apache.calcite.rel.type.RelDataTypeSystem;
import org.apache.calcite.runtime.CalciteContextException;
+import org.apache.calcite.sql.SqlBasicCall;
import org.apache.calcite.sql.SqlNode;
import org.apache.calcite.sql.SqlSelect;
import org.apache.calcite.sql.type.SqlTypeFactoryImpl;
@@ -221,6 +222,52 @@ class TransformParserTest {
testFilterExpression("lower(id)", "lower(id)");
testFilterExpression("concat(a,b)", "concat(a, b)");
testFilterExpression("SUBSTR(a,1)", "substr(a, 1)");
+ testFilterExpression("OVERLAY(id PLACING 'x' FROM 2)", "overlay(id,
\"x\", 2)");
+ testFilterExpression("OVERLAY(id PLACING 'x' FROM 2 FOR 3)",
"overlay(id, \"x\", 2, 3)");
+ testFilterExpression("POSITION('b' IN id)", "position(\"b\", id)");
+ testFilterExpression("POSITION('b' IN id FROM 2)", "position(\"b\",
id, 2)");
+ testFilterExpression("LOCATE('b', id)", "locate(\"b\", id)");
+ testFilterExpression("LOCATE('b', id, 2)", "locate(\"b\", id, 2)");
+ testFilterExpression("INSTR(id, 'b')", "instr(id, \"b\")");
+ testFilterExpression("LTRIM(id)", "ltrim(id)");
+ testFilterExpression("LTRIM(id, 'x')", "ltrim(id, \"x\")");
+ testFilterExpression("RTRIM(id)", "rtrim(id)");
+ testFilterExpression("RTRIM(id, 'x')", "rtrim(id, \"x\")");
+ testFilterExpression("BTRIM(id)", "btrim(id)");
+ testFilterExpression("BTRIM(id, 'x')", "btrim(id, \"x\")");
+ testFilterExpression("CONCAT_WS(',', a, b)", "concatWs(\",\", a, b)");
+ testFilterExpression("LPAD(id, 5, 'x')", "lpad(id, 5, \"x\")");
+ testFilterExpression("RPAD(id, 5, 'x')", "rpad(id, 5, \"x\")");
+ testFilterExpression("REPLACE(id, 'a', 'b')", "replace(id, \"a\",
\"b\")");
+ testFilterExpression("REPEAT(id, 2)", "repeat(id, 2)");
+ testFilterExpression("LEFT(id, 2)", "left(id, 2)");
+ testFilterExpression("RIGHT(id, 2)", "right(id, 2)");
+ testFilterExpression("STARTSWITH(id, 'a')", "startswith(id, \"a\")");
+ testFilterExpression("ENDSWITH(id, 'a')", "endswith(id, \"a\")");
+ List<Column> binaryColumns =
+ Arrays.asList(
+ Column.physicalColumn("binary_value",
DataTypes.VARBINARY(16)),
+ Column.physicalColumn("binary_part",
DataTypes.BINARY(2)));
+ testFilterExpressionWithColumns(
+ "STARTSWITH(binary_value, binary_part)",
+ "startswith(binary_value, binary_part)",
+ binaryColumns);
+ testFilterExpressionWithColumns(
+ "ENDSWITH(binary_value, binary_part)",
+ "endswith(binary_value, binary_part)",
+ binaryColumns);
+ testFilterExpressionWithColumns(
+ "POSITION(binary_part IN binary_value)",
+ "position(binary_part, binary_value)",
+ binaryColumns);
+ testFilterExpressionWithColumns(
+ "OVERLAY(binary_value PLACING binary_part FROM 1)",
+ "overlay(binary_value, binary_part, 1)",
+ binaryColumns);
+ testFilterExpression("TO_BASE64(id)", "toBase64(id)");
+ testFilterExpression("TO_BASE64(binary_value)",
"toBase64(binary_value)");
+ testFilterExpression("FROM_BASE64(id)", "fromBase64(id)");
+ testFilterExpression("FROM_BASE64_BINARY(id)", "fromBase64Binary(id)");
testFilterExpression("id like '^[a-zA-Z]'", "like(id, \"^[a-zA-Z]\")");
testFilterExpression("id not like '^[a-zA-Z]'", "notLike(id,
\"^[a-zA-Z]\")");
testFilterExpression("id like 'A$%' escape '$'", "like(id, \"A$%\",
\"$\")");
@@ -419,6 +466,80 @@ class TransformParserTest {
testFilterExpression("try_parse_json(jsonStr)",
"tryParseJson(jsonStr)");
}
+ @Test
+ void testStringFunctionArgumentValidation() {
+ List<Column> columns =
+ Arrays.asList(
+ Column.physicalColumn("s", DataTypes.STRING()),
+ Column.physicalColumn("binary_value",
DataTypes.VARBINARY(16)),
+ Column.physicalColumn("binary_part",
DataTypes.BINARY(2)));
+ SqlSelect positionSelect =
+ TransformParser.parseSelect(
+ "SELECT POSITION(binary_part IN binary_value) AS
position_value FROM TB");
+ SqlBasicCall positionAlias = (SqlBasicCall)
positionSelect.getSelectList().get(0);
+ SqlBasicCall positionCall = (SqlBasicCall) positionAlias.operand(0);
+
+ Assertions.assertThat(positionCall.getOperator())
+ .isSameAs(TransformSqlOperatorTable.POSITION);
+
+ Assertions.assertThatThrownBy(
+ () ->
+ TransformParser.generateProjectionColumns(
+ "INSTR(s, 'a', 1) AS invalid_value",
+ columns,
+ Collections.emptyList(),
+ new SupportedMetadataColumn[0]))
+ .isExactlyInstanceOf(CalciteContextException.class)
+ .hasMessageContaining("INSTR");
+ Assertions.assertThatThrownBy(
+ () ->
+ TransformParser.generateProjectionColumns(
+ "LPAD(s, true, 'x') AS invalid_value",
+ columns,
+ Collections.emptyList(),
+ new SupportedMetadataColumn[0]))
+ .isExactlyInstanceOf(CalciteContextException.class)
+ .hasMessageContaining("LPAD");
+ Assertions.assertThatThrownBy(
+ () ->
+ TransformParser.generateProjectionColumns(
+ "OVERLAY(s PLACING binary_part FROM 1)
AS invalid_value",
+ columns,
+ Collections.emptyList(),
+ new SupportedMetadataColumn[0]))
+ .isExactlyInstanceOf(CalciteContextException.class)
+ .hasMessageContaining("is not comparable to");
+ Assertions.assertThat(
+ TransformParser.generateProjectionColumns(
+ "POSITION(binary_part IN binary_value) AS
position_value",
+ columns,
+ Collections.emptyList(),
+ new SupportedMetadataColumn[0]))
+ .hasSize(1);
+ Assertions.assertThat(
+ TransformParser.generateProjectionColumns(
+ "POSITION('a' IN s FROM 2) AS position_value",
+ columns,
+ Collections.emptyList(),
+ new SupportedMetadataColumn[0]))
+ .hasSize(1);
+ Assertions.assertThat(
+ TransformParser.generateProjectionColumns(
+ "OVERLAY(binary_value PLACING binary_part FROM
1) AS overlaid_value, "
+ + "POSITION(binary_part IN
binary_value) AS position_value, "
+ + "TO_BASE64(binary_value) AS
encoded_value, "
+ + "FROM_BASE64_BINARY(s) AS
decoded_value",
+ columns,
+ Collections.emptyList(),
+ new SupportedMetadataColumn[0]))
+ .extracting(ProjectionColumn::getDataType)
+ .containsExactly(
+ DataTypes.VARBINARY(18),
+ DataTypes.INT(),
+ DataTypes.STRING(),
+ DataTypes.VARBINARY(65536));
+ }
+
@Test
void testTranslateLogicalFilterToJaninoExpressionByNullability() {
List<Column> columns =