This is an automated email from the ASF dual-hosted git repository.
RocMarshal 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 21567fb986f [FLINK-40170][table-planner] Infer update-producing
changelog mode for early-fire interval join (#28877)
21567fb986f is described below
commit 21567fb986fdc91d4f41b21233a622d2540b77c5
Author: Weiqing Yang <[email protected]>
AuthorDate: Mon Aug 10 06:34:31 2026 -0700
[FLINK-40170][table-planner] Infer update-producing changelog mode for
early-fire interval join (#28877)
---
.../stream/StreamPhysicalIntervalJoin.scala | 10 +++
.../FlinkChangelogModeInferenceProgram.scala | 24 +++++-
.../plan/hints/stream/EarlyFireJoinHintTest.java | 61 +++++++++++++-
.../plan/hints/stream/EarlyFireJoinHintTest.xml | 96 ++++++++++++++++++++++
.../testEarlyFireJsonPlanRoundTrip.out | 5 +-
5 files changed, 190 insertions(+), 6 deletions(-)
diff --git
a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalIntervalJoin.scala
b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalIntervalJoin.scala
index d3401783989..e47cad9522c 100644
---
a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalIntervalJoin.scala
+++
b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalIntervalJoin.scala
@@ -64,6 +64,16 @@ class StreamPhysicalIntervalJoin(
override def requireWatermark: Boolean = windowBounds.isEventTime
+ /**
+ * Whether this interval join produces update changes because of the
EARLY_FIRE hint. Only an
+ * outer join with a non-negative window can speculatively emit a padded row
and later correct it;
+ * a negative-window join only ever emits inserts, so it must stay
insert-only even with the hint
+ * set.
+ */
+ def produceEarlyFireUpdates: Boolean =
+ earlyFireDelay != null && getJoinType.isOuterJoin &&
+ windowBounds.getLeftUpperBound >= windowBounds.getLeftLowerBound
+
override def copy(
traitSet: RelTraitSet,
conditionExpr: RexNode,
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 5d578e2b076..27719c67cdd 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
@@ -362,9 +362,27 @@ class FlinkChangelogModeInferenceProgram extends
FlinkOptimizeProgram[StreamOpti
val providedTrait = new ModifyKindSetTrait(builder.build())
createNewNode(over, children, providedTrait, requiredTrait, requester)
- case _: StreamPhysicalTemporalSort | _: StreamPhysicalIntervalJoin |
- _: StreamPhysicalPythonOverAggregate =>
- // TemporalSort, IntervalJoin only support consuming insert-only
+ case intervalJoin: StreamPhysicalIntervalJoin =>
+ // The interval join consumes insert-only input. Without the
EARLY_FIRE hint it also only
+ // produces insert-only changes; an early-firing outer join
additionally produces update
+ // changes, because it speculatively emits a padded row and later
corrects it on a match.
+ val children = visitChildren(intervalJoin,
ModifyKindSetTrait.INSERT_ONLY)
+ val builder =
ModifyKindSet.newBuilder().addContainedKind(ModifyKind.INSERT)
+ if (intervalJoin.produceEarlyFireUpdates) {
+ builder.addContainedKind(ModifyKind.UPDATE)
+ }
+ val providedTrait = new ModifyKindSetTrait(builder.build())
+ if (intervalJoin.produceEarlyFireUpdates &&
!providedTrait.satisfies(requiredTrait)) {
+ throw new TableException(
+ s"$requester doesn't support consuming update changes, but the
EARLY_FIRE hint " +
+ "makes this outer interval join produce update changes (a padded
row is emitted " +
+ "speculatively and later corrected on a match). Remove the
EARLY_FIRE hint, or " +
+ "write into a downstream/sink that accepts update changes.")
+ }
+ createNewNode(intervalJoin, children, providedTrait, requiredTrait,
requester)
+
+ case _: StreamPhysicalTemporalSort | _:
StreamPhysicalPythonOverAggregate =>
+ // TemporalSort and PythonOverAggregate only support consuming
insert-only
// and producing insert-only changes
val children = visitChildren(rel, ModifyKindSetTrait.INSERT_ONLY)
createNewNode(rel, children, ModifyKindSetTrait.INSERT_ONLY,
requiredTrait, requester)
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.java
index 1db3dc6e8c2..05879701b33 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.java
@@ -27,6 +27,8 @@ import org.apache.flink.table.planner.utils.TableTestBase;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import java.util.Collections;
+
import scala.Enumeration;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -71,7 +73,8 @@ class EarlyFireJoinHintTest extends TableTestBase {
+ " a INT,\n"
+ " b VARCHAR\n"
+ ") WITH (\n"
- + " 'connector' = 'values'\n"
+ + " 'connector' = 'values',\n"
+ + " 'sink-insert-only' = 'false'\n"
+ ")");
}
@@ -212,6 +215,58 @@ class EarlyFireJoinHintTest extends TableTestBase {
verify(sql);
}
+ @Test
+ void testEarlyFireOuterJoinProducesUpdates() {
+ String sql =
+ "SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b\n"
+ + "FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON\n"
+ + " t1.a = t2.a AND\n"
+ + " t1.rowtime BETWEEN t2.rowtime - INTERVAL '10'
SECOND AND t2.rowtime + INTERVAL '1' HOUR";
+ verifyChangelogMode(sql);
+ }
+
+ @Test
+ void testEarlyFireOuterJoinIntoInsertOnlySinkFails() {
+ util.tableEnv()
+ .executeSql(
+ "CREATE TABLE InsertOnlySink (\n"
+ + " a INT,\n"
+ + " b VARCHAR\n"
+ + ") WITH (\n"
+ + " 'connector' = 'values',\n"
+ + " 'sink-insert-only' = 'true'\n"
+ + ")");
+ String insert =
+ "INSERT INTO InsertOnlySink\n"
+ + "SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b\n"
+ + "FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON\n"
+ + " t1.a = t2.a AND\n"
+ + " t1.rowtime BETWEEN t2.rowtime - INTERVAL '10'
SECOND AND t2.rowtime + INTERVAL '1' HOUR";
+ assertThatThrownBy(() -> util.verifyRelPlanInsert(insert))
+ .hasMessageContaining(
+ "the EARLY_FIRE hint makes this outer interval join
produce update");
+ }
+
+ @Test
+ void testEarlyFireNegativeWindowStaysInsertOnly() {
+ String sql =
+ "SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b\n"
+ + "FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON\n"
+ + " t1.a = t2.a AND\n"
+ + " t1.rowtime BETWEEN t2.rowtime + INTERVAL '10'
SECOND AND t2.rowtime + INTERVAL '5' SECOND";
+ verifyChangelogMode(sql);
+ }
+
+ @Test
+ void testEarlyFireInnerJoinStaysInsertOnly() {
+ String sql =
+ "SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b\n"
+ + "FROM MyTable t1 JOIN MyTable2 t2 ON\n"
+ + " t1.a = t2.a AND\n"
+ + " t1.rowtime BETWEEN t2.rowtime - INTERVAL '10'
SECOND AND t2.rowtime + INTERVAL '1' HOUR";
+ verifyChangelogMode(sql);
+ }
+
@Test
void testEarlyFireJsonPlanRoundTrip() {
String insert =
@@ -231,4 +286,8 @@ class EarlyFireJoinHintTest extends TableTestBase {
new Enumeration.Value[] {PlanKind.AST(), PlanKind.OPT_EXEC()},
false);
}
+
+ private void verifyChangelogMode(String sql) {
+ util.verifyRelPlan(sql,
Collections.singletonList(ExplainDetail.CHANGELOG_MODE));
+ }
}
diff --git
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.xml
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.xml
index 8daae577559..8faa17c8bea 100644
---
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.xml
+++
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.xml
@@ -45,6 +45,38 @@ Calc(select=[a, b])
+- Exchange(distribution=[hash[a]])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime])
+- TableSourceScan(table=[[default_catalog, default_database,
MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime])
+]]>
+ </Resource>
+ </TestCase>
+ <TestCase name="testEarlyFireInnerJoinStaysInsertOnly">
+ <Resource name="sql">
+ <![CDATA[SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b
+FROM MyTable t1 JOIN MyTable2 t2 ON
+ t1.a = t2.a AND
+ t1.rowtime BETWEEN t2.rowtime - INTERVAL '10' SECOND AND t2.rowtime +
INTERVAL '1' HOUR]]>
+ </Resource>
+ <Resource name="ast">
+ <![CDATA[
+LogicalProject(a=[$0], b=[$6])
++- LogicalJoin(condition=[AND(=($0, $5), >=($4, -($9, 10000:INTERVAL SECOND)),
<=($4, +($9, 3600000:INTERVAL HOUR)))], joinType=[inner],
joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s}]]])
+ :- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+ : +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()],
rowtime=[$3])
+ : +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable]])
+ +- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+ +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()],
rowtime=[$3])
+ +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable2]])
+]]>
+ </Resource>
+ <Resource name="optimized rel plan">
+ <![CDATA[
+Calc(select=[a, b], changelogMode=[I])
++- IntervalJoin(joinType=[InnerJoin], windowBounds=[isRowTime=true,
leftLowerBound=-10000, leftUpperBound=3600000, leftTimeIndex=1,
rightTimeIndex=2], where=[AND(=(a, a0), >=(rowtime, -(rowtime0, 10000:INTERVAL
SECOND)), <=(rowtime, +(rowtime0, 3600000:INTERVAL HOUR)))], select=[a,
rowtime, a0, b, rowtime0], earlyFireDelay=[5000], earlyFireTimeMode=[ROWTIME],
changelogMode=[I])
+ :- Exchange(distribution=[hash[a]], changelogMode=[I])
+ : +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime],
changelogMode=[I])
+ : +- TableSourceScan(table=[[default_catalog, default_database,
MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime],
changelogMode=[I])
+ +- Exchange(distribution=[hash[a]], changelogMode=[I])
+ +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime],
changelogMode=[I])
+ +- TableSourceScan(table=[[default_catalog, default_database,
MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime],
changelogMode=[I])
]]>
</Resource>
</TestCase>
@@ -77,6 +109,38 @@ Calc(select=[a, b])
+- Exchange(distribution=[hash[a]])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime])
+- TableSourceScan(table=[[default_catalog, default_database,
MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime])
+]]>
+ </Resource>
+ </TestCase>
+ <TestCase name="testEarlyFireNegativeWindowStaysInsertOnly">
+ <Resource name="sql">
+ <![CDATA[SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b
+FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON
+ t1.a = t2.a AND
+ t1.rowtime BETWEEN t2.rowtime + INTERVAL '10' SECOND AND t2.rowtime +
INTERVAL '5' SECOND]]>
+ </Resource>
+ <Resource name="ast">
+ <![CDATA[
+LogicalProject(a=[$0], b=[$6])
++- LogicalJoin(condition=[AND(=($0, $5), >=($4, +($9, 10000:INTERVAL SECOND)),
<=($4, +($9, 5000:INTERVAL SECOND)))], joinType=[left], joinHints=[[[EARLY_FIRE
inheritPath:[0] options:{delay=5s}]]])
+ :- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+ : +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()],
rowtime=[$3])
+ : +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable]])
+ +- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+ +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()],
rowtime=[$3])
+ +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable2]])
+]]>
+ </Resource>
+ <Resource name="optimized rel plan">
+ <![CDATA[
+Calc(select=[a, b], changelogMode=[I])
++- IntervalJoin(joinType=[LeftOuterJoin], windowBounds=[isRowTime=true,
leftLowerBound=10000, leftUpperBound=5000, leftTimeIndex=1, rightTimeIndex=2],
where=[AND(=(a, a0), >=(rowtime, +(rowtime0, 10000:INTERVAL SECOND)),
<=(rowtime, +(rowtime0, 5000:INTERVAL SECOND)))], select=[a, rowtime, a0, b,
rowtime0], earlyFireDelay=[5000], earlyFireTimeMode=[ROWTIME],
changelogMode=[I])
+ :- Exchange(distribution=[hash[a]], changelogMode=[I])
+ : +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime],
changelogMode=[I])
+ : +- TableSourceScan(table=[[default_catalog, default_database,
MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime],
changelogMode=[I])
+ +- Exchange(distribution=[hash[a]], changelogMode=[I])
+ +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime],
changelogMode=[I])
+ +- TableSourceScan(table=[[default_catalog, default_database,
MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime],
changelogMode=[I])
]]>
</Resource>
</TestCase>
@@ -145,6 +209,38 @@ Calc(select=[a, b])
+- Exchange(distribution=[hash[a]])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime])
+- TableSourceScan(table=[[default_catalog, default_database,
MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime])
+]]>
+ </Resource>
+ </TestCase>
+ <TestCase name="testEarlyFireOuterJoinProducesUpdates">
+ <Resource name="sql">
+ <![CDATA[SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b
+FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON
+ t1.a = t2.a AND
+ t1.rowtime BETWEEN t2.rowtime - INTERVAL '10' SECOND AND t2.rowtime +
INTERVAL '1' HOUR]]>
+ </Resource>
+ <Resource name="ast">
+ <![CDATA[
+LogicalProject(a=[$0], b=[$6])
++- LogicalJoin(condition=[AND(=($0, $5), >=($4, -($9, 10000:INTERVAL SECOND)),
<=($4, +($9, 3600000:INTERVAL HOUR)))], joinType=[left],
joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s}]]])
+ :- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+ : +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()],
rowtime=[$3])
+ : +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable]])
+ +- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+ +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()],
rowtime=[$3])
+ +- LogicalTableScan(table=[[default_catalog, default_database,
MyTable2]])
+]]>
+ </Resource>
+ <Resource name="optimized rel plan">
+ <![CDATA[
+Calc(select=[a, b], changelogMode=[I,UA])
++- IntervalJoin(joinType=[LeftOuterJoin], windowBounds=[isRowTime=true,
leftLowerBound=-10000, leftUpperBound=3600000, leftTimeIndex=1,
rightTimeIndex=2], where=[AND(=(a, a0), >=(rowtime, -(rowtime0, 10000:INTERVAL
SECOND)), <=(rowtime, +(rowtime0, 3600000:INTERVAL HOUR)))], select=[a,
rowtime, a0, b, rowtime0], earlyFireDelay=[5000], earlyFireTimeMode=[ROWTIME],
changelogMode=[I,UA])
+ :- Exchange(distribution=[hash[a]], changelogMode=[I])
+ : +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime],
changelogMode=[I])
+ : +- TableSourceScan(table=[[default_catalog, default_database,
MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime],
changelogMode=[I])
+ +- Exchange(distribution=[hash[a]], changelogMode=[I])
+ +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime],
changelogMode=[I])
+ +- TableSourceScan(table=[[default_catalog, default_database,
MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime],
changelogMode=[I])
]]>
</Resource>
</TestCase>
diff --git
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest_jsonplan/testEarlyFireJsonPlanRoundTrip.out
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest_jsonplan/testEarlyFireJsonPlanRoundTrip.out
index 2b98b28751e..8b6af42a64e 100644
---
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest_jsonplan/testEarlyFireJsonPlanRoundTrip.out
+++
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest_jsonplan/testEarlyFireJsonPlanRoundTrip.out
@@ -370,12 +370,13 @@
} ]
},
"options" : {
- "connector" : "values"
+ "connector" : "values",
+ "sink-insert-only" : "false"
}
}
}
},
- "inputChangelogMode" : [ "INSERT" ],
+ "inputChangelogMode" : [ "INSERT", "UPDATE_BEFORE", "UPDATE_AFTER" ],
"inputProperties" : [ {
"requiredDistribution" : {
"type" : "UNKNOWN"