Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -515,7 +515,7 @@ public static Set<ImmutableBitSet> toPartitionColumns(RexCall call) {
// f(t1 PARTITION BY (k1, k2), t2 PARTITION BY (k3, k4))
// -> [k1, k2, k3, k4, function out...]
final List<Integer> partitionColumns =
IntStream.range(pos, partitionKeyCount)
IntStream.range(pos, pos + partitionKeyCount)
.boxed()
.collect(Collectors.toList());
pos += partitionKeyCount;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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()
Expand All @@ -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) "
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -466,6 +466,35 @@ LogicalProject(out=[$0])
<![CDATA[
ProcessTableFunction(invocation=[f(1, true, DEFAULT(), DEFAULT())], uid=[null], select=[out], rowType=[RecordType(VARCHAR(2147483647) out)])
+- Values(tuples=[[{ }]])
]]>
</Resource>
</TestCase>
<TestCase name="testUpsertKeyWithMultipleTableArgs">
<Resource name="sql">
<![CDATA[INSERT INTO t_key_only_delete_sink SELECT * FROM v]]>
</Resource>
<Resource name="ast">
<![CDATA[
LogicalSink(table=[default_catalog.default_database.t_key_only_delete_sink], fields=[name, name0, out])
+- LogicalProject(name=[$0], name0=[$1], out=[$2])
+- LogicalProject(name=[$0], name0=[$1], out=[$2])
+- LogicalTableFunctionScan(invocation=[f(TABLE(#0) PARTITION BY($0), TABLE(#1) PARTITION BY($0), DEFAULT(), DEFAULT())], rowType=[RecordType(VARCHAR(5) name, VARCHAR(5) name0, VARCHAR(2147483647) out)])
:- LogicalProject(name=[$0], score=[$1])
: +- LogicalProject(name=[$0], score=[$1])
: +- 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'Tokyo' }, { _UTF-16LE'Alice', _UTF-16LE'Berlin' }]])
]]>
</Resource>
<Resource name="optimized rel plan">
<![CDATA[
Sink(table=[default_catalog.default_database.t_key_only_delete_sink], fields=[name, name0, out])
+- ProcessTableFunction(invocation=[f(TABLE(#0) PARTITION BY($0), TABLE(#1) PARTITION BY($0), DEFAULT(), DEFAULT())], uid=[f], select=[name,name0,out], rowType=[RecordType(VARCHAR(5) name, VARCHAR(5) name0, VARCHAR(2147483647) out)])
:- 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'Tokyo' }, { _UTF-16LE'Alice', _UTF-16LE'Berlin' }]])
]]>
</Resource>
</TestCase>
Expand Down