This is an automated email from the ASF dual-hosted git repository.
gustavodemorais pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/master by this push:
new 2245026dff5 [FLINK-40447][table] Fall back to retract when input
upsert key is empty or mixed
2245026dff5 is described below
commit 2245026dff55e1348356ff472ee42c716710bbba
Author: Gustavo de Morais <[email protected]>
AuthorDate: Fri Aug 21 16:23:59 2026 +0200
[FLINK-40447][table] Fall back to retract when input upsert key is empty or
mixed
This closes #29000.
---
.../stream/StreamPhysicalProcessTableFunction.java | 2 +-
.../FlinkChangelogModeInferenceProgram.scala | 11 ++-
.../physical/stream/ChangelogModeInferenceTest.xml | 82 ++++++++++++++++
.../plan/stream/sql/DagOptimizationTest.xml | 58 ++++++------
.../stream/ChangelogModeInferenceTest.scala | 104 +++++++++++++++++++++
.../plan/stream/sql/DagOptimizationTest.scala | 28 +++---
6 files changed, 239 insertions(+), 46 deletions(-)
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 a337c250860..abc41fc106e 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 class StreamPhysicalProcessTableFunction extends
AbstractRelNode
// 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;
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 27719c67cdd..8d22466c3ee 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
@@ -1131,16 +1131,19 @@ class FlinkChangelogModeInferenceProgram extends
FlinkOptimizeProgram[StreamOpti
* when the upsert key has columns outside sink pk. This differs from
batch job's unique key
* inference.
*
- * <p>A sink without a primary key is satisfied whenever the input carries
any upsert key.
+ * <p>A sink without a primary key is satisfied whenever the input carries
a real (non-empty)
+ * upsert key; an empty candidate ("at most one row") never counts, even
alongside a real one.
*/
private def canUpsertKeysWithImmutableColsSatisfyPk(sink:
StreamPhysicalSink): Boolean = {
val sinkDefinedPks =
sink.contextResolvedTable.getResolvedSchema.getPrimaryKeyIndexes
val fmq =
FlinkRelMetadataQuery.reuseOrCreate(sink.getCluster.getMetadataQuery)
val changeLogUpsertKeys = fmq.getUpsertKeys(sink.getInput)
if (sinkDefinedPks.isEmpty) {
- // A keyless sink cannot apply UPDATE_AFTER in place, so it can only
accept upsert when the
- // input itself carries an upsert key; otherwise fall back to
beforeAndAfter.
- return changeLogUpsertKeys != null && !changeLogUpsertKeys.isEmpty
+ // A keyless sink can only stay upsert when the input has a real,
column-based upsert
+ // key. An empty candidate means "at most one row" (e.g. a global
aggregate), not columns
+ // to match on - UpsertKeyUtil.getSmallestKey would otherwise prefer
it over a real one.
+ return changeLogUpsertKeys != null && changeLogUpsertKeys.nonEmpty &&
+ !changeLogUpsertKeys.exists(_.isEmpty)
}
val sinkPks = ImmutableBitSet.of(sinkDefinedPks: _*)
// if upsert key is null, pk cannot be satisfied, should fall back to
beforeAndAfter
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 1e453bc0c9e..e0bd6a0a2fd 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
@@ -172,6 +172,27 @@ GroupAggregate(groupBy=[cnt], select=[cnt,
COUNT_RETRACT(cnt) AS frequency], cha
: +- LegacyTableSourceScan(table=[[default_catalog,
default_database, MyTable, source: [CollectionTableSource(word, number)]]],
fields=[word, number], changelogMode=[I])
+- Calc(select=[CAST(cnt AS BIGINT) AS cnt], changelogMode=[I])
+- LegacyTableSourceScan(table=[[default_catalog, default_database,
MyTable2, source: [CollectionTableSource(word, cnt)]]], fields=[word, cnt],
changelogMode=[I])
+]]>
+ </Resource>
+ </TestCase>
+ <TestCase name="testKeylessUpsertSinkFallsBackToRetractOnGlobalAggregate">
+ <Resource name="sql">
+ <![CDATA[INSERT INTO keyless_upsert_sink_count SELECT COUNT(*) FROM
MyTable]]>
+ </Resource>
+ <Resource name="ast">
+ <![CDATA[
+LogicalSink(table=[default_catalog.default_database.keyless_upsert_sink_count],
fields=[EXPR$0])
++- LogicalAggregate(group=[{}], EXPR$0=[COUNT()])
+ +- LogicalTableScan(table=[[default_catalog, default_database, MyTable,
source: [CollectionTableSource(word, number)]]])
+]]>
+ </Resource>
+ <Resource name="optimized rel plan">
+ <![CDATA[
+Sink(table=[default_catalog.default_database.keyless_upsert_sink_count],
fields=[EXPR$0], changelogMode=[NONE])
++- GroupAggregate(select=[COUNT(*) AS EXPR$0], changelogMode=[I,UB,UA])
+ +- Exchange(distribution=[single], changelogMode=[I])
+ +- Calc(select=[0 AS $f0], changelogMode=[I])
+ +- LegacyTableSourceScan(table=[[default_catalog, default_database,
MyTable, source: [CollectionTableSource(word, number)]]], fields=[word,
number], changelogMode=[I])
]]>
</Resource>
</TestCase>
@@ -224,6 +245,67 @@
Sink(table=[default_catalog.default_database.keyless_upsert_sink], fields=[curre
+- Exchange(distribution=[hash[currency]], changelogMode=[I])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime],
changelogMode=[I])
+- LegacyTableSourceScan(table=[[default_catalog,
default_database, ratesHistory, source: [CollectionTableSource(currency, rate,
rowtime)]]], fields=[currency, rate, rowtime], changelogMode=[I])
+]]>
+ </Resource>
+ </TestCase>
+ <TestCase
name="testKeylessUpsertSinkWithLookupJoinOnGlobalAggregateProbeSide">
+ <Resource name="sql">
+ <![CDATA[
+INSERT INTO lookup_sink
+SELECT dim.id, dim.v
+FROM (SELECT COUNT(*) AS cnt, PROCTIME() AS pt FROM MyTable) g
+JOIN LookupDim FOR SYSTEM_TIME AS OF g.pt AS dim
+ON g.cnt = dim.id
+ ]]>
+ </Resource>
+ <Resource name="ast">
+ <![CDATA[
+LogicalSink(table=[default_catalog.default_database.lookup_sink], fields=[id,
v])
++- LogicalProject(id=[$2], v=[$3])
+ +- LogicalCorrelate(correlation=[$cor0], joinType=[inner],
requiredColumns=[{0, 1}])
+ :- LogicalProject(cnt=[$0], pt=[PROCTIME()])
+ : +- LogicalAggregate(group=[{}], cnt=[COUNT()])
+ : +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable, source: [CollectionTableSource(word, number)]]])
+ +- LogicalFilter(condition=[=($cor0.cnt, $0)])
+ +- LogicalSnapshot(period=[$cor0.pt])
+ +- LogicalTableScan(table=[[default_catalog, default_database,
LookupDim]])
+]]>
+ </Resource>
+ <Resource name="optimized rel plan">
+ <![CDATA[
+Sink(table=[default_catalog.default_database.lookup_sink], fields=[id, v],
changelogMode=[NONE])
++- Calc(select=[id, v], changelogMode=[I,UB,UA])
+ +- LookupJoin(table=[default_catalog.default_database.LookupDim],
joinType=[InnerJoin], lookup=[id=cnt], select=[cnt, id, v], upsertKey=[[]],
changelogMode=[I,UB,UA])
+ +- GroupAggregate(select=[COUNT(*) AS cnt], changelogMode=[I,UB,UA])
+ +- Exchange(distribution=[single], changelogMode=[I])
+ +- Calc(select=[0 AS $f0], changelogMode=[I])
+ +- LegacyTableSourceScan(table=[[default_catalog,
default_database, MyTable, source: [CollectionTableSource(word, number)]]],
fields=[word, number], changelogMode=[I])
+]]>
+ </Resource>
+ </TestCase>
+ <TestCase
name="testKeylessUpsertSinkWithMultiTableArgUpsertPtfComputesDistinctKeyPerArg">
+ <Resource name="sql">
+ <![CDATA[INSERT INTO keyless_upsert_sink_ptf_probe SELECT `name`, name0,
`out` FROM f(scoreTable => TABLE scores_ptf_probe PARTITION BY name, cityTable
=> TABLE city_ptf_probe PARTITION BY name)]]>
+ </Resource>
+ <Resource name="ast">
+ <![CDATA[
+LogicalSink(table=[default_catalog.default_database.keyless_upsert_sink_ptf_probe],
fields=[name, name0, out])
++- 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(2147483647) name, VARCHAR(2147483647) name0,
VARCHAR(2147483647) out)])
+ :- LogicalProject(name=[$0], score=[$1])
+ : +- LogicalTableScan(table=[[default_catalog, default_database,
scores_ptf_probe]])
+ +- LogicalProject(name=[$0], city=[$1])
+ +- LogicalTableScan(table=[[default_catalog, default_database,
city_ptf_probe]])
+]]>
+ </Resource>
+ <Resource name="optimized rel plan">
+ <![CDATA[
+Sink(table=[default_catalog.default_database.keyless_upsert_sink_ptf_probe],
fields=[name, name0, out], changelogMode=[NONE])
++- ProcessTableFunction(invocation=[f(TABLE(#0) PARTITION BY($0), TABLE(#1)
PARTITION BY($0), DEFAULT(), DEFAULT())], uid=[f], select=[name,name0,out],
rowType=[RecordType(VARCHAR(2147483647) name, VARCHAR(2147483647) name0,
VARCHAR(2147483647) out)], changelogMode=[I,UA,PD])
+ :- Exchange(distribution=[hash[name]], changelogMode=[I,UA,PD])
+ : +- TableSourceScan(table=[[default_catalog, default_database,
scores_ptf_probe]], fields=[name, score], changelogMode=[I,UA,PD])
+ +- Exchange(distribution=[hash[name]], changelogMode=[I,UA,PD])
+ +- TableSourceScan(table=[[default_catalog, default_database,
city_ptf_probe]], fields=[name, city], changelogMode=[I,UA,PD])
]]>
</Resource>
</TestCase>
diff --git
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.xml
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.xml
index 67d54af4a83..b03570f9415 100644
---
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.xml
+++
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.xml
@@ -353,9 +353,9 @@ Sink(table=[default_catalog.default_database.retractSink2],
fields=[total_min],
<TestCase name="testMultiSinksSplitOnUnion1">
<Resource name="ast">
<![CDATA[
-LogicalSink(table=[default_catalog.default_database.upsertSink],
fields=[total_sum])
-+- LogicalAggregate(group=[{}], total_sum=[SUM($0)])
- +- LogicalProject(a=[$0])
+LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[c,
total_sum])
++- LogicalAggregate(group=[{0}], total_sum=[SUM($1)])
+ +- LogicalProject(c=[$1], a=[$0])
+- LogicalUnion(all=[true])
:- LogicalProject(a=[$0], c=[$2])
: +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable]])
@@ -374,14 +374,14 @@
LogicalSink(table=[default_catalog.default_database.retractSink], fields=[total_
</Resource>
<Resource name="optimized rel plan">
<![CDATA[
-Sink(table=[default_catalog.default_database.upsertSink], fields=[total_sum],
changelogMode=[NONE])
-+- GroupAggregate(select=[SUM(a) AS total_sum], changelogMode=[I,UA])
- +- Exchange(distribution=[single], changelogMode=[I])
- +- Union(all=[true], union=[a], changelogMode=[I])
- :- Calc(select=[a], changelogMode=[I])
+Sink(table=[default_catalog.default_database.upsertSink], fields=[c,
total_sum], changelogMode=[NONE])
++- GroupAggregate(groupBy=[c], select=[c, SUM(a) AS total_sum],
changelogMode=[I,UA])
+ +- Exchange(distribution=[hash[c]], changelogMode=[I])
+ +- Union(all=[true], union=[c, a], changelogMode=[I])
+ :- Calc(select=[c, a], changelogMode=[I])
: +- Calc(select=[a, c], changelogMode=[I])
: +- TableSourceScan(table=[[default_catalog, default_database,
MyTable]], fields=[a, b, c], changelogMode=[I])
- +- Calc(select=[d AS a], changelogMode=[I])
+ +- Calc(select=[f AS c, d AS a], changelogMode=[I])
+- Calc(select=[d, f], changelogMode=[I])
+- TableSourceScan(table=[[default_catalog, default_database,
MyTable1]], fields=[d, e, f], changelogMode=[I])
@@ -498,9 +498,9 @@
LogicalSink(table=[default_catalog.default_database.retractSink], fields=[total_
+- LogicalProject(a=[$0], c=[$2])
+- LogicalTableScan(table=[[default_catalog, default_database,
MyTable2]])
-LogicalSink(table=[default_catalog.default_database.upsertSink],
fields=[total_min])
-+- LogicalAggregate(group=[{}], total_min=[MIN($0)])
- +- LogicalProject(a=[$0])
+LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[c,
total_min])
++- LogicalAggregate(group=[{0}], total_min=[MIN($1)])
+ +- LogicalProject(c=[$1], a=[$0])
+- LogicalUnion(all=[true])
:- LogicalProject(a=[$0], c=[$1])
: +- LogicalUnion(all=[true])
@@ -535,17 +535,17 @@
Sink(table=[default_catalog.default_database.retractSink], fields=[total_sum], c
+- Calc(select=[a, c], changelogMode=[I])
+- TableSourceScan(table=[[default_catalog, default_database,
MyTable2]], fields=[a, b, c], changelogMode=[I])
-Sink(table=[default_catalog.default_database.upsertSink], fields=[total_min],
changelogMode=[NONE])
-+- GroupAggregate(select=[MIN(a) AS total_min], changelogMode=[I,UA])
- +- Exchange(distribution=[single], changelogMode=[I])
- +- Union(all=[true], union=[a], changelogMode=[I])
- :- Calc(select=[a], changelogMode=[I])
+Sink(table=[default_catalog.default_database.upsertSink], fields=[c,
total_min], changelogMode=[NONE])
++- GroupAggregate(groupBy=[c], select=[c, MIN(a) AS total_min],
changelogMode=[I,UA])
+ +- Exchange(distribution=[hash[c]], changelogMode=[I])
+ +- Union(all=[true], union=[c, a], changelogMode=[I])
+ :- Calc(select=[c, a], changelogMode=[I])
: +- Union(all=[true], union=[a, c], changelogMode=[I])
: :- Calc(select=[a, c], changelogMode=[I])
: : +- TableSourceScan(table=[[default_catalog,
default_database, MyTable]], fields=[a, b, c], changelogMode=[I])
: +- Calc(select=[d, f], changelogMode=[I])
: +- TableSourceScan(table=[[default_catalog,
default_database, MyTable1]], fields=[d, e, f], changelogMode=[I])
- +- Calc(select=[a], changelogMode=[I])
+ +- Calc(select=[c, a], changelogMode=[I])
+- Calc(select=[a, c], changelogMode=[I])
+- TableSourceScan(table=[[default_catalog, default_database,
MyTable2]], fields=[a, b, c], changelogMode=[I])
]]>
@@ -554,9 +554,9 @@ Sink(table=[default_catalog.default_database.upsertSink],
fields=[total_min], ch
<TestCase name="testMultiSinksSplitOnUnion4">
<Resource name="ast">
<![CDATA[
-LogicalSink(table=[default_catalog.default_database.upsertSink],
fields=[total_sum])
-+- LogicalAggregate(group=[{}], total_sum=[SUM($0)])
- +- LogicalProject(a=[$0])
+LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[c,
total_sum])
++- LogicalAggregate(group=[{0}], total_sum=[SUM($1)])
+ +- LogicalProject(c=[$1], a=[$0])
+- LogicalUnion(all=[true])
:- LogicalUnion(all=[true])
: :- LogicalProject(a=[$0], c=[$2])
@@ -581,18 +581,18 @@
LogicalSink(table=[default_catalog.default_database.retractSink], fields=[total_
</Resource>
<Resource name="optimized rel plan">
<![CDATA[
-Sink(table=[default_catalog.default_database.upsertSink], fields=[total_sum],
changelogMode=[NONE])
-+- GroupAggregate(select=[SUM(a) AS total_sum], changelogMode=[I,UA])
- +- Exchange(distribution=[single], changelogMode=[I])
- +- Union(all=[true], union=[a], changelogMode=[I])
- :- Union(all=[true], union=[a], changelogMode=[I])
- : :- Calc(select=[a], changelogMode=[I])
+Sink(table=[default_catalog.default_database.upsertSink], fields=[c,
total_sum], changelogMode=[NONE])
++- GroupAggregate(groupBy=[c], select=[c, SUM(a) AS total_sum],
changelogMode=[I,UA])
+ +- Exchange(distribution=[hash[c]], changelogMode=[I])
+ +- Union(all=[true], union=[c, a], changelogMode=[I])
+ :- Union(all=[true], union=[c, a], changelogMode=[I])
+ : :- Calc(select=[c, a], changelogMode=[I])
: : +- Calc(select=[a, c], changelogMode=[I])
: : +- TableSourceScan(table=[[default_catalog,
default_database, MyTable]], fields=[a, b, c], changelogMode=[I])
- : +- Calc(select=[d AS a], changelogMode=[I])
+ : +- Calc(select=[f AS c, d AS a], changelogMode=[I])
: +- Calc(select=[d, f], changelogMode=[I])
: +- TableSourceScan(table=[[default_catalog,
default_database, MyTable1]], fields=[d, e, f], changelogMode=[I])
- +- Calc(select=[a], changelogMode=[I])
+ +- Calc(select=[c, a], changelogMode=[I])
+- Calc(select=[a, c], changelogMode=[I])
+- TableSourceScan(table=[[default_catalog, default_database,
MyTable2]], fields=[a, b, c], changelogMode=[I])
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 787cbf39116..9b7e0d6b00d 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
@@ -19,6 +19,7 @@ package
org.apache.flink.table.planner.plan.rules.physical.stream
import org.apache.flink.table.api.ExplainDetail
import org.apache.flink.table.api.config.{AggregatePhaseStrategy,
OptimizerConfigOptions}
+import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.UpdatingJoinFunction
import
org.apache.flink.table.planner.plan.optimize.program.FlinkChangelogModeInferenceProgram
import
org.apache.flink.table.planner.utils.ImmutableColConstraintTestUtils.addImmutableColConstraint
import org.apache.flink.table.planner.utils.TableTestBase
@@ -88,6 +89,17 @@ class ChangelogModeInferenceTest extends TableTestBase {
| 'changelog-mode' = 'I,UA,UB,D'
|)
""".stripMargin)
+
+ util.addTable("""
+ |CREATE TABLE LookupDim (
+ | id BIGINT,
+ | v STRING,
+ | PRIMARY KEY (id) NOT ENFORCED
+ |) WITH (
+ | 'connector' = 'values',
+ | 'bounded' = 'true'
+ |)
+ """.stripMargin)
}
@Test
@@ -583,4 +595,96 @@ class ChangelogModeInferenceTest extends TableTestBase {
"INSERT INTO keyless_upsert_sink_no_key SELECT rate FROM
DeduplicatedView",
ExplainDetail.CHANGELOG_MODE)
}
+
+ @Test
+ def testKeylessUpsertSinkFallsBackToRetractOnGlobalAggregate(): Unit = {
+ // A global aggregation reports an empty upsert key, meaning "at most one
row" rather than a
+ // set of columns. The keyless sink must fall back to retract, or
consumers receive bare
+ // UPDATE_AFTER rows with no key to apply them by.
+ util.addTable("""
+ |CREATE TABLE keyless_upsert_sink_count (
+ | cnt BIGINT
+ |) WITH (
+ | 'connector' = 'values',
+ | 'sink-insert-only' = 'false',
+ | 'sink-changelog-mode-enforced' = 'I,UA,D'
+ |)
+ |""".stripMargin)
+ util.verifyRelPlanInsert(
+ "INSERT INTO keyless_upsert_sink_count SELECT COUNT(*) FROM MyTable",
+ ExplainDetail.CHANGELOG_MODE)
+ }
+
+ @Test
+ def testKeylessUpsertSinkWithLookupJoinOnGlobalAggregateProbeSide(): Unit = {
+ // The reported key (dim.id) isn't invariant across updates - it really
tracks the volatile
+ // join value (cnt) - but lookup joins don't filter that out, unlike
regular joins.
+ util.addTable("""
+ |CREATE TABLE lookup_sink (
+ | id BIGINT,
+ | v STRING
+ |) WITH (
+ | 'connector' = 'values',
+ | 'sink-insert-only' = 'false',
+ | 'sink-changelog-mode-enforced' = 'I,UA,D'
+ |)
+ |""".stripMargin)
+ val sql =
+ """
+ |INSERT INTO lookup_sink
+ |SELECT dim.id, dim.v
+ |FROM (SELECT COUNT(*) AS cnt, PROCTIME() AS pt FROM MyTable) g
+ |JOIN LookupDim FOR SYSTEM_TIME AS OF g.pt AS dim
+ |ON g.cnt = dim.id
+ """.stripMargin
+ util.verifyRelPlanInsert(sql, ExplainDetail.CHANGELOG_MODE)
+ }
+
+ @Test
+ def
testKeylessUpsertSinkWithMultiTableArgUpsertPtfComputesDistinctKeyPerArg():
Unit = {
+ // Multi-arg PTF with a partitioned table argument after the first: the
2nd+ arg's partition
+ // columns must not collapse to an empty key, or a real key gets mixed
with a spurious empty
+ // candidate and the sink would wrongly stay upsert.
+ util.addTemporarySystemFunction("f", classOf[UpdatingJoinFunction])
+ util.addTable("""
+ |CREATE TABLE scores_ptf_probe (
+ | name STRING,
+ | score INT,
+ | PRIMARY KEY(name) NOT ENFORCED
+ |) WITH (
+ | 'connector' = 'values',
+ | 'changelog-mode' = 'I,UA,D',
+ | 'source.produces-delete-by-key' = 'true'
+ |)
+ |""".stripMargin)
+ util.addTable("""
+ |CREATE TABLE city_ptf_probe (
+ | name STRING,
+ | city STRING,
+ | PRIMARY KEY(name) NOT ENFORCED
+ |) WITH (
+ | 'connector' = 'values',
+ | 'changelog-mode' = 'I,UA,D',
+ | 'source.produces-delete-by-key' = 'true'
+ |)
+ |""".stripMargin)
+ util.addTable("""
+ |CREATE TABLE keyless_upsert_sink_ptf_probe (
+ | name STRING,
+ | name0 STRING,
+ | `out` STRING
+ |) WITH (
+ | 'connector' = 'values',
+ | 'sink-insert-only' = 'false',
+ | 'sink-changelog-mode-enforced' = 'I,UA,D',
+ | 'sink.supports-delete-by-key' = 'true'
+ |)
+ |""".stripMargin)
+ util.verifyRelPlanInsert(
+ "INSERT INTO keyless_upsert_sink_ptf_probe SELECT `name`, name0, `out`
FROM f("
+ + "scoreTable => TABLE scores_ptf_probe PARTITION BY name, "
+ + "cityTable => TABLE city_ptf_probe PARTITION BY name)",
+ ExplainDetail.CHANGELOG_MODE
+ )
+ }
}
diff --git
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.scala
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.scala
index 0daee92964e..c9d0c88b7e2 100644
---
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.scala
+++
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.scala
@@ -394,13 +394,14 @@ class DagOptimizationTest extends TableTestBase {
util.tableEnv.sqlQuery("SELECT a, c FROM MyTable UNION ALL SELECT d, f
FROM MyTable1")
util.tableEnv.createTemporaryView("TempTable", table)
- val table1 = util.tableEnv.sqlQuery("SELECT SUM(a) AS total_sum FROM
TempTable")
+ val table1 = util.tableEnv.sqlQuery("SELECT c, SUM(a) AS total_sum FROM
TempTable GROUP BY c")
TestSinkUtil.addValuesSink(
util.tableEnv,
"upsertSink",
- List("total_sum"),
- List(INT),
- ChangelogMode.upsert()
+ List("c", "total_sum"),
+ List(STRING, INT),
+ ChangelogMode.upsert(),
+ pk = List("c")
)
stmtSet.addInsert("upsertSink", table1)
@@ -508,13 +509,15 @@ class DagOptimizationTest extends TableTestBase {
)
stmtSet.addInsert("retractSink", table2)
- val table3 = util.tableEnv.sqlQuery("SELECT MIN(a) AS total_min FROM
TempTable1")
+ val table3 =
+ util.tableEnv.sqlQuery("SELECT c, MIN(a) AS total_min FROM TempTable1
GROUP BY c")
TestSinkUtil.addValuesSink(
util.tableEnv,
"upsertSink",
- List("total_min"),
- List(INT),
- ChangelogMode.upsert()
+ List("c", "total_min"),
+ List(STRING, INT),
+ ChangelogMode.upsert(),
+ pk = List("c")
)
stmtSet.addInsert("upsertSink", table3)
@@ -540,13 +543,14 @@ class DagOptimizationTest extends TableTestBase {
val table = util.tableEnv.sqlQuery(sqlQuery)
util.tableEnv.createTemporaryView("TempTable", table)
- val table1 = util.tableEnv.sqlQuery("SELECT SUM(a) AS total_sum FROM
TempTable")
+ val table1 = util.tableEnv.sqlQuery("SELECT c, SUM(a) AS total_sum FROM
TempTable GROUP BY c")
TestSinkUtil.addValuesSink(
util.tableEnv,
"upsertSink",
- List("total_sum"),
- List(INT),
- ChangelogMode.upsert()
+ List("c", "total_sum"),
+ List(STRING, INT),
+ ChangelogMode.upsert(),
+ pk = List("c")
)
stmtSet.addInsert("upsertSink", table1)