diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala index 1c91e960a8a180..843ab885962ac9 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala @@ -467,8 +467,21 @@ class FlinkChangelogModeInferenceProgram extends FlinkOptimizeProgram[StreamOpti case normalize: StreamPhysicalChangelogNormalize => // changelog normalize support update&delete input val children = visitChildren(normalize, ModifyKindSetTrait.ALL_CHANGES) - // changelog normalize will output all changes - val providedTrait = ModifyKindSetTrait.ALL_CHANGES + // A filter can turn an update into a delete when the updated row stops matching. + val inputModifyKindSet = getModifyKindSet(children.head) + val providedTrait = + if ( + normalize.filterCondition == null && !inputModifyKindSet.contains(ModifyKind.DELETE) + ) { + new ModifyKindSetTrait( + ModifyKindSet + .newBuilder() + .addContainedKind(ModifyKind.INSERT) + .addContainedKind(ModifyKind.UPDATE) + .build()) + } else { + ModifyKindSetTrait.ALL_CHANGES + } createNewNode(normalize, children, providedTrait, requiredTrait, requester) case ts: StreamPhysicalTableSourceScan => diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.xml index 1e453bc0c9e18e..45144823d72902 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.xml @@ -16,6 +16,47 @@ See the License for the specific language governing permissions and limitations under the License. --> + + + + + + + + + + + + + + 'x']]> + + + ($1, _UTF-16LE'x')]) + +- LogicalTableScan(table=[[default_catalog, default_database, upsertSource]]) +]]> + + + (payload, 'x')], changelogMode=[I,UB,UA,D]) + +- Exchange(distribution=[hash[id]], changelogMode=[I,UA]) + +- TableSourceScan(table=[[default_catalog, default_database, upsertSource, filter=[]]], fields=[id, payload], changelogMode=[I,UA]) +]]> + + diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableScanTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableScanTest.xml index 9a1909432cdaff..609856f5db8f56 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableScanTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableScanTest.xml @@ -391,12 +391,12 @@ LogicalProject(currency_name=[$2], amount=[$0], rate=[$5], EXPR$3=[*($0, $5)]) diff --git a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.scala b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.scala index 787cbf39116c6f..6f2fd0bc84390c 100644 --- a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.scala +++ b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.scala @@ -88,6 +88,26 @@ class ChangelogModeInferenceTest extends TableTestBase { | 'changelog-mode' = 'I,UA,UB,D' |) """.stripMargin) + + util.addTable(""" + |CREATE TABLE upsertSource ( + | id INT PRIMARY KEY NOT ENFORCED, + | payload STRING + |) WITH ( + | 'connector' = 'values', + | 'changelog-mode' = 'I,UA' + |) + |""".stripMargin) + + util.addTable(""" + |CREATE TABLE changelogSink ( + | id INT, + | payload STRING + |) WITH ( + | 'connector' = 'values', + | 'sink-insert-only' = 'false' + |) + |""".stripMargin) } @Test @@ -583,4 +603,18 @@ class ChangelogModeInferenceTest extends TableTestBase { "INSERT INTO keyless_upsert_sink_no_key SELECT rate FROM DeduplicatedView", ExplainDetail.CHANGELOG_MODE) } + + @Test + def testChangelogNormalizeDoesNotInferDeleteWithoutFilterOrInputDelete(): Unit = { + util.verifyRelPlanInsert( + "INSERT INTO changelogSink SELECT * FROM upsertSource", + ExplainDetail.CHANGELOG_MODE) + } + + @Test + def testChangelogNormalizeInfersDeleteWithFilterWithoutInputDelete(): Unit = { + util.verifyRelPlanInsert( + "INSERT INTO changelogSink SELECT * FROM upsertSource WHERE payload <> 'x'", + ExplainDetail.CHANGELOG_MODE) + } } diff --git a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DeltaJoinTest.scala b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DeltaJoinTest.scala index 474146437519a6..d477131a2e76dc 100644 --- a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DeltaJoinTest.scala +++ b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DeltaJoinTest.scala @@ -852,8 +852,6 @@ class DeltaJoinTest extends TableTestBase { @Test def testPKContainsJoinKeyAndSourceNoUBAndD(): Unit = { - // FLINK-38489 Currently, ChangelogNormalize will always generate changelog mode with D, - // and Join with D cannot be optimized into Delta Join replaceTable( "no_delete_src1", "no_delete_and_update_before_src1",