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"

Reply via email to