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 a337c25086089..abc41fc106e56 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 7500bef9c98ba..e45963ae78b83 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', '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() @@ -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 7cd0f83548725..78a39a25a758c 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]) + + + + + + + + + + +