From 3b6473a20c37df15cf7d3c9472815a9c36160f9a Mon Sep 17 00:00:00 2001 From: Au_Miner <358671982@qq.com> Date: Thu, 30 Jul 2026 17:35:00 +0800 Subject: [PATCH] [FLINK-40301][table] Merge Calc nodes before ChangelogNormalize projection pushdown --- .../plan/rules/FlinkStreamRuleSets.scala | 1 + ...nerChangelogNormalizeTransposeRuleTest.xml | 32 ++++++++++ ...rChangelogNormalizeTransposeRuleTest.scala | 61 +++++++++++-------- 3 files changed, 69 insertions(+), 25 deletions(-) diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/FlinkStreamRuleSets.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/FlinkStreamRuleSets.scala index de637d07a244b..6beb889c17a19 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/FlinkStreamRuleSets.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/FlinkStreamRuleSets.scala @@ -540,6 +540,7 @@ object FlinkStreamRuleSets { val CHANGELOG_NORMALIZE_TRANSPOSE_RULES: RuleSet = RuleSets.ofList( WatermarkAssignerChangelogNormalizeTransposeRule.WITH_CALC, WatermarkAssignerChangelogNormalizeTransposeRule.WITHOUT_CALC, + FlinkCalcMergeRule.STREAM_PHYSICAL_INSTANCE, // reduce state size in ChangelogNormalize PushCalcPastChangelogNormalizeRule.INSTANCE ) diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.xml index 7a9052875746f..f34a5c38b3116 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.xml @@ -122,6 +122,38 @@ Calc(select=[a, b, f], where=[f], changelogMode=[I]) +- WatermarkAssigner(rowtime=[ingestion_time], watermark=[ingestion_time], changelogMode=[I,UA,D]) +- Calc(select=[CAST(ingestion_time AS TIMESTAMP(3) *ROWTIME*) AS ingestion_time, a, f], changelogMode=[I,UA,D]) +- TableSourceScan(table=[[default_catalog, default_database, t2, project=[a, f], metadata=[ts]]], fields=[a, f, ingestion_time], changelogMode=[I,UA,D]) +]]> + + + + + + + + + + + diff --git a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.scala b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.scala index 9d37d60bcd5f0..ce5630b06ee92 100644 --- a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.scala +++ b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/rules/physical/stream/WatermarkAssignerChangelogNormalizeTransposeRuleTest.scala @@ -75,6 +75,31 @@ class WatermarkAssignerChangelogNormalizeTransposeRuleTest extends TableTestBase | 'enable-watermark-push-down' = 'true' |) |""".stripMargin) + util.addTable(""" + |CREATE TABLE t1 ( + | ingestion_time TIMESTAMP(3) METADATA FROM 'ts', + | a VARCHAR NOT NULL, + | b VARCHAR NOT NULL, + | WATERMARK FOR ingestion_time AS ingestion_time + |) WITH ( + | 'connector' = 'values', + | 'readable-metadata' = 'ts:TIMESTAMP(3)' + |) + """.stripMargin) + util.addTable(""" + |CREATE TABLE t2 ( + | k VARBINARY, + | ingestion_time TIMESTAMP(3) METADATA FROM 'ts', + | a VARCHAR NOT NULL, + | f BOOLEAN NOT NULL, + | WATERMARK FOR `ingestion_time` AS `ingestion_time`, + | PRIMARY KEY (`a`) NOT ENFORCED + |) WITH ( + | 'connector' = 'values', + | 'readable-metadata' = 'ts:TIMESTAMP(3)', + | 'changelog-mode' = 'I,UA,D' + |) + """.stripMargin) } // ---------------------------------------------------------------------------------------- @@ -169,31 +194,6 @@ class WatermarkAssignerChangelogNormalizeTransposeRuleTest extends TableTestBase @Test def testPushdownCalcNotAffectChangelogNormalizeKey(): Unit = { - util.addTable(""" - |CREATE TABLE t1 ( - | ingestion_time TIMESTAMP(3) METADATA FROM 'ts', - | a VARCHAR NOT NULL, - | b VARCHAR NOT NULL, - | WATERMARK FOR ingestion_time AS ingestion_time - |) WITH ( - | 'connector' = 'values', - | 'readable-metadata' = 'ts:TIMESTAMP(3)' - |) - """.stripMargin) - util.addTable(""" - |CREATE TABLE t2 ( - | k VARBINARY, - | ingestion_time TIMESTAMP(3) METADATA FROM 'ts', - | a VARCHAR NOT NULL, - | f BOOLEAN NOT NULL, - | WATERMARK FOR `ingestion_time` AS `ingestion_time`, - | PRIMARY KEY (`a`) NOT ENFORCED - |) WITH ( - | 'connector' = 'values', - | 'readable-metadata' = 'ts:TIMESTAMP(3)', - | 'changelog-mode' = 'I,UA,D' - |) - """.stripMargin) // After FLINK-28988 applied, the filter will not be pushed down into left input of join and get // a more optimal plan (upsert mode without ChangelogNormalize). val sql = @@ -204,4 +204,15 @@ class WatermarkAssignerChangelogNormalizeTransposeRuleTest extends TableTestBase |""".stripMargin util.verifyRelPlan(sql, ExplainDetail.CHANGELOG_MODE) } + + @Test + def testPushdownCalcNotAffectChangelogNormalizeKey2(): Unit = { + val sql = + """ + |SELECT f, count(*) + |FROM (select ingestion_time, f, a from t2) t2 + |group by f + |""".stripMargin + util.verifyRelPlan(sql, ExplainDetail.CHANGELOG_MODE) + } }