From ebef58a31364caf352ea809ed9252d2c3f18cd30 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=98=A5=E6=A0=96?= Date: Thu, 23 Jul 2026 13:44:48 +0800 Subject: [PATCH 1/3] [FLINK-40203][runtime] Add an explicit option for Flink SQL-compatible YAML Transform expression semantics --- .../docs/core-concept/data-pipeline.md | 17 ++++ .../docs/core-concept/data-pipeline.md | 17 ++++ .../YamlPipelineDefinitionParserTest.java | 23 +++++ .../cdc/common/pipeline/PipelineOptions.java | 8 ++ .../TransformExpressionSemantics.java | 30 ++++++ .../composer/flink/FlinkPipelineComposer.java | 4 + .../flink/translator/TransformTranslator.java | 5 + .../translator/TransformTranslatorTest.java | 3 + .../composer/specs/TransformSpecsITCase.java | 10 ++ .../src/test/resources/specs/comparison.yaml | 29 ++++++ .../functions/impl/ComparisonFunctions.java | 64 ++++++++++++ .../functions/impl/StringFunctions.java | 11 +++ .../transform/PostTransformOperator.java | 13 ++- .../PostTransformOperatorBuilder.java | 15 ++- .../transform/TransformFilterProcessor.java | 26 ++++- .../TransformProjectionProcessor.java | 9 +- .../async/AsyncPostTransformFunction.java | 13 ++- .../AsyncPostTransformFunctionBuilder.java | 9 ++ .../cdc/runtime/parser/JaninoCompiler.java | 99 +++++++++++++++++-- .../cdc/runtime/parser/TransformParser.java | 79 ++++++++++++++- .../impl/ComparisonFunctionsTest.java | 29 ++++++ .../functions/impl/StringFunctionsTest.java | 16 +++ .../runtime/parser/TransformParserTest.java | 30 ++++++ 23 files changed, 534 insertions(+), 25 deletions(-) create mode 100644 flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java diff --git a/docs/content.zh/docs/core-concept/data-pipeline.md b/docs/content.zh/docs/core-concept/data-pipeline.md index 2e228574c3d..26e484bf142 100644 --- a/docs/content.zh/docs/core-concept/data-pipeline.md +++ b/docs/content.zh/docs/core-concept/data-pipeline.md @@ -119,6 +119,7 @@ under the License. | `local-time-zone` | 作业级别的本地时区。 | optional | | `execution.runtime-mode` | pipeline 的运行模式,包含 STREAMING 和 BATCH,默认值是 STREAMING。 | optional | | `route-mode` | [路由]({{< ref "docs/core-concept/route" >}}#路由模式)规则的匹配模式。可选值:`ALL_MATCH`(默认,应用所有匹配的规则)或 `FIRST_MATCH`(只应用第一个匹配的规则)。 | optional | +| `transform.expression.semantics` | Transform 表达式的求值语义。可选值:`LEGACY`(默认,保持现有行为)或 `FLINK_SQL`(对支持的谓词使用与 Flink SQL 兼容的语义)。 | optional | | `schema.change.behavior` | 如何处理 [schema 变更]({{< ref "docs/core-concept/schema-evolution" >}})。可选值:[`exception`]({{< ref "docs/core-concept/schema-evolution" >}}#exception-mode)、[`evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#evolve-mode)、[`try_evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#tryevolve-mode)、[`lenient`]({{< ref "docs/core-concept/schema-evolution" >}}#lenient-mode)(默认值)或 [`ignore`]({{< ref "docs/core-concept/schema-evolution" >}}#ignore-mode)。 | optional | | `schema.operator.uid` | Schema 算子的唯一 ID。此 ID 用于算子间通信,必须在所有算子中保持唯一。**已废弃**:请使用 `operator.uid.prefix` 代替。 | optional | | `schema-operator.rpc-timeout` | SchemaOperator 等待下游 SchemaChangeEvent 应用完成的超时时间,默认值是 3 分钟。 | optional | @@ -135,3 +136,19 @@ under the License. 异步 PostTransform 适用于 AI 模型调用等 I/O 密集型表达式。DataChangeEvent 可能并发调用 UDF 和 AI 模型客户端,但输出和所有 SchemaChangeEvent 仍保持有序,因此 UDF 和 AI 模型客户端实现必须是线程安全的。为保证 schema 状态一致,checkpoint 或 savepoint 前会等待尚未完成的异步请求,长时间运行的请求可能会延长 checkpoint 时间。savepoint 只支持在并行度不变时恢复,并且不能在从已有 savepoint 恢复时开启或关闭此选项。 注意:虽然上述参数都是可选的,但至少需要指定其中一个。`pipeline` 部分是必需的,不能为空。 + +## Transform 表达式语义 + +`transform.expression.semantics` 配置项控制行级 Transform 过滤和投影表达式的求值语义: + +* `LEGACY` 是默认值,用于保持现有 Pipeline 行为。comparison、`BETWEEN` 和 `IN` 在操作数为 `NULL` 时仍返回布尔值;两参数 `LIKE` 将 pattern 作为 Java 正则表达式,并使用子串匹配。 +* `FLINK_SQL` 将 `=`、`<>`、`<`、`<=`、`>`、`>=`、`BETWEEN` / `NOT BETWEEN`、`IN` / `NOT IN` 和两参数 `LIKE` / `NOT LIKE` 与 Flink SQL 对齐。这些谓词使用 SQL 三值逻辑,因此结果可能为 `UNKNOWN`(以 `NULL` 表示)。两参数 `LIKE` 使用 SQL `%` 和 `_` 通配符,并对整个字符串进行匹配。 + +配置示例: + +```yaml +pipeline: + transform.expression.semantics: FLINK_SQL +``` + +该配置会影响行级 Transform 的结果。过滤表达式仅在结果为 `TRUE` 时保留数据行,因此 `FALSE` 和 `UNKNOWN` 都会被过滤;投影表达式会将 `UNKNOWN` 保留为 `NULL`。 diff --git a/docs/content/docs/core-concept/data-pipeline.md b/docs/content/docs/core-concept/data-pipeline.md index 664c5c5db4e..acb75e0188f 100644 --- a/docs/content/docs/core-concept/data-pipeline.md +++ b/docs/content/docs/core-concept/data-pipeline.md @@ -121,6 +121,7 @@ Note that whilst the parameters are each individually optional, at least one of | `local-time-zone` | The local time zone defines current session time zone id. | optional | | `execution.runtime-mode` | The runtime mode of the pipeline includes STREAMING and BATCH, with the default value being STREAMING. | optional | | `route-mode` | The matching mode for [route]({{< ref "docs/core-concept/route" >}}#route-mode) rules. One of: `ALL_MATCH` (default, apply all matching rules) or `FIRST_MATCH` (apply only the first matching rule). | optional | +| `transform.expression.semantics` | The semantics used to evaluate transform expressions. One of: `LEGACY` (default, preserves existing behavior) or `FLINK_SQL` (uses Flink SQL-compatible semantics for supported predicates). | optional | | `schema.change.behavior` | How to handle [changes in schema]({{< ref "docs/core-concept/schema-evolution" >}}). One of: [`exception`]({{< ref "docs/core-concept/schema-evolution" >}}#exception-mode), [`evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#evolve-mode), [`try_evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#tryevolve-mode), [`lenient`]({{< ref "docs/core-concept/schema-evolution" >}}#lenient-mode) (default) or [`ignore`]({{< ref "docs/core-concept/schema-evolution" >}}#ignore-mode). | optional | | `schema.operator.uid` | The unique ID for schema operator. This ID will be used for inter-operator communications and must be unique across operators. **Deprecated**: use `operator.uid.prefix` instead. | optional | | `schema-operator.rpc-timeout` | The timeout time for SchemaOperator to wait downstream SchemaChangeEvent applying finished, the default value is 3 minutes. | optional | @@ -137,3 +138,19 @@ Note that whilst the parameters are each individually optional, at least one of Asynchronous post-transform execution is intended for I/O-bound expressions such as AI model calls. Data change events may invoke UDFs and AI model clients concurrently, while their output and all schema changes remain ordered. UDF and AI model client implementations must therefore be thread-safe. Pending asynchronous requests are completed before a checkpoint or savepoint to keep schema state consistent, so long-running requests may extend checkpoint duration. Savepoint restore is supported only with unchanged parallelism, and this option must not be enabled or disabled when restoring an existing savepoint. NOTE: Whilst the above parameters are each individually optional, at least one of them must be specified. The `pipeline` section is mandatory and cannot be empty. + +## Transform Expression Semantics + +The `transform.expression.semantics` option controls the evaluation semantics of row-level transform filters and projections: + +* `LEGACY` is the default and preserves existing pipeline behavior. Comparisons, `BETWEEN`, and `IN` return a boolean result when an operand is `NULL`. Two-argument `LIKE` treats its pattern as a Java regular expression and uses substring matching. +* `FLINK_SQL` aligns `=`, `<>`, `<`, `<=`, `>`, `>=`, `BETWEEN` / `NOT BETWEEN`, `IN` / `NOT IN`, and two-argument `LIKE` / `NOT LIKE` with Flink SQL. These predicates use SQL three-valued logic, so a result may be `UNKNOWN` (represented as `NULL`). Two-argument `LIKE` uses SQL `%` and `_` wildcards and matches the entire string. + +For example: + +```yaml +pipeline: + transform.expression.semantics: FLINK_SQL +``` + +This option can change row-level transform results. A filter keeps a row only when its expression evaluates to `TRUE`, so both `FALSE` and `UNKNOWN` are filtered out. A projected predicate preserves `UNKNOWN` as `NULL`. diff --git a/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java b/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java index a1673723c17..0d4842cdf3f 100644 --- a/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java +++ b/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java @@ -20,6 +20,7 @@ import org.apache.flink.cdc.common.configuration.Configuration; import org.apache.flink.cdc.common.event.SchemaChangeEventType; import org.apache.flink.cdc.common.pipeline.PipelineOptions; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.composer.definition.ModelDef; import org.apache.flink.cdc.composer.definition.PipelineDef; import org.apache.flink.cdc.composer.definition.RouteDef; @@ -105,6 +106,28 @@ void testEvaluateDefaultLocalTimeZone() throws Exception { .isNotEqualTo(PIPELINE_LOCAL_TIME_ZONE.defaultValue()); } + @Test + void testTransformExpressionSemantics() throws Exception { + YamlPipelineDefinitionParser parser = new YamlPipelineDefinitionParser(); + PipelineDef pipelineDef = + parser.parse( + "source:\n" + + " type: foo\n" + + "sink:\n" + + " type: bar\n" + + "pipeline:\n" + + " transform.expression.semantics: FLINK_SQL\n", + new Configuration()); + + assertThat( + pipelineDef + .getConfig() + .get(PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS)) + .isEqualTo(TransformExpressionSemantics.FLINK_SQL); + assertThat(new Configuration().get(PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS)) + .isEqualTo(TransformExpressionSemantics.LEGACY); + } + @Test void testEvaluateDefaultRpcTimeOut() throws Exception { URL resource = Resources.getResource("definitions/pipeline-definition-minimized.yaml"); diff --git a/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java index 3b39645d096..95d8827ad89 100644 --- a/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java +++ b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java @@ -203,5 +203,13 @@ public class PipelineOptions { .withDescription( "The number of worker threads used by each async post-transform task."); + public static final ConfigOption + PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS = + ConfigOptions.key("transform.expression.semantics") + .enumType(TransformExpressionSemantics.class) + .defaultValue(TransformExpressionSemantics.LEGACY) + .withDescription( + "Semantics used to evaluate transform expressions. LEGACY preserves historical behavior, while FLINK_SQL enables Flink SQL-compatible semantics for supported predicates."); + private PipelineOptions() {} } diff --git a/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java new file mode 100644 index 00000000000..5dabcea294f --- /dev/null +++ b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java @@ -0,0 +1,30 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.cdc.common.pipeline; + +import org.apache.flink.cdc.common.annotation.PublicEvolving; + +/** Defines the semantics used to evaluate pipeline transform expressions. */ +@PublicEvolving +public enum TransformExpressionSemantics { + /** Preserves the historical pipeline transform expression behavior. */ + LEGACY, + + /** Uses Flink SQL-compatible semantics for supported transform expressions. */ + FLINK_SQL +} diff --git a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java index 3b73a283e50..79a6c580644 100644 --- a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java +++ b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java @@ -200,6 +200,8 @@ private void translate(StreamExecutionEnvironment env, PipelineDef pipelineDef) pipelineDefConfig.get(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE), pipelineDefConfig.get( PipelineOptions.PIPELINE_TRANSFORM_DECIMAL_PRECISION_MODE), + pipelineDefConfig.get( + PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS), pipelineDef.getUdfs(), pipelineDef.getModels(), dataSource.supportedMetadataColumns(), @@ -220,6 +222,8 @@ private void translate(StreamExecutionEnvironment env, PipelineDef pipelineDef) pipelineDefConfig.get(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE), pipelineDefConfig.get( PipelineOptions.PIPELINE_TRANSFORM_DECIMAL_PRECISION_MODE), + pipelineDefConfig.get( + PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS), pipelineDef.getUdfs(), pipelineDef.getModels(), dataSource.supportedMetadataColumns(), diff --git a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java index a8c650c07b2..b8b65f82163 100644 --- a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java +++ b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java @@ -25,6 +25,7 @@ import org.apache.flink.cdc.common.model.AiModelClientFactory; import org.apache.flink.cdc.common.model.ModelContext; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import org.apache.flink.cdc.composer.definition.ModelDef; import org.apache.flink.cdc.composer.definition.TransformDef; @@ -115,6 +116,7 @@ public DataStream translatePostTransform( List transforms, String timezone, DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics, List udfFunctions, List models, SupportedMetadataColumn[] supportedMetadataColumns, @@ -140,6 +142,7 @@ public DataStream translatePostTransform( } postTransformFunctionBuilder.addTimezone(timezone); postTransformFunctionBuilder.addDecimalPrecisionMode(decimalPrecisionMode); + postTransformFunctionBuilder.addExpressionSemantics(expressionSemantics); postTransformFunctionBuilder.addUdfFunctions( udfFunctions.stream().map(this::udfDefToUDFTuple).collect(Collectors.toList())); postTransformFunctionBuilder.addUdfFunctions( @@ -159,6 +162,7 @@ public DataStream translateAsyncPostTransform( List transforms, String timezone, DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics, List udfFunctions, List models, SupportedMetadataColumn[] supportedMetadataColumns, @@ -195,6 +199,7 @@ public DataStream translateAsyncPostTransform( } asyncPostTransformFunctionBuilder.addTimezone(timezone); asyncPostTransformFunctionBuilder.addDecimalPrecisionMode(decimalPrecisionMode); + asyncPostTransformFunctionBuilder.addExpressionSemantics(expressionSemantics); asyncPostTransformFunctionBuilder.addUdfFunctions( udfFunctions.stream().map(this::udfDefToUDFTuple).collect(Collectors.toList())); asyncPostTransformFunctionBuilder.addUdfFunctions( diff --git a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslatorTest.java b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslatorTest.java index 12aa4c3ca81..4206efd0cc0 100644 --- a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslatorTest.java +++ b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslatorTest.java @@ -21,6 +21,7 @@ import org.apache.flink.cdc.common.event.Event; import org.apache.flink.cdc.common.event.TableId; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.PipelineOptions; import org.apache.flink.cdc.common.schema.Schema; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import org.apache.flink.cdc.common.types.DataTypes; @@ -60,6 +61,8 @@ void testTranslateAsyncPostTransformUsesStateConsistentOperator() { Collections.singletonList(transform), "UTC", DecimalPrecisionMode.UP_TO_19, + PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS + .defaultValue(), Collections.emptyList(), Collections.emptyList(), new SupportedMetadataColumn[0], diff --git a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java index 7c5f22e23df..973ffa71bf7 100644 --- a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java +++ b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java @@ -35,6 +35,7 @@ import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; import org.apache.flink.cdc.common.pipeline.PipelineOptions; import org.apache.flink.cdc.common.pipeline.SchemaChangeBehavior; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.schema.Schema; import org.apache.flink.cdc.common.types.DataType; import org.apache.flink.cdc.common.types.variant.Variant; @@ -365,6 +366,11 @@ private static Stream loadTestSpec(Path specPath) { DecimalPrecisionMode.valueOf( specNode.get("decimal-precision-mode").asText().toUpperCase()); } + if (specNode.has("expression-semantics")) { + spec.expressionSemantics = + TransformExpressionSemantics.valueOf( + specNode.get("expression-semantics").asText().toUpperCase()); + } if (specNode.has("projection")) { spec.projectionRules = List.of( @@ -409,6 +415,8 @@ static class TestSpec { public String ignore; public String timeZone = "UTC"; public DecimalPrecisionMode decimalPrecisionMode = DecimalPrecisionMode.UP_TO_19; + public TransformExpressionSemantics expressionSemantics = + TransformExpressionSemantics.LEGACY; public List projectionRules = new ArrayList<>(); public @Nullable String filterRule; public @Nullable String primaryKey; @@ -493,6 +501,8 @@ void runTransformSpecs(String group, String name, TestSpec spec) throws Exceptio Configuration pipelineConfig = new Configuration(); pipelineConfig.set(PipelineOptions.PIPELINE_PARALLELISM, 1); pipelineConfig.set(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE, spec.timeZone); + pipelineConfig.set( + PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS, spec.expressionSemantics); pipelineConfig.set( PipelineOptions.PIPELINE_SCHEMA_CHANGE_BEHAVIOR, SchemaChangeBehavior.EVOLVE); pipelineConfig.set( diff --git a/flink-cdc-composer/src/test/resources/specs/comparison.yaml b/flink-cdc-composer/src/test/resources/specs/comparison.yaml index 85d7e80b49b..0f7cb560682 100644 --- a/flink-cdc-composer/src/test/resources/specs/comparison.yaml +++ b/flink-cdc-composer/src/test/resources/specs/comparison.yaml @@ -370,3 +370,32 @@ DataChangeEvent{tableId=foo.bar.baz, before=[-1, true, true, true, true, true, true, true, true, true, true, true], after=[], op=DELETE, meta=()} DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, true, true, true, true, true, true, true, true, true, true, true], op=INSERT, meta=()} DataChangeEvent{tableId=foo.bar.baz, before=[0, true, true, true, true, true, true, true, true, true, true, true], after=[], op=DELETE, meta=()} +- do: Flink SQL Predicate Semantics + expression-semantics: FLINK_SQL + projection: |- + id_ + int_ = 4 AS equals_ + int_ <> 4 AS not_equals_ + int_ BETWEEN 3 AND 5 AS between_ + int_ NOT BETWEEN 3 AND 5 AS not_between_ + int_ IN (4, CAST(NULL AS INT)) AS in_ + int_ NOT IN (-4, CAST(NULL AS INT)) AS not_in_ + string_ LIKE 'From % Lie' AS like_ + string_ NOT LIKE '%Alice%' AS not_like_ + primary-key: id_ + expect: |- + CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT NULL 'Identifier',`equals_` BOOLEAN,`not_equals_` BOOLEAN,`between_` BOOLEAN,`not_between_` BOOLEAN,`in_` BOOLEAN,`not_in_` BOOLEAN,`like_` BOOLEAN,`not_like_` BOOLEAN}, primaryKeys=id_, options=()} + DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, true, false, true, false, true, null, true, true], op=INSERT, meta=()} + DataChangeEvent{tableId=foo.bar.baz, before=[1, true, false, true, false, true, null, true, true], after=[-1, false, true, false, true, null, false, false, true], op=UPDATE, meta=()} + DataChangeEvent{tableId=foo.bar.baz, before=[-1, false, true, false, true, null, false, false, true], after=[], op=DELETE, meta=()} + DataChangeEvent{tableId=foo.bar.baz, before=[], after=[0, null, null, null, null, null, null, null, null], op=INSERT, meta=()} + DataChangeEvent{tableId=foo.bar.baz, before=[0, null, null, null, null, null, null, null, null], after=[], op=DELETE, meta=()} +- do: Flink SQL Unknown Filter Semantics + expression-semantics: FLINK_SQL + projection: id_, int_ + filter: int_ IN (4, CAST(NULL AS INT)) + primary-key: id_ + expect: |- + CreateTableEvent{tableId=foo.bar.baz, schema=columns={`id_` BIGINT NOT NULL 'Identifier',`int_` INT}, primaryKeys=id_, options=()} + DataChangeEvent{tableId=foo.bar.baz, before=[], after=[1, 4], op=INSERT, meta=()} + DataChangeEvent{tableId=foo.bar.baz, before=[1, 4], after=[], op=DELETE, meta=()} diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java index e318e45e1b3..a87ce16a798 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java @@ -27,6 +27,17 @@ public static boolean valueEquals(Object object1, Object object2) { return (object1 != null && object2 != null) && object1.equals(object2); } + public static Boolean sqlValueEquals(Object object1, Object object2) { + if (object1 == null || object2 == null) { + return null; + } + return object1.equals(object2); + } + + public static Boolean sqlNotEquals(Object object1, Object object2) { + return LogicalFunctions.not(sqlValueEquals(object1, object2)); + } + public static boolean isDistinctFrom(Object object1, Object object2) { if (object1 == null || object2 == null) { return object1 != object2; @@ -61,6 +72,13 @@ public static boolean greaterThan(Object lhs, Object rhs) { return universalCompares(lhs, rhs) > 0; } + public static Boolean sqlGreaterThan(Object lhs, Object rhs) { + if (lhs == null || rhs == null) { + return null; + } + return universalCompares(lhs, rhs) > 0; + } + public static boolean greaterThanOrEqual(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return false; @@ -68,6 +86,13 @@ public static boolean greaterThanOrEqual(Object lhs, Object rhs) { return universalCompares(lhs, rhs) >= 0; } + public static Boolean sqlGreaterThanOrEqual(Object lhs, Object rhs) { + if (lhs == null || rhs == null) { + return null; + } + return universalCompares(lhs, rhs) >= 0; + } + public static boolean lessThan(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return false; @@ -75,6 +100,13 @@ public static boolean lessThan(Object lhs, Object rhs) { return universalCompares(lhs, rhs) < 0; } + public static Boolean sqlLessThan(Object lhs, Object rhs) { + if (lhs == null || rhs == null) { + return null; + } + return universalCompares(lhs, rhs) < 0; + } + public static boolean lessThanOrEqual(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return false; @@ -82,6 +114,38 @@ public static boolean lessThanOrEqual(Object lhs, Object rhs) { return universalCompares(lhs, rhs) <= 0; } + public static Boolean sqlLessThanOrEqual(Object lhs, Object rhs) { + if (lhs == null || rhs == null) { + return null; + } + return universalCompares(lhs, rhs) <= 0; + } + + public static Boolean sqlBetween(Object value, Object minValue, Object maxValue) { + return LogicalFunctions.and( + sqlGreaterThanOrEqual(value, minValue), sqlLessThanOrEqual(value, maxValue)); + } + + public static Boolean sqlNotBetween(Object value, Object minValue, Object maxValue) { + return LogicalFunctions.not(sqlBetween(value, minValue, maxValue)); + } + + public static Boolean sqlIn(Object value, Object... values) { + boolean containsUnknown = false; + for (Object candidate : values) { + Boolean equals = sqlValueEquals(value, candidate); + if (Boolean.TRUE.equals(equals)) { + return true; + } + containsUnknown |= equals == null; + } + return containsUnknown ? null : false; + } + + public static Boolean sqlNotIn(Object value, Object... values) { + return LogicalFunctions.not(sqlIn(value, values)); + } + public static boolean betweenAsymmetric(String value, String minValue, String maxValue) { if (value == null) { return false; 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 c665b44df01..c219072c6cc 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 @@ -163,6 +163,13 @@ public static boolean like(String str, String regex) { return Pattern.compile(regex).matcher(str).find(); } + public static Boolean sqlLike(String str, String pattern) { + if (str == null || pattern == null) { + return null; + } + return SqlFunctions.like(str, pattern); + } + public static Boolean like(String str, String pattern, String escape) { if (str == null || pattern == null || escape == null) { return null; @@ -174,6 +181,10 @@ public static boolean notLike(String str, String regex) { return !like(str, regex); } + public static Boolean sqlNotLike(String str, String pattern) { + return LogicalFunctions.not(sqlLike(str, pattern)); + } + public static Boolean notLike(String str, String pattern, String escape) { return LogicalFunctions.not(like(str, pattern, escape)); } diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperator.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperator.java index 36deef753fb..f3206f116c7 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperator.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperator.java @@ -31,6 +31,7 @@ import org.apache.flink.cdc.common.event.TableId; import org.apache.flink.cdc.common.model.AiModelClient; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.schema.Schema; import org.apache.flink.cdc.common.schema.Selectors; import org.apache.flink.cdc.common.udf.UserDefinedFunctionContext; @@ -79,6 +80,7 @@ public class PostTransformOperator extends AbstractStreamOperatorAdapter private final String timezone; private final DecimalPrecisionMode decimalPrecisionMode; + private final TransformExpressionSemantics expressionSemantics; private final List transformRules; private final Map hasAsteriskMap; private final Map> projectedColumnsMap; @@ -110,10 +112,12 @@ public static PostTransformOperatorBuilder newBuilder() { List transformRules, String timezone, DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics, List>> udfFunctions, Map modelClients) { this.timezone = timezone; this.decimalPrecisionMode = decimalPrecisionMode; + this.expressionSemantics = expressionSemantics; this.transformRules = transformRules; this.hasAsteriskMap = new HashMap<>(); this.projectedColumnsMap = new HashMap<>(); @@ -388,7 +392,8 @@ private Schema transformSchema(Schema preSchema, PostTransformer transformer) { preSchema.getColumns(), udfDescriptors, transformer.getSupportedMetadataColumns(), - decimalPrecisionMode); + decimalPrecisionMode, + expressionSemantics); return preSchema.copy( projectionColumns.stream() .map(ProjectionColumn::getColumn) @@ -469,7 +474,8 @@ private TransformProjectionProcessor getProjectionProcessor( udfDescriptors, udfFunctionInstances, postTransformer.getSupportedMetadataColumns(), - modelClients)); + modelClients, + expressionSemantics)); } return projectionProcessors.get(tableId, postTransformer); } @@ -499,7 +505,8 @@ private TransformFilterProcessor getFilterProcessor( udfDescriptors, udfFunctionInstances, postTransformer.getSupportedMetadataColumns(), - modelClients)); + modelClients, + expressionSemantics)); } } return filterProcessors.get(tableId, postTransformer); diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java index a5de56c7597..8653ad7fc67 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java @@ -21,6 +21,7 @@ import org.apache.flink.cdc.common.model.AiModelClient; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; import org.apache.flink.cdc.common.pipeline.PipelineOptions; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import javax.annotation.Nullable; @@ -36,6 +37,7 @@ public class PostTransformOperatorBuilder { private final List transformRules = new ArrayList<>(); private String timezone; private DecimalPrecisionMode decimalPrecisionMode = DecimalPrecisionMode.UP_TO_19; + private TransformExpressionSemantics expressionSemantics = TransformExpressionSemantics.LEGACY; private final List>> udfFunctions = new ArrayList<>(); private final Map modelClients = new LinkedHashMap<>(); @@ -116,6 +118,12 @@ public PostTransformOperatorBuilder addDecimalPrecisionMode( return this; } + public PostTransformOperatorBuilder addExpressionSemantics( + TransformExpressionSemantics expressionSemantics) { + this.expressionSemantics = expressionSemantics; + return this; + } + public PostTransformOperatorBuilder addUdfFunctions( List>> udfFunctions) { this.udfFunctions.addAll(udfFunctions); @@ -129,6 +137,11 @@ public PostTransformOperatorBuilder addModelClients(Map c public PostTransformOperator build() { return new PostTransformOperator( - transformRules, timezone, decimalPrecisionMode, udfFunctions, modelClients); + transformRules, + timezone, + decimalPrecisionMode, + expressionSemantics, + udfFunctions, + modelClients); } } diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java index 6e4eac45f36..79c1ae11a68 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java @@ -21,6 +21,7 @@ import org.apache.flink.cdc.common.converter.JavaClassConverter; import org.apache.flink.cdc.common.model.AiModelClient; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.schema.Column; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import org.apache.flink.cdc.runtime.parser.JaninoCompiler; @@ -52,6 +53,7 @@ public class TransformFilterProcessor { private final List udfFunctionInstances; private final Map supportedMetadataColumns; private final Map modelClients; + private final TransformExpressionSemantics expressionSemantics; private final TransformExpressionKey transformExpressionKey; private final ExpressionEvaluator expressionEvaluator; @@ -65,7 +67,8 @@ protected TransformFilterProcessor( List udfDescriptors, List udfFunctionInstances, Map supportedMetadataColumns, - Map modelClients) { + Map modelClients, + TransformExpressionSemantics expressionSemantics) { this.isNoOp = isNoOp; this.tableInfo = tableInfo; this.transformFilter = transformFilter; @@ -74,6 +77,7 @@ protected TransformFilterProcessor( this.udfFunctionInstances = udfFunctionInstances; this.supportedMetadataColumns = supportedMetadataColumns; this.modelClients = modelClients == null ? Collections.emptyMap() : modelClients; + this.expressionSemantics = expressionSemantics; if (isNoOp) { this.transformExpressionKey = null; @@ -94,7 +98,16 @@ protected TransformFilterProcessor( public static TransformFilterProcessor ofNoOp(DecimalPrecisionMode decimalPrecisionMode) { return new TransformFilterProcessor( - true, null, null, null, decimalPrecisionMode, null, null, null, null); + true, + null, + null, + null, + decimalPrecisionMode, + null, + null, + null, + null, + TransformExpressionSemantics.LEGACY); } public static TransformFilterProcessor of( @@ -105,7 +118,8 @@ public static TransformFilterProcessor of( List udfDescriptors, List udfFunctionInstances, SupportedMetadataColumn[] supportedMetadataColumns, - Map modelClients) { + Map modelClients, + TransformExpressionSemantics expressionSemantics) { Map supportedMetadataColumnsMap = new HashMap<>(); for (SupportedMetadataColumn supportedMetadataColumn : supportedMetadataColumns) { supportedMetadataColumnsMap.put( @@ -120,7 +134,8 @@ public static TransformFilterProcessor of( udfDescriptors, udfFunctionInstances, supportedMetadataColumnsMap, - modelClients); + modelClients, + expressionSemantics); } public boolean test(Object[] preRow, Object[] postRow, TransformContext context) { @@ -247,7 +262,8 @@ private TransformExpressionKey generateTransformExpressionKey( udfDescriptors, supportedMetadataColumns, transformFilter.getColumnNameMap(), - decimalPrecisionMode); + decimalPrecisionMode, + expressionSemantics); return TransformExpressionKey.of( transformFilter.getExpression(), diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformProjectionProcessor.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformProjectionProcessor.java index 84b4d3a9ad7..35a044f63ee 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformProjectionProcessor.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformProjectionProcessor.java @@ -19,6 +19,7 @@ import org.apache.flink.cdc.common.model.AiModelClient; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import org.apache.flink.cdc.common.utils.Preconditions; import org.apache.flink.cdc.runtime.parser.TransformParser; @@ -54,6 +55,7 @@ public class TransformProjectionProcessor { private final List udfFunctionInstances; private final List columnProcessors; private final SupportedMetadataColumn[] supportedMetadataColumns; + private final TransformExpressionSemantics expressionSemantics; private final Map supportedMetadataColumnsMap; private final Map modelClients; @@ -65,7 +67,8 @@ public TransformProjectionProcessor( List udfDescriptors, List udfFunctionInstances, SupportedMetadataColumn[] supportedMetadataColumns, - Map modelClients) { + Map modelClients, + TransformExpressionSemantics expressionSemantics) { this.changeInfo = changeInfo; this.projectionExpression = projectionExpression; this.timezone = timezone; @@ -74,6 +77,7 @@ public TransformProjectionProcessor( this.udfFunctionInstances = udfFunctionInstances; this.supportedMetadataColumns = supportedMetadataColumns; this.modelClients = modelClients; + this.expressionSemantics = expressionSemantics; // Construct a mapping table ad-hoc to accelerate looking-up Map supportedMetadataColumnsMap = new HashMap<>(); @@ -102,7 +106,8 @@ private List createProjectionColumnProcessors() { changeInfo.getPreTransformedSchema().getColumns(), udfDescriptors, supportedMetadataColumns, - decimalPrecisionMode); + decimalPrecisionMode, + expressionSemantics); List columnProcessors = projectionColumns.stream() diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunction.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunction.java index 531fde01be7..12ba10cb6c2 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunction.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunction.java @@ -34,6 +34,7 @@ import org.apache.flink.cdc.common.event.TableId; import org.apache.flink.cdc.common.model.AiModelClient; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.schema.Schema; import org.apache.flink.cdc.common.schema.Selectors; import org.apache.flink.cdc.common.udf.UserDefinedFunctionContext; @@ -122,6 +123,7 @@ public class AsyncPostTransformFunction extends RichAsyncFunction private final String timezone; private final DecimalPrecisionMode decimalPrecisionMode; + private final TransformExpressionSemantics expressionSemantics; private final List transformRules; private final Map tableInfoMap; @@ -160,6 +162,7 @@ public static AsyncPostTransformFunctionBuilder newBuilder() { List transformRules, String timezone, DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics, List>> udfFunctions, Map modelClients, int asyncWorkerThreads) { @@ -167,6 +170,7 @@ public static AsyncPostTransformFunctionBuilder newBuilder() { asyncWorkerThreads > 0, "Async worker threads must be greater than 0."); this.timezone = timezone; this.decimalPrecisionMode = decimalPrecisionMode; + this.expressionSemantics = expressionSemantics; this.transformRules = transformRules; this.tableInfoMap = new ConcurrentHashMap<>(); this.udfFunctions = udfFunctions; @@ -689,7 +693,8 @@ private Schema transformSchema(Schema preSchema, PostTransformer transformer) { preSchema.getColumns(), udfDescriptors, transformer.getSupportedMetadataColumns(), - decimalPrecisionMode); + decimalPrecisionMode, + expressionSemantics); return preSchema.copy( projectionColumns.stream() .map(ProjectionColumn::getColumn) @@ -760,7 +765,8 @@ private TransformProjectionProcessor getProjectionProcessor( udfDescriptors, udfFunctionInstances, postTransformer.getSupportedMetadataColumns(), - modelClients)); + modelClients, + expressionSemantics)); } return processors.get(tableId, postTransformer); } @@ -789,7 +795,8 @@ private TransformFilterProcessor getFilterProcessor( udfDescriptors, udfFunctionInstances, postTransformer.getSupportedMetadataColumns(), - modelClients)); + modelClients, + expressionSemantics)); } } return processors.get(tableId, postTransformer); diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java index affbe9cd76c..06fb48fc354 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java @@ -21,6 +21,7 @@ import org.apache.flink.cdc.common.model.AiModelClient; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; import org.apache.flink.cdc.common.pipeline.PipelineOptions; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import org.apache.flink.cdc.runtime.operators.transform.TransformRule; @@ -38,6 +39,7 @@ public class AsyncPostTransformFunctionBuilder { private final List transformRules = new ArrayList<>(); private String timezone; private DecimalPrecisionMode decimalPrecisionMode = DecimalPrecisionMode.UP_TO_19; + private TransformExpressionSemantics expressionSemantics = TransformExpressionSemantics.LEGACY; private final List>> udfFunctions = new ArrayList<>(); private final Map modelClients = new LinkedHashMap<>(); @@ -109,6 +111,12 @@ public AsyncPostTransformFunctionBuilder addUdfFunctions( return this; } + public AsyncPostTransformFunctionBuilder addExpressionSemantics( + TransformExpressionSemantics expressionSemantics) { + this.expressionSemantics = expressionSemantics; + return this; + } + public AsyncPostTransformFunctionBuilder addModelClients(Map clients) { this.modelClients.putAll(clients); return this; @@ -124,6 +132,7 @@ public AsyncPostTransformFunction build() { transformRules, timezone, decimalPrecisionMode, + expressionSemantics, udfFunctions, modelClients, asyncWorkerThreads); 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 340e7e61eb8..e8a04ef8947 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 @@ -22,6 +22,7 @@ import org.apache.flink.cdc.common.annotation.VisibleForTesting; import org.apache.flink.cdc.common.converter.JavaClassConverter; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.schema.Column; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import org.apache.flink.cdc.common.types.DataType; @@ -428,6 +429,9 @@ private static Java.Rvalue sqlBasicCallToJaninoRvalue( case EQUALS: return generateEqualsOperation(context, sqlBasicCall, atoms); case NOT_EQUALS: + if (usesFlinkSqlSemantics(context)) { + return generateFunctionOperation("sqlNotEquals", atoms); + } return generateUnaryOperation( context, "!", generateEqualsOperation(context, sqlBasicCall, atoms)); case IS_DISTINCT_FROM: @@ -454,6 +458,10 @@ private static Java.Rvalue sqlBasicCallToJaninoRvalue( case IN: case NOT_IN: case LIKE: + if (usesFlinkSqlSemantics(context)) { + return generateFlinkSqlPredicateOperation(sqlBasicCall, atoms); + } + return generateOtherFunctionOperation(context, sqlBasicCall, atoms); case SIMILAR: case POSITION: case CEIL: @@ -692,6 +700,10 @@ private static boolean isExpressionNullable(Context context, SqlNode sqlNode) { if (sqlNode instanceof SqlBasicCall) { return isBasicCallNullable(context, (SqlBasicCall) sqlNode); } + if (sqlNode instanceof SqlNodeList) { + return ((SqlNodeList) sqlNode) + .getList().stream().anyMatch(operand -> isExpressionNullable(context, operand)); + } return true; } @@ -742,6 +754,7 @@ private static boolean isBasicCallNullable(Context context, SqlBasicCall sqlBasi case IS_UNKNOWN: case IS_DISTINCT_FROM: case IS_NOT_DISTINCT_FROM: + return false; case EQUALS: case NOT_EQUALS: case LESS_THAN: @@ -751,7 +764,9 @@ private static boolean isBasicCallNullable(Context context, SqlBasicCall sqlBasi case BETWEEN: case IN: case NOT_IN: - return false; + return usesFlinkSqlSemantics(context) + && sqlBasicCall.getOperandList().stream() + .anyMatch(operand -> isExpressionNullable(context, operand)); case LIKE: case SIMILAR: return sqlBasicCall.getOperandList().stream() @@ -794,8 +809,8 @@ private static Java.Rvalue generateEqualsOperation( if (atoms.length != 2) { throw new ParseException("Unrecognized expression: " + sqlBasicCall.toString()); } - return new Java.MethodInvocation( - Location.NOWHERE, null, StringUtils.convertToCamelCase("VALUE_EQUALS"), atoms); + return generateFunctionOperation( + usesFlinkSqlSemantics(context) ? "sqlValueEquals" : "valueEquals", atoms); } private static Java.Rvalue generateCastOperation( @@ -859,7 +874,39 @@ private static Java.Rvalue generateCompareOperation( + sqlBasicCall.getKind().toString()); } return new Java.MethodInvocation( - Location.NOWHERE, null, StringUtils.convertToCamelCase(compareMethodName), atoms); + Location.NOWHERE, + null, + StringUtils.convertToCamelCase( + (usesFlinkSqlSemantics(context) ? "SQL_" : "") + compareMethodName), + atoms); + } + + private static Java.Rvalue generateFlinkSqlPredicateOperation( + SqlBasicCall sqlBasicCall, Java.Rvalue[] atoms) { + String operationName = sqlBasicCall.getOperator().getName().toUpperCase(); + switch (sqlBasicCall.getKind()) { + case BETWEEN: + return generateFunctionOperation( + operationName.startsWith("NOT") ? "sqlNotBetween" : "sqlBetween", atoms); + case IN: + return generateFunctionOperation("sqlIn", atoms); + case NOT_IN: + return generateFunctionOperation("sqlNotIn", atoms); + case LIKE: + if (atoms.length == 2) { + return generateFunctionOperation( + operationName.startsWith("NOT") ? "sqlNotLike" : "sqlLike", atoms); + } + return generateFunctionOperation( + StringUtils.convertToCamelCase(sqlBasicCall.getOperator().getName()), + atoms); + default: + throw new ParseException("Unsupported Flink SQL predicate: " + sqlBasicCall); + } + } + + private static boolean usesFlinkSqlSemantics(Context context) { + return TransformExpressionSemantics.FLINK_SQL.equals(context.expressionSemantics); } private static Java.Rvalue generateTimestampDiffOperation( @@ -1279,17 +1326,22 @@ public static class Context { // Maximum precision mode for DECIMAL type evaluation public final DecimalPrecisionMode decimalPrecisionMode; + // Semantics used to generate supported transform predicates + public final TransformExpressionSemantics expressionSemantics; + private Context( List columns, Map columnNameMap, List udfDescriptors, SupportedMetadataColumn[] supportedMetadataColumns, - DecimalPrecisionMode decimalPrecisionMode) { + DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics) { this.columns = columns; this.columnNameMap = columnNameMap; this.udfDescriptors = udfDescriptors; this.supportedMetadataColumns = supportedMetadataColumns; this.decimalPrecisionMode = decimalPrecisionMode; + this.expressionSemantics = expressionSemantics; } public static Context of( @@ -1302,7 +1354,8 @@ public static Context of( columnNameMap, udfDescriptors, supportedMetadataColumns, - DecimalPrecisionMode.UP_TO_19); + DecimalPrecisionMode.UP_TO_19, + TransformExpressionSemantics.LEGACY); } public static Context of( @@ -1311,12 +1364,44 @@ public static Context of( List udfDescriptors, SupportedMetadataColumn[] supportedMetadataColumns, DecimalPrecisionMode decimalPrecisionMode) { + return of( + columns, + columnNameMap, + udfDescriptors, + supportedMetadataColumns, + decimalPrecisionMode, + TransformExpressionSemantics.LEGACY); + } + + public static Context of( + List columns, + Map columnNameMap, + List udfDescriptors, + SupportedMetadataColumn[] supportedMetadataColumns, + TransformExpressionSemantics expressionSemantics) { + return of( + columns, + columnNameMap, + udfDescriptors, + supportedMetadataColumns, + DecimalPrecisionMode.UP_TO_19, + expressionSemantics); + } + + public static Context of( + List columns, + Map columnNameMap, + List udfDescriptors, + SupportedMetadataColumn[] supportedMetadataColumns, + DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics) { return new Context( columns, columnNameMap, udfDescriptors, supportedMetadataColumns, - decimalPrecisionMode); + decimalPrecisionMode, + expressionSemantics); } } } diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java index ebddd1de8a6..83fc2e048e6 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java @@ -19,6 +19,7 @@ import org.apache.flink.api.common.io.ParseException; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.schema.Column; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; import org.apache.flink.cdc.common.types.DataType; @@ -406,7 +407,8 @@ public static List generateProjectionColumns( columns, udfDescriptors, supportedMetadataColumns, - DecimalPrecisionMode.UP_TO_19); + DecimalPrecisionMode.UP_TO_19, + TransformExpressionSemantics.LEGACY); } public static List generateProjectionColumns( @@ -415,6 +417,37 @@ public static List generateProjectionColumns( List udfDescriptors, SupportedMetadataColumn[] supportedMetadataColumns, DecimalPrecisionMode decimalPrecisionMode) { + return generateProjectionColumns( + projectionExpression, + columns, + udfDescriptors, + supportedMetadataColumns, + decimalPrecisionMode, + TransformExpressionSemantics.LEGACY); + } + + public static List generateProjectionColumns( + String projectionExpression, + List columns, + List udfDescriptors, + SupportedMetadataColumn[] supportedMetadataColumns, + TransformExpressionSemantics expressionSemantics) { + return generateProjectionColumns( + projectionExpression, + columns, + udfDescriptors, + supportedMetadataColumns, + DecimalPrecisionMode.UP_TO_19, + expressionSemantics); + } + + public static List generateProjectionColumns( + String projectionExpression, + List columns, + List udfDescriptors, + SupportedMetadataColumn[] supportedMetadataColumns, + DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics) { if (isNullOrWhitespaceOnly(projectionExpression)) { return new ArrayList<>(); } @@ -504,7 +537,8 @@ public static List generateProjectionColumns( columnNameMap, udfDescriptors, supportedMetadataColumns, - decimalPrecisionMode), + decimalPrecisionMode, + expressionSemantics), exprNode), originalColumnNames, columnNameMap); @@ -743,7 +777,8 @@ public static String translateFilterExpressionToJaninoExpression( udfDescriptors, supportedMetadataColumns, columnNameMap, - DecimalPrecisionMode.UP_TO_19); + DecimalPrecisionMode.UP_TO_19, + TransformExpressionSemantics.LEGACY); } public static String translateFilterExpressionToJaninoExpression( @@ -753,6 +788,41 @@ public static String translateFilterExpressionToJaninoExpression( SupportedMetadataColumn[] supportedMetadataColumns, Map columnNameMap, DecimalPrecisionMode decimalPrecisionMode) { + return translateFilterExpressionToJaninoExpression( + filterExpression, + columns, + udfDescriptors, + supportedMetadataColumns, + columnNameMap, + decimalPrecisionMode, + TransformExpressionSemantics.LEGACY); + } + + public static String translateFilterExpressionToJaninoExpression( + String filterExpression, + List columns, + List udfDescriptors, + SupportedMetadataColumn[] supportedMetadataColumns, + Map columnNameMap, + TransformExpressionSemantics expressionSemantics) { + return translateFilterExpressionToJaninoExpression( + filterExpression, + columns, + udfDescriptors, + supportedMetadataColumns, + columnNameMap, + DecimalPrecisionMode.UP_TO_19, + expressionSemantics); + } + + public static String translateFilterExpressionToJaninoExpression( + String filterExpression, + List columns, + List udfDescriptors, + SupportedMetadataColumn[] supportedMetadataColumns, + Map columnNameMap, + DecimalPrecisionMode decimalPrecisionMode, + TransformExpressionSemantics expressionSemantics) { if (isNullOrWhitespaceOnly(filterExpression)) { return ""; } @@ -767,7 +837,8 @@ public static String translateFilterExpressionToJaninoExpression( columnNameMap, udfDescriptors, supportedMetadataColumns, - decimalPrecisionMode), + decimalPrecisionMode, + expressionSemantics), where); } diff --git a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java index a7a1c3f39c2..65847425bc5 100644 --- a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java +++ b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java @@ -34,6 +34,18 @@ void testLegacyNullComparisonBehavior() { assertThat(ComparisonFunctions.greaterThanOrEqual(1, null)).isFalse(); } + @Test + void testFlinkSqlNullComparisonBehavior() { + assertThat(ComparisonFunctions.sqlValueEquals(null, 1)).isNull(); + assertThat(ComparisonFunctions.sqlNotEquals(1, null)).isNull(); + assertThat(ComparisonFunctions.sqlLessThan(null, 1)).isNull(); + assertThat(ComparisonFunctions.sqlLessThanOrEqual(1, null)).isNull(); + assertThat(ComparisonFunctions.sqlGreaterThan(null, 1)).isNull(); + assertThat(ComparisonFunctions.sqlGreaterThanOrEqual(1, null)).isNull(); + assertThat(ComparisonFunctions.sqlValueEquals(1, 1)).isTrue(); + assertThat(ComparisonFunctions.sqlGreaterThan(2, 1)).isTrue(); + } + @Test void testThreeValuedLogicalOperators() { assertThat(LogicalFunctions.and(true, (Boolean) null)).isNull(); @@ -55,6 +67,14 @@ void testLegacyBetweenNullBehavior() { assertThat(ComparisonFunctions.notBetweenAsymmetric((Integer) null, 1, 3)).isTrue(); } + @Test + void testFlinkSqlBetweenNullBehavior() { + assertThat(ComparisonFunctions.sqlBetween(null, 1, 3)).isNull(); + assertThat(ComparisonFunctions.sqlBetween(2, null, 3)).isNull(); + assertThat(ComparisonFunctions.sqlBetween(4, null, 3)).isFalse(); + assertThat(ComparisonFunctions.sqlNotBetween(null, 1, 3)).isNull(); + } + @Test void testLegacyInNullBehavior() { assertThat(ComparisonFunctions.in(1, 1, null)).isTrue(); @@ -63,6 +83,15 @@ void testLegacyInNullBehavior() { assertThat(ComparisonFunctions.notIn(1, 2, null)).isTrue(); } + @Test + void testFlinkSqlInNullBehavior() { + assertThat(ComparisonFunctions.sqlIn(1, 1, null)).isTrue(); + assertThat(ComparisonFunctions.sqlIn(1, 2, null)).isNull(); + assertThat(ComparisonFunctions.sqlIn(null, 1, 2)).isNull(); + assertThat(ComparisonFunctions.sqlIn(1, 2, 3)).isFalse(); + assertThat(ComparisonFunctions.sqlNotIn(1, 2, null)).isNull(); + } + @Test void testDistinctFromNeverReturnsUnknown() { assertThat(ComparisonFunctions.isDistinctFrom(null, null)).isFalse(); 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 4994dc0a1df..baa58973a94 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 @@ -137,6 +137,22 @@ void testLegacyLikeUsesJavaRegex() { assertThat(StringFunctions.notLike("Alice", "A.*")).isFalse(); } + @Test + void testFlinkSqlLikeUsesSqlWildcardsAndMatchesWholeString() { + assertThat(StringFunctions.sqlLike("Alice", "A%")).isTrue(); + assertThat(StringFunctions.sqlLike("Alice", "A_i_e")).isTrue(); + assertThat(StringFunctions.sqlLike("xabcy", "abc")).isFalse(); + assertThat(StringFunctions.sqlLike("Alice", "A.*")).isFalse(); + assertThat(StringFunctions.sqlNotLike("Alice", "A%")).isFalse(); + } + + @Test + void testFlinkSqlLikeNullReturnsUnknown() { + assertThat(StringFunctions.sqlLike(null, "A%")).isNull(); + assertThat(StringFunctions.sqlLike("Alice", null)).isNull(); + assertThat(StringFunctions.sqlNotLike(null, "A%")).isNull(); + } + @Test void testLikeEscape() { assertThat(StringFunctions.like("A%", "A$%", "$")).isTrue(); 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 06db8f9f472..e37a81ed218 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 @@ -19,6 +19,7 @@ import org.apache.flink.api.common.io.ParseException; import org.apache.flink.cdc.common.pipeline.DecimalPrecisionMode; +import org.apache.flink.cdc.common.pipeline.TransformExpressionSemantics; import org.apache.flink.cdc.common.schema.Column; import org.apache.flink.cdc.common.schema.Schema; import org.apache.flink.cdc.common.source.SupportedMetadataColumn; @@ -748,6 +749,23 @@ void testStringFunctionArgumentValidation() { DataTypes.VARBINARY(65536)); } + @Test + void testTranslateFilterWithFlinkSqlSemantics() { + testFlinkSqlFilterExpression("id = 1", "sqlValueEquals(id, 1)"); + testFlinkSqlFilterExpression("id <> 1", "sqlNotEquals(id, 1)"); + testFlinkSqlFilterExpression("id < 1", "sqlLessThan(id, 1)"); + testFlinkSqlFilterExpression("id <= 1", "sqlLessThanOrEqual(id, 1)"); + testFlinkSqlFilterExpression("id > 1", "sqlGreaterThan(id, 1)"); + testFlinkSqlFilterExpression("id >= 1", "sqlGreaterThanOrEqual(id, 1)"); + testFlinkSqlFilterExpression("id between 1 and 3", "sqlBetween(id, 1, 3)"); + testFlinkSqlFilterExpression("id not between 1 and 3", "sqlNotBetween(id, 1, 3)"); + testFlinkSqlFilterExpression("id in (1, 2)", "sqlIn(id, 1, 2)"); + testFlinkSqlFilterExpression("id not in (1, 2)", "sqlNotIn(id, 1, 2)"); + testFlinkSqlFilterExpression("id like 'A%'", "sqlLike(id, \"A%\")"); + testFlinkSqlFilterExpression("id not like 'A%'", "sqlNotLike(id, \"A%\")"); + testFlinkSqlFilterExpression("id like 'A$%' escape '$'", "like(id, \"A$%\", \"$\")"); + } + @Test void testTranslateLogicalFilterToJaninoExpressionByNullability() { List columns = @@ -1427,6 +1445,18 @@ private void testFilterExpression(String expression, String expressionExpect) { Assertions.assertThat(janinoExpression).isEqualTo(expressionExpect); } + private void testFlinkSqlFilterExpression(String expression, String expressionExpect) { + String janinoExpression = + TransformParser.translateFilterExpressionToJaninoExpression( + expression, + DUMMY_COLUMNS, + Collections.emptyList(), + new SupportedMetadataColumn[0], + Collections.emptyMap(), + TransformExpressionSemantics.FLINK_SQL); + Assertions.assertThat(janinoExpression).isEqualTo(expressionExpect); + } + private void testFilterExpressionWithColumns( String expression, String expressionExpect, List columns) { String janinoExpression = From 361ef249f7db4c0eac6c2f157705ce6093766999 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=98=A5=E6=A0=96?= Date: Tue, 28 Jul 2026 16:50:28 +0800 Subject: [PATCH 2/3] Rename LEGACY to DEFAULT for transform expression semantics --- docs/content.zh/docs/core-concept/data-pipeline.md | 4 ++-- docs/content/docs/core-concept/data-pipeline.md | 4 ++-- .../cdc/cli/parser/YamlPipelineDefinitionParserTest.java | 2 +- .../apache/flink/cdc/common/pipeline/PipelineOptions.java | 4 ++-- .../cdc/common/pipeline/TransformExpressionSemantics.java | 4 ++-- .../flink/cdc/composer/specs/TransformSpecsITCase.java | 2 +- .../operators/transform/PostTransformOperatorBuilder.java | 2 +- .../operators/transform/TransformFilterProcessor.java | 2 +- .../async/AsyncPostTransformFunctionBuilder.java | 2 +- .../apache/flink/cdc/runtime/parser/JaninoCompiler.java | 4 ++-- .../apache/flink/cdc/runtime/parser/TransformParser.java | 8 ++++---- 11 files changed, 19 insertions(+), 19 deletions(-) diff --git a/docs/content.zh/docs/core-concept/data-pipeline.md b/docs/content.zh/docs/core-concept/data-pipeline.md index 26e484bf142..a54e0ad56b1 100644 --- a/docs/content.zh/docs/core-concept/data-pipeline.md +++ b/docs/content.zh/docs/core-concept/data-pipeline.md @@ -119,7 +119,7 @@ under the License. | `local-time-zone` | 作业级别的本地时区。 | optional | | `execution.runtime-mode` | pipeline 的运行模式,包含 STREAMING 和 BATCH,默认值是 STREAMING。 | optional | | `route-mode` | [路由]({{< ref "docs/core-concept/route" >}}#路由模式)规则的匹配模式。可选值:`ALL_MATCH`(默认,应用所有匹配的规则)或 `FIRST_MATCH`(只应用第一个匹配的规则)。 | optional | -| `transform.expression.semantics` | Transform 表达式的求值语义。可选值:`LEGACY`(默认,保持现有行为)或 `FLINK_SQL`(对支持的谓词使用与 Flink SQL 兼容的语义)。 | optional | +| `transform.expression.semantics` | Transform 表达式的求值语义。可选值:`DEFAULT`(默认,保持当前行为)或 `FLINK_SQL`(对支持的谓词使用与 Flink SQL 兼容的语义)。 | optional | | `schema.change.behavior` | 如何处理 [schema 变更]({{< ref "docs/core-concept/schema-evolution" >}})。可选值:[`exception`]({{< ref "docs/core-concept/schema-evolution" >}}#exception-mode)、[`evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#evolve-mode)、[`try_evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#tryevolve-mode)、[`lenient`]({{< ref "docs/core-concept/schema-evolution" >}}#lenient-mode)(默认值)或 [`ignore`]({{< ref "docs/core-concept/schema-evolution" >}}#ignore-mode)。 | optional | | `schema.operator.uid` | Schema 算子的唯一 ID。此 ID 用于算子间通信,必须在所有算子中保持唯一。**已废弃**:请使用 `operator.uid.prefix` 代替。 | optional | | `schema-operator.rpc-timeout` | SchemaOperator 等待下游 SchemaChangeEvent 应用完成的超时时间,默认值是 3 分钟。 | optional | @@ -141,7 +141,7 @@ under the License. `transform.expression.semantics` 配置项控制行级 Transform 过滤和投影表达式的求值语义: -* `LEGACY` 是默认值,用于保持现有 Pipeline 行为。comparison、`BETWEEN` 和 `IN` 在操作数为 `NULL` 时仍返回布尔值;两参数 `LIKE` 将 pattern 作为 Java 正则表达式,并使用子串匹配。 +* `DEFAULT` 是默认值,用于保持当前 Pipeline 行为。comparison、`BETWEEN` 和 `IN` 在操作数为 `NULL` 时仍返回布尔值;两参数 `LIKE` 将 pattern 作为 Java 正则表达式,并使用子串匹配。 * `FLINK_SQL` 将 `=`、`<>`、`<`、`<=`、`>`、`>=`、`BETWEEN` / `NOT BETWEEN`、`IN` / `NOT IN` 和两参数 `LIKE` / `NOT LIKE` 与 Flink SQL 对齐。这些谓词使用 SQL 三值逻辑,因此结果可能为 `UNKNOWN`(以 `NULL` 表示)。两参数 `LIKE` 使用 SQL `%` 和 `_` 通配符,并对整个字符串进行匹配。 配置示例: diff --git a/docs/content/docs/core-concept/data-pipeline.md b/docs/content/docs/core-concept/data-pipeline.md index acb75e0188f..e353ab0a509 100644 --- a/docs/content/docs/core-concept/data-pipeline.md +++ b/docs/content/docs/core-concept/data-pipeline.md @@ -121,7 +121,7 @@ Note that whilst the parameters are each individually optional, at least one of | `local-time-zone` | The local time zone defines current session time zone id. | optional | | `execution.runtime-mode` | The runtime mode of the pipeline includes STREAMING and BATCH, with the default value being STREAMING. | optional | | `route-mode` | The matching mode for [route]({{< ref "docs/core-concept/route" >}}#route-mode) rules. One of: `ALL_MATCH` (default, apply all matching rules) or `FIRST_MATCH` (apply only the first matching rule). | optional | -| `transform.expression.semantics` | The semantics used to evaluate transform expressions. One of: `LEGACY` (default, preserves existing behavior) or `FLINK_SQL` (uses Flink SQL-compatible semantics for supported predicates). | optional | +| `transform.expression.semantics` | The semantics used to evaluate transform expressions. One of: `DEFAULT` (default, preserves current behavior) or `FLINK_SQL` (uses Flink SQL-compatible semantics for supported predicates). | optional | | `schema.change.behavior` | How to handle [changes in schema]({{< ref "docs/core-concept/schema-evolution" >}}). One of: [`exception`]({{< ref "docs/core-concept/schema-evolution" >}}#exception-mode), [`evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#evolve-mode), [`try_evolve`]({{< ref "docs/core-concept/schema-evolution" >}}#tryevolve-mode), [`lenient`]({{< ref "docs/core-concept/schema-evolution" >}}#lenient-mode) (default) or [`ignore`]({{< ref "docs/core-concept/schema-evolution" >}}#ignore-mode). | optional | | `schema.operator.uid` | The unique ID for schema operator. This ID will be used for inter-operator communications and must be unique across operators. **Deprecated**: use `operator.uid.prefix` instead. | optional | | `schema-operator.rpc-timeout` | The timeout time for SchemaOperator to wait downstream SchemaChangeEvent applying finished, the default value is 3 minutes. | optional | @@ -143,7 +143,7 @@ NOTE: Whilst the above parameters are each individually optional, at least one o The `transform.expression.semantics` option controls the evaluation semantics of row-level transform filters and projections: -* `LEGACY` is the default and preserves existing pipeline behavior. Comparisons, `BETWEEN`, and `IN` return a boolean result when an operand is `NULL`. Two-argument `LIKE` treats its pattern as a Java regular expression and uses substring matching. +* `DEFAULT` is the default and preserves current pipeline behavior. Comparisons, `BETWEEN`, and `IN` return a boolean result when an operand is `NULL`. Two-argument `LIKE` treats its pattern as a Java regular expression and uses substring matching. * `FLINK_SQL` aligns `=`, `<>`, `<`, `<=`, `>`, `>=`, `BETWEEN` / `NOT BETWEEN`, `IN` / `NOT IN`, and two-argument `LIKE` / `NOT LIKE` with Flink SQL. These predicates use SQL three-valued logic, so a result may be `UNKNOWN` (represented as `NULL`). Two-argument `LIKE` uses SQL `%` and `_` wildcards and matches the entire string. For example: diff --git a/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java b/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java index 0d4842cdf3f..f0f44313014 100644 --- a/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java +++ b/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java @@ -125,7 +125,7 @@ void testTransformExpressionSemantics() throws Exception { .get(PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS)) .isEqualTo(TransformExpressionSemantics.FLINK_SQL); assertThat(new Configuration().get(PipelineOptions.PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS)) - .isEqualTo(TransformExpressionSemantics.LEGACY); + .isEqualTo(TransformExpressionSemantics.DEFAULT); } @Test diff --git a/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java index 95d8827ad89..848714e1c81 100644 --- a/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java +++ b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/PipelineOptions.java @@ -207,9 +207,9 @@ public class PipelineOptions { PIPELINE_TRANSFORM_EXPRESSION_SEMANTICS = ConfigOptions.key("transform.expression.semantics") .enumType(TransformExpressionSemantics.class) - .defaultValue(TransformExpressionSemantics.LEGACY) + .defaultValue(TransformExpressionSemantics.DEFAULT) .withDescription( - "Semantics used to evaluate transform expressions. LEGACY preserves historical behavior, while FLINK_SQL enables Flink SQL-compatible semantics for supported predicates."); + "Semantics used to evaluate transform expressions. DEFAULT preserves current behavior, while FLINK_SQL enables Flink SQL-compatible semantics for supported predicates."); private PipelineOptions() {} } diff --git a/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java index 5dabcea294f..3c2f8908e99 100644 --- a/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java +++ b/flink-cdc-common/src/main/java/org/apache/flink/cdc/common/pipeline/TransformExpressionSemantics.java @@ -22,8 +22,8 @@ /** Defines the semantics used to evaluate pipeline transform expressions. */ @PublicEvolving public enum TransformExpressionSemantics { - /** Preserves the historical pipeline transform expression behavior. */ - LEGACY, + /** Preserves the current pipeline transform expression behavior. */ + DEFAULT, /** Uses Flink SQL-compatible semantics for supported transform expressions. */ FLINK_SQL diff --git a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java index 973ffa71bf7..15e27315bd2 100644 --- a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java +++ b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/specs/TransformSpecsITCase.java @@ -416,7 +416,7 @@ static class TestSpec { public String timeZone = "UTC"; public DecimalPrecisionMode decimalPrecisionMode = DecimalPrecisionMode.UP_TO_19; public TransformExpressionSemantics expressionSemantics = - TransformExpressionSemantics.LEGACY; + TransformExpressionSemantics.DEFAULT; public List projectionRules = new ArrayList<>(); public @Nullable String filterRule; public @Nullable String primaryKey; diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java index 8653ad7fc67..36d8e39e62b 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorBuilder.java @@ -37,7 +37,7 @@ public class PostTransformOperatorBuilder { private final List transformRules = new ArrayList<>(); private String timezone; private DecimalPrecisionMode decimalPrecisionMode = DecimalPrecisionMode.UP_TO_19; - private TransformExpressionSemantics expressionSemantics = TransformExpressionSemantics.LEGACY; + private TransformExpressionSemantics expressionSemantics = TransformExpressionSemantics.DEFAULT; private final List>> udfFunctions = new ArrayList<>(); private final Map modelClients = new LinkedHashMap<>(); diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java index 79c1ae11a68..cb00e6a41d4 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java @@ -107,7 +107,7 @@ public static TransformFilterProcessor ofNoOp(DecimalPrecisionMode decimalPrecis null, null, null, - TransformExpressionSemantics.LEGACY); + TransformExpressionSemantics.DEFAULT); } public static TransformFilterProcessor of( diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java index 06fb48fc354..fce7c557445 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/async/AsyncPostTransformFunctionBuilder.java @@ -39,7 +39,7 @@ public class AsyncPostTransformFunctionBuilder { private final List transformRules = new ArrayList<>(); private String timezone; private DecimalPrecisionMode decimalPrecisionMode = DecimalPrecisionMode.UP_TO_19; - private TransformExpressionSemantics expressionSemantics = TransformExpressionSemantics.LEGACY; + private TransformExpressionSemantics expressionSemantics = TransformExpressionSemantics.DEFAULT; private final List>> udfFunctions = new ArrayList<>(); private final Map modelClients = new LinkedHashMap<>(); 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 e8a04ef8947..6d978cf30a7 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 @@ -1355,7 +1355,7 @@ public static Context of( udfDescriptors, supportedMetadataColumns, DecimalPrecisionMode.UP_TO_19, - TransformExpressionSemantics.LEGACY); + TransformExpressionSemantics.DEFAULT); } public static Context of( @@ -1370,7 +1370,7 @@ public static Context of( udfDescriptors, supportedMetadataColumns, decimalPrecisionMode, - TransformExpressionSemantics.LEGACY); + TransformExpressionSemantics.DEFAULT); } public static Context of( diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java index 83fc2e048e6..f8bd70b1cf7 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java @@ -408,7 +408,7 @@ public static List generateProjectionColumns( udfDescriptors, supportedMetadataColumns, DecimalPrecisionMode.UP_TO_19, - TransformExpressionSemantics.LEGACY); + TransformExpressionSemantics.DEFAULT); } public static List generateProjectionColumns( @@ -423,7 +423,7 @@ public static List generateProjectionColumns( udfDescriptors, supportedMetadataColumns, decimalPrecisionMode, - TransformExpressionSemantics.LEGACY); + TransformExpressionSemantics.DEFAULT); } public static List generateProjectionColumns( @@ -778,7 +778,7 @@ public static String translateFilterExpressionToJaninoExpression( supportedMetadataColumns, columnNameMap, DecimalPrecisionMode.UP_TO_19, - TransformExpressionSemantics.LEGACY); + TransformExpressionSemantics.DEFAULT); } public static String translateFilterExpressionToJaninoExpression( @@ -795,7 +795,7 @@ public static String translateFilterExpressionToJaninoExpression( supportedMetadataColumns, columnNameMap, decimalPrecisionMode, - TransformExpressionSemantics.LEGACY); + TransformExpressionSemantics.DEFAULT); } public static String translateFilterExpressionToJaninoExpression( From b62f47262ba0da7e23263ac2ebc943e46ddfd1fb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=98=A5=E6=A0=96?= Date: Wed, 5 Aug 2026 11:53:04 +0800 Subject: [PATCH 3/3] address review comment --- .../functions/impl/ComparisonFunctions.java | 30 +++++++++-- .../impl/ComparisonFunctionsTest.java | 53 +++++++++++++++++++ 2 files changed, 79 insertions(+), 4 deletions(-) diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java index a87ce16a798..f5256b98f12 100644 --- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java +++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctions.java @@ -31,6 +31,9 @@ public static Boolean sqlValueEquals(Object object1, Object object2) { if (object1 == null || object2 == null) { return null; } + if (object1 instanceof Number && object2 instanceof Number) { + return sqlNumericCompare((Number) object1, (Number) object2) == 0; + } return object1.equals(object2); } @@ -65,6 +68,25 @@ private static int universalCompares(Object lhs, Object rhs) { } } + private static int sqlCompares(Object lhs, Object rhs) { + if (lhs instanceof Number && rhs instanceof Number) { + return sqlNumericCompare((Number) lhs, (Number) rhs); + } + return universalCompares(lhs, rhs); + } + + private static int sqlNumericCompare(Number lhs, Number rhs) { + if (isNonFinite(lhs) || isNonFinite(rhs)) { + return Double.compare(lhs.doubleValue(), rhs.doubleValue()); + } + return new BigDecimal(lhs.toString()).compareTo(new BigDecimal(rhs.toString())); + } + + private static boolean isNonFinite(Number value) { + return (value instanceof Double && !Double.isFinite(value.doubleValue())) + || (value instanceof Float && !Float.isFinite(value.floatValue())); + } + public static boolean greaterThan(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return false; @@ -76,7 +98,7 @@ public static Boolean sqlGreaterThan(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return null; } - return universalCompares(lhs, rhs) > 0; + return sqlCompares(lhs, rhs) > 0; } public static boolean greaterThanOrEqual(Object lhs, Object rhs) { @@ -90,7 +112,7 @@ public static Boolean sqlGreaterThanOrEqual(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return null; } - return universalCompares(lhs, rhs) >= 0; + return sqlCompares(lhs, rhs) >= 0; } public static boolean lessThan(Object lhs, Object rhs) { @@ -104,7 +126,7 @@ public static Boolean sqlLessThan(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return null; } - return universalCompares(lhs, rhs) < 0; + return sqlCompares(lhs, rhs) < 0; } public static boolean lessThanOrEqual(Object lhs, Object rhs) { @@ -118,7 +140,7 @@ public static Boolean sqlLessThanOrEqual(Object lhs, Object rhs) { if (lhs == null || rhs == null) { return null; } - return universalCompares(lhs, rhs) <= 0; + return sqlCompares(lhs, rhs) <= 0; } public static Boolean sqlBetween(Object value, Object minValue, Object maxValue) { diff --git a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java index 65847425bc5..7b0bab2e8c7 100644 --- a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java +++ b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/functions/impl/ComparisonFunctionsTest.java @@ -19,6 +19,8 @@ import org.junit.jupiter.api.Test; +import java.math.BigDecimal; + import static org.assertj.core.api.Assertions.assertThat; /** Unit tests for {@link ComparisonFunctions}. */ @@ -46,6 +48,37 @@ void testFlinkSqlNullComparisonBehavior() { assertThat(ComparisonFunctions.sqlGreaterThan(2, 1)).isTrue(); } + @Test + void testDefaultNumericComparisonBehaviorRemainsUnchanged() { + assertThat(ComparisonFunctions.valueEquals(1L, 1)).isFalse(); + assertThat(ComparisonFunctions.valueEquals(new BigDecimal("1.0"), new BigDecimal("1.00"))) + .isFalse(); + assertThat( + ComparisonFunctions.greaterThan( + 9007199254740993L, new BigDecimal("9007199254740992"))) + .isFalse(); + } + + @Test + void testFlinkSqlNumericComparisonBehavior() { + assertThat(ComparisonFunctions.sqlValueEquals(1L, 1)).isTrue(); + assertThat( + ComparisonFunctions.sqlValueEquals( + new BigDecimal("1.0"), new BigDecimal("1.00"))) + .isTrue(); + assertThat(ComparisonFunctions.sqlNotEquals(1L, 2)).isTrue(); + assertThat( + ComparisonFunctions.sqlGreaterThan( + 9007199254740993L, new BigDecimal("9007199254740992"))) + .isTrue(); + assertThat( + ComparisonFunctions.sqlLessThan( + new BigDecimal("9007199254740992"), 9007199254740993L)) + .isTrue(); + assertThat(ComparisonFunctions.sqlGreaterThanOrEqual(new BigDecimal("1.00"), 1L)).isTrue(); + assertThat(ComparisonFunctions.sqlLessThanOrEqual(1, new BigDecimal("1.0"))).isTrue(); + } + @Test void testThreeValuedLogicalOperators() { assertThat(LogicalFunctions.and(true, (Boolean) null)).isNull(); @@ -75,6 +108,20 @@ void testFlinkSqlBetweenNullBehavior() { assertThat(ComparisonFunctions.sqlNotBetween(null, 1, 3)).isNull(); } + @Test + void testFlinkSqlBetweenUsesNumericComparison() { + assertThat( + ComparisonFunctions.sqlBetween( + 9007199254740993L, + BigDecimal.ZERO, + new BigDecimal("9007199254740992"))) + .isFalse(); + assertThat( + ComparisonFunctions.sqlBetween( + 1L, new BigDecimal("1.00"), new BigDecimal("1.0"))) + .isTrue(); + } + @Test void testLegacyInNullBehavior() { assertThat(ComparisonFunctions.in(1, 1, null)).isTrue(); @@ -92,6 +139,12 @@ void testFlinkSqlInNullBehavior() { assertThat(ComparisonFunctions.sqlNotIn(1, 2, null)).isNull(); } + @Test + void testFlinkSqlInUsesNumericComparison() { + assertThat(ComparisonFunctions.sqlIn(1L, 2, new BigDecimal("1.00"))).isTrue(); + assertThat(ComparisonFunctions.sqlNotIn(new BigDecimal("1.0"), 1L, null)).isFalse(); + } + @Test void testDistinctFromNeverReturnsUnknown() { assertThat(ComparisonFunctions.isDistinctFrom(null, null)).isFalse();