From fa3b6d002d3cbbe7298b9bc52522cfd12d7d2fe2 Mon Sep 17 00:00:00 2001 From: Au_Miner Date: Thu, 30 Jul 2026 17:55:12 +0800 Subject: [PATCH 1/2] [FLINK-40304][table] Fix PTF partition columns for multiple table arguments AI-Contributed/Feature: 0/2 AI-Contributed/UT: 0/48 --- .../StreamPhysicalProcessTableFunction.java | 2 +- .../stream/sql/ProcessTableFunctionTest.java | 19 ++++++++++++ .../stream/sql/ProcessTableFunctionTest.xml | 29 +++++++++++++++++++ 3 files changed, 49 insertions(+), 1 deletion(-) diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalProcessTableFunction.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalProcessTableFunction.java index a337c25086089d..abc41fc106e569 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalProcessTableFunction.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalProcessTableFunction.java @@ -515,7 +515,7 @@ public static Set toPartitionColumns(RexCall call) { // f(t1 PARTITION BY (k1, k2), t2 PARTITION BY (k3, k4)) // -> [k1, k2, k3, k4, function out...] final List partitionColumns = - IntStream.range(pos, partitionKeyCount) + IntStream.range(pos, pos + partitionKeyCount) .boxed() .collect(Collectors.toList()); pos += partitionKeyCount; diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java index 7500bef9c98ba6..490adb04629f33 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java @@ -46,6 +46,7 @@ import org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.SetSemanticTableRetractArgFunction; import org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.TypedRowSemanticTableFunction; import org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.TypedSetSemanticTableFunction; +import org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.UpdatingJoinFunction; import org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.UpdatingRetractRowSemanticFunction; import org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.UpdatingUpsertFunction; import org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.User; @@ -91,6 +92,9 @@ void setup() { util.tableEnv() .executeSql( "CREATE VIEW t AS SELECT * FROM (VALUES ('Bob', 12), ('Alice', 42)) AS T(name, score)"); + util.tableEnv() + .executeSql( + "CREATE VIEW t_city AS SELECT * FROM (VALUES ('Bob', 'London'), ('Alice', 'Berlin')) AS T(name, city)"); util.tableEnv() .executeSql("CREATE VIEW t_name_diff AS SELECT 'Bob' AS name, 12 AS different"); util.tableEnv() @@ -111,6 +115,10 @@ void setup() { util.tableEnv() .executeSql( "CREATE TABLE t_keyed_sink (`name` STRING, `out` STRING) WITH ('connector' = 'blackhole')"); + util.tableEnv() + .executeSql( + "CREATE TABLE t_key_only_delete_sink (`name` STRING, `name0` STRING PRIMARY KEY NOT ENFORCED, `out` STRING) " + + "WITH ('connector' = 'values', 'sink-insert-only' = 'false', 'sink.supports-delete-by-key' = 'true')"); util.tableEnv() .executeSql( "CREATE TABLE t_full_delete_sink (`name` STRING PRIMARY KEY NOT ENFORCED, `name0` STRING, `count` BIGINT, `mode` STRING) " @@ -138,6 +146,17 @@ void testFunctionWithMultipleTableArgs() { util.verifyRelPlan("SELECT * FROM v"); } + @Test + void testUpsertKeyWithMultipleTableArgs() { + util.addTemporarySystemFunction("f", UpdatingJoinFunction.class); + util.tableEnv() + .executeSql( + "CREATE VIEW v AS SELECT * FROM f(" + + "scoreTable => TABLE t PARTITION BY name," + + "cityTable => TABLE t_city PARTITION BY name)"); + util.verifyRelPlanInsert("INSERT INTO t_key_only_delete_sink SELECT * FROM v"); + } + @Test void testScalarArgsWithUid() { util.addTemporarySystemFunction("f", ScalarArgsFunction.class); diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml index 7cd0f83548725b..b4794a6992a648 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml @@ -466,6 +466,35 @@ LogicalProject(out=[$0]) + + + + + + + + + + + From 26e3d840705272f2ac3ab3e01cdb61a442d192c3 Mon Sep 17 00:00:00 2001 From: Au_Miner <358671982@qq.com> Date: Tue, 4 Aug 2026 14:55:27 +0800 Subject: [PATCH 2/2] trigger ci AI-Contributed/Feature: 0/0 AI-Contributed/UT: 0/6 --- .../planner/plan/stream/sql/ProcessTableFunctionTest.java | 2 +- .../planner/plan/stream/sql/ProcessTableFunctionTest.xml | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java index 490adb04629f33..e45963ae78b838 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java @@ -94,7 +94,7 @@ void setup() { "CREATE VIEW t AS SELECT * FROM (VALUES ('Bob', 12), ('Alice', 42)) AS T(name, score)"); util.tableEnv() .executeSql( - "CREATE VIEW t_city AS SELECT * FROM (VALUES ('Bob', 'London'), ('Alice', 'Berlin')) AS T(name, city)"); + "CREATE VIEW t_city AS SELECT * FROM (VALUES ('Bob', 'Tokyo'), ('Alice', 'Berlin')) AS T(name, city)"); util.tableEnv() .executeSql("CREATE VIEW t_name_diff AS SELECT 'Bob' AS name, 12 AS different"); util.tableEnv() diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml index b4794a6992a648..78a39a25a758c5 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.xml @@ -484,7 +484,7 @@ LogicalSink(table=[default_catalog.default_database.t_key_only_delete_sink], fie : +- LogicalValues(tuples=[[{ _UTF-16LE'Bob', 12 }, { _UTF-16LE'Alice', 42 }]]) +- LogicalProject(name=[$0], city=[$1]) +- LogicalProject(name=[$0], city=[$1]) - +- LogicalValues(tuples=[[{ _UTF-16LE'Bob', _UTF-16LE'London' }, { _UTF-16LE'Alice', _UTF-16LE'Berlin' }]]) + +- LogicalValues(tuples=[[{ _UTF-16LE'Bob', _UTF-16LE'Tokyo' }, { _UTF-16LE'Alice', _UTF-16LE'Berlin' }]]) ]]> @@ -494,7 +494,7 @@ Sink(table=[default_catalog.default_database.t_key_only_delete_sink], fields=[na :- Exchange(distribution=[hash[name]]) : +- Values(tuples=[[{ _UTF-16LE'Bob', 12 }, { _UTF-16LE'Alice', 42 }]]) +- Exchange(distribution=[hash[name]]) - +- Values(tuples=[[{ _UTF-16LE'Bob', _UTF-16LE'London' }, { _UTF-16LE'Alice', _UTF-16LE'Berlin' }]]) + +- Values(tuples=[[{ _UTF-16LE'Bob', _UTF-16LE'Tokyo' }, { _UTF-16LE'Alice', _UTF-16LE'Berlin' }]]) ]]>