Skip to content

Commit 2245026

Browse files
[FLINK-40447][table] Fall back to retract when input upsert key is empty or mixed
This closes apache#29000.
1 parent 9ee1037 commit 2245026

6 files changed

Lines changed: 239 additions & 46 deletions

File tree

flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalProcessTableFunction.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -515,7 +515,7 @@ public static Set<ImmutableBitSet> toPartitionColumns(RexCall call) {
515515
// f(t1 PARTITION BY (k1, k2), t2 PARTITION BY (k3, k4))
516516
// -> [k1, k2, k3, k4, function out...]
517517
final List<Integer> partitionColumns =
518-
IntStream.range(pos, partitionKeyCount)
518+
IntStream.range(pos, pos + partitionKeyCount)
519519
.boxed()
520520
.collect(Collectors.toList());
521521
pos += partitionKeyCount;

flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1131,16 +1131,19 @@ class FlinkChangelogModeInferenceProgram extends FlinkOptimizeProgram[StreamOpti
11311131
* when the upsert key has columns outside sink pk. This differs from batch job's unique key
11321132
* inference.
11331133
*
1134-
* <p>A sink without a primary key is satisfied whenever the input carries any upsert key.
1134+
* <p>A sink without a primary key is satisfied whenever the input carries a real (non-empty)
1135+
* upsert key; an empty candidate ("at most one row") never counts, even alongside a real one.
11351136
*/
11361137
private def canUpsertKeysWithImmutableColsSatisfyPk(sink: StreamPhysicalSink): Boolean = {
11371138
val sinkDefinedPks = sink.contextResolvedTable.getResolvedSchema.getPrimaryKeyIndexes
11381139
val fmq = FlinkRelMetadataQuery.reuseOrCreate(sink.getCluster.getMetadataQuery)
11391140
val changeLogUpsertKeys = fmq.getUpsertKeys(sink.getInput)
11401141
if (sinkDefinedPks.isEmpty) {
1141-
// A keyless sink cannot apply UPDATE_AFTER in place, so it can only accept upsert when the
1142-
// input itself carries an upsert key; otherwise fall back to beforeAndAfter.
1143-
return changeLogUpsertKeys != null && !changeLogUpsertKeys.isEmpty
1142+
// A keyless sink can only stay upsert when the input has a real, column-based upsert
1143+
// key. An empty candidate means "at most one row" (e.g. a global aggregate), not columns
1144+
// to match on - UpsertKeyUtil.getSmallestKey would otherwise prefer it over a real one.
1145+
return changeLogUpsertKeys != null && changeLogUpsertKeys.nonEmpty &&
1146+
!changeLogUpsertKeys.exists(_.isEmpty)
11441147
}
11451148
val sinkPks = ImmutableBitSet.of(sinkDefinedPks: _*)
11461149
// if upsert key is null, pk cannot be satisfied, should fall back to beforeAndAfter

flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/physical/stream/ChangelogModeInferenceTest.xml

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -172,6 +172,27 @@ GroupAggregate(groupBy=[cnt], select=[cnt, COUNT_RETRACT(cnt) AS frequency], cha
172172
: +- LegacyTableSourceScan(table=[[default_catalog, default_database, MyTable, source: [CollectionTableSource(word, number)]]], fields=[word, number], changelogMode=[I])
173173
+- Calc(select=[CAST(cnt AS BIGINT) AS cnt], changelogMode=[I])
174174
+- LegacyTableSourceScan(table=[[default_catalog, default_database, MyTable2, source: [CollectionTableSource(word, cnt)]]], fields=[word, cnt], changelogMode=[I])
175+
]]>
176+
</Resource>
177+
</TestCase>
178+
<TestCase name="testKeylessUpsertSinkFallsBackToRetractOnGlobalAggregate">
179+
<Resource name="sql">
180+
<![CDATA[INSERT INTO keyless_upsert_sink_count SELECT COUNT(*) FROM MyTable]]>
181+
</Resource>
182+
<Resource name="ast">
183+
<![CDATA[
184+
LogicalSink(table=[default_catalog.default_database.keyless_upsert_sink_count], fields=[EXPR$0])
185+
+- LogicalAggregate(group=[{}], EXPR$0=[COUNT()])
186+
+- LogicalTableScan(table=[[default_catalog, default_database, MyTable, source: [CollectionTableSource(word, number)]]])
187+
]]>
188+
</Resource>
189+
<Resource name="optimized rel plan">
190+
<![CDATA[
191+
Sink(table=[default_catalog.default_database.keyless_upsert_sink_count], fields=[EXPR$0], changelogMode=[NONE])
192+
+- GroupAggregate(select=[COUNT(*) AS EXPR$0], changelogMode=[I,UB,UA])
193+
+- Exchange(distribution=[single], changelogMode=[I])
194+
+- Calc(select=[0 AS $f0], changelogMode=[I])
195+
+- LegacyTableSourceScan(table=[[default_catalog, default_database, MyTable, source: [CollectionTableSource(word, number)]]], fields=[word, number], changelogMode=[I])
175196
]]>
176197
</Resource>
177198
</TestCase>
@@ -224,6 +245,67 @@ Sink(table=[default_catalog.default_database.keyless_upsert_sink], fields=[curre
224245
+- Exchange(distribution=[hash[currency]], changelogMode=[I])
225246
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime], changelogMode=[I])
226247
+- LegacyTableSourceScan(table=[[default_catalog, default_database, ratesHistory, source: [CollectionTableSource(currency, rate, rowtime)]]], fields=[currency, rate, rowtime], changelogMode=[I])
248+
]]>
249+
</Resource>
250+
</TestCase>
251+
<TestCase name="testKeylessUpsertSinkWithLookupJoinOnGlobalAggregateProbeSide">
252+
<Resource name="sql">
253+
<![CDATA[
254+
INSERT INTO lookup_sink
255+
SELECT dim.id, dim.v
256+
FROM (SELECT COUNT(*) AS cnt, PROCTIME() AS pt FROM MyTable) g
257+
JOIN LookupDim FOR SYSTEM_TIME AS OF g.pt AS dim
258+
ON g.cnt = dim.id
259+
]]>
260+
</Resource>
261+
<Resource name="ast">
262+
<![CDATA[
263+
LogicalSink(table=[default_catalog.default_database.lookup_sink], fields=[id, v])
264+
+- LogicalProject(id=[$2], v=[$3])
265+
+- LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{0, 1}])
266+
:- LogicalProject(cnt=[$0], pt=[PROCTIME()])
267+
: +- LogicalAggregate(group=[{}], cnt=[COUNT()])
268+
: +- LogicalTableScan(table=[[default_catalog, default_database, MyTable, source: [CollectionTableSource(word, number)]]])
269+
+- LogicalFilter(condition=[=($cor0.cnt, $0)])
270+
+- LogicalSnapshot(period=[$cor0.pt])
271+
+- LogicalTableScan(table=[[default_catalog, default_database, LookupDim]])
272+
]]>
273+
</Resource>
274+
<Resource name="optimized rel plan">
275+
<![CDATA[
276+
Sink(table=[default_catalog.default_database.lookup_sink], fields=[id, v], changelogMode=[NONE])
277+
+- Calc(select=[id, v], changelogMode=[I,UB,UA])
278+
+- LookupJoin(table=[default_catalog.default_database.LookupDim], joinType=[InnerJoin], lookup=[id=cnt], select=[cnt, id, v], upsertKey=[[]], changelogMode=[I,UB,UA])
279+
+- GroupAggregate(select=[COUNT(*) AS cnt], changelogMode=[I,UB,UA])
280+
+- Exchange(distribution=[single], changelogMode=[I])
281+
+- Calc(select=[0 AS $f0], changelogMode=[I])
282+
+- LegacyTableSourceScan(table=[[default_catalog, default_database, MyTable, source: [CollectionTableSource(word, number)]]], fields=[word, number], changelogMode=[I])
283+
]]>
284+
</Resource>
285+
</TestCase>
286+
<TestCase name="testKeylessUpsertSinkWithMultiTableArgUpsertPtfComputesDistinctKeyPerArg">
287+
<Resource name="sql">
288+
<![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)]]>
289+
</Resource>
290+
<Resource name="ast">
291+
<![CDATA[
292+
LogicalSink(table=[default_catalog.default_database.keyless_upsert_sink_ptf_probe], fields=[name, name0, out])
293+
+- LogicalProject(name=[$0], name0=[$1], out=[$2])
294+
+- 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)])
295+
:- LogicalProject(name=[$0], score=[$1])
296+
: +- LogicalTableScan(table=[[default_catalog, default_database, scores_ptf_probe]])
297+
+- LogicalProject(name=[$0], city=[$1])
298+
+- LogicalTableScan(table=[[default_catalog, default_database, city_ptf_probe]])
299+
]]>
300+
</Resource>
301+
<Resource name="optimized rel plan">
302+
<![CDATA[
303+
Sink(table=[default_catalog.default_database.keyless_upsert_sink_ptf_probe], fields=[name, name0, out], changelogMode=[NONE])
304+
+- 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])
305+
:- Exchange(distribution=[hash[name]], changelogMode=[I,UA,PD])
306+
: +- TableSourceScan(table=[[default_catalog, default_database, scores_ptf_probe]], fields=[name, score], changelogMode=[I,UA,PD])
307+
+- Exchange(distribution=[hash[name]], changelogMode=[I,UA,PD])
308+
+- TableSourceScan(table=[[default_catalog, default_database, city_ptf_probe]], fields=[name, city], changelogMode=[I,UA,PD])
227309
]]>
228310
</Resource>
229311
</TestCase>

flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/DagOptimizationTest.xml

Lines changed: 29 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -353,9 +353,9 @@ Sink(table=[default_catalog.default_database.retractSink2], fields=[total_min],
353353
<TestCase name="testMultiSinksSplitOnUnion1">
354354
<Resource name="ast">
355355
<![CDATA[
356-
LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[total_sum])
357-
+- LogicalAggregate(group=[{}], total_sum=[SUM($0)])
358-
+- LogicalProject(a=[$0])
356+
LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[c, total_sum])
357+
+- LogicalAggregate(group=[{0}], total_sum=[SUM($1)])
358+
+- LogicalProject(c=[$1], a=[$0])
359359
+- LogicalUnion(all=[true])
360360
:- LogicalProject(a=[$0], c=[$2])
361361
: +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
@@ -374,14 +374,14 @@ LogicalSink(table=[default_catalog.default_database.retractSink], fields=[total_
374374
</Resource>
375375
<Resource name="optimized rel plan">
376376
<![CDATA[
377-
Sink(table=[default_catalog.default_database.upsertSink], fields=[total_sum], changelogMode=[NONE])
378-
+- GroupAggregate(select=[SUM(a) AS total_sum], changelogMode=[I,UA])
379-
+- Exchange(distribution=[single], changelogMode=[I])
380-
+- Union(all=[true], union=[a], changelogMode=[I])
381-
:- Calc(select=[a], changelogMode=[I])
377+
Sink(table=[default_catalog.default_database.upsertSink], fields=[c, total_sum], changelogMode=[NONE])
378+
+- GroupAggregate(groupBy=[c], select=[c, SUM(a) AS total_sum], changelogMode=[I,UA])
379+
+- Exchange(distribution=[hash[c]], changelogMode=[I])
380+
+- Union(all=[true], union=[c, a], changelogMode=[I])
381+
:- Calc(select=[c, a], changelogMode=[I])
382382
: +- Calc(select=[a, c], changelogMode=[I])
383383
: +- TableSourceScan(table=[[default_catalog, default_database, MyTable]], fields=[a, b, c], changelogMode=[I])
384-
+- Calc(select=[d AS a], changelogMode=[I])
384+
+- Calc(select=[f AS c, d AS a], changelogMode=[I])
385385
+- Calc(select=[d, f], changelogMode=[I])
386386
+- TableSourceScan(table=[[default_catalog, default_database, MyTable1]], fields=[d, e, f], changelogMode=[I])
387387
@@ -498,9 +498,9 @@ LogicalSink(table=[default_catalog.default_database.retractSink], fields=[total_
498498
+- LogicalProject(a=[$0], c=[$2])
499499
+- LogicalTableScan(table=[[default_catalog, default_database, MyTable2]])
500500
501-
LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[total_min])
502-
+- LogicalAggregate(group=[{}], total_min=[MIN($0)])
503-
+- LogicalProject(a=[$0])
501+
LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[c, total_min])
502+
+- LogicalAggregate(group=[{0}], total_min=[MIN($1)])
503+
+- LogicalProject(c=[$1], a=[$0])
504504
+- LogicalUnion(all=[true])
505505
:- LogicalProject(a=[$0], c=[$1])
506506
: +- LogicalUnion(all=[true])
@@ -535,17 +535,17 @@ Sink(table=[default_catalog.default_database.retractSink], fields=[total_sum], c
535535
+- Calc(select=[a, c], changelogMode=[I])
536536
+- TableSourceScan(table=[[default_catalog, default_database, MyTable2]], fields=[a, b, c], changelogMode=[I])
537537
538-
Sink(table=[default_catalog.default_database.upsertSink], fields=[total_min], changelogMode=[NONE])
539-
+- GroupAggregate(select=[MIN(a) AS total_min], changelogMode=[I,UA])
540-
+- Exchange(distribution=[single], changelogMode=[I])
541-
+- Union(all=[true], union=[a], changelogMode=[I])
542-
:- Calc(select=[a], changelogMode=[I])
538+
Sink(table=[default_catalog.default_database.upsertSink], fields=[c, total_min], changelogMode=[NONE])
539+
+- GroupAggregate(groupBy=[c], select=[c, MIN(a) AS total_min], changelogMode=[I,UA])
540+
+- Exchange(distribution=[hash[c]], changelogMode=[I])
541+
+- Union(all=[true], union=[c, a], changelogMode=[I])
542+
:- Calc(select=[c, a], changelogMode=[I])
543543
: +- Union(all=[true], union=[a, c], changelogMode=[I])
544544
: :- Calc(select=[a, c], changelogMode=[I])
545545
: : +- TableSourceScan(table=[[default_catalog, default_database, MyTable]], fields=[a, b, c], changelogMode=[I])
546546
: +- Calc(select=[d, f], changelogMode=[I])
547547
: +- TableSourceScan(table=[[default_catalog, default_database, MyTable1]], fields=[d, e, f], changelogMode=[I])
548-
+- Calc(select=[a], changelogMode=[I])
548+
+- Calc(select=[c, a], changelogMode=[I])
549549
+- Calc(select=[a, c], changelogMode=[I])
550550
+- TableSourceScan(table=[[default_catalog, default_database, MyTable2]], fields=[a, b, c], changelogMode=[I])
551551
]]>
@@ -554,9 +554,9 @@ Sink(table=[default_catalog.default_database.upsertSink], fields=[total_min], ch
554554
<TestCase name="testMultiSinksSplitOnUnion4">
555555
<Resource name="ast">
556556
<![CDATA[
557-
LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[total_sum])
558-
+- LogicalAggregate(group=[{}], total_sum=[SUM($0)])
559-
+- LogicalProject(a=[$0])
557+
LogicalSink(table=[default_catalog.default_database.upsertSink], fields=[c, total_sum])
558+
+- LogicalAggregate(group=[{0}], total_sum=[SUM($1)])
559+
+- LogicalProject(c=[$1], a=[$0])
560560
+- LogicalUnion(all=[true])
561561
:- LogicalUnion(all=[true])
562562
: :- LogicalProject(a=[$0], c=[$2])
@@ -581,18 +581,18 @@ LogicalSink(table=[default_catalog.default_database.retractSink], fields=[total_
581581
</Resource>
582582
<Resource name="optimized rel plan">
583583
<![CDATA[
584-
Sink(table=[default_catalog.default_database.upsertSink], fields=[total_sum], changelogMode=[NONE])
585-
+- GroupAggregate(select=[SUM(a) AS total_sum], changelogMode=[I,UA])
586-
+- Exchange(distribution=[single], changelogMode=[I])
587-
+- Union(all=[true], union=[a], changelogMode=[I])
588-
:- Union(all=[true], union=[a], changelogMode=[I])
589-
: :- Calc(select=[a], changelogMode=[I])
584+
Sink(table=[default_catalog.default_database.upsertSink], fields=[c, total_sum], changelogMode=[NONE])
585+
+- GroupAggregate(groupBy=[c], select=[c, SUM(a) AS total_sum], changelogMode=[I,UA])
586+
+- Exchange(distribution=[hash[c]], changelogMode=[I])
587+
+- Union(all=[true], union=[c, a], changelogMode=[I])
588+
:- Union(all=[true], union=[c, a], changelogMode=[I])
589+
: :- Calc(select=[c, a], changelogMode=[I])
590590
: : +- Calc(select=[a, c], changelogMode=[I])
591591
: : +- TableSourceScan(table=[[default_catalog, default_database, MyTable]], fields=[a, b, c], changelogMode=[I])
592-
: +- Calc(select=[d AS a], changelogMode=[I])
592+
: +- Calc(select=[f AS c, d AS a], changelogMode=[I])
593593
: +- Calc(select=[d, f], changelogMode=[I])
594594
: +- TableSourceScan(table=[[default_catalog, default_database, MyTable1]], fields=[d, e, f], changelogMode=[I])
595-
+- Calc(select=[a], changelogMode=[I])
595+
+- Calc(select=[c, a], changelogMode=[I])
596596
+- Calc(select=[a, c], changelogMode=[I])
597597
+- TableSourceScan(table=[[default_catalog, default_database, MyTable2]], fields=[a, b, c], changelogMode=[I])
598598

0 commit comments

Comments
 (0)