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 75d978ef1d9 [FLINK-40169][table] Add target option to the EARLY_FIRE 
hint (#28827)
75d978ef1d9 is described below

commit 75d978ef1d90cf16aacc6715390e1faa24a01256
Author: Weiqing Yang <[email protected]>
AuthorDate: Fri Aug 7 18:09:06 2026 -0700

    [FLINK-40169][table] Add target option to the EARLY_FIRE hint (#28827)
---
 .../table/api/config/EarlyFireJoinHintOptions.java | 13 +++++++++
 .../table/planner/hint/FlinkHintStrategies.java    |  8 ++++++
 .../stream/StreamPhysicalIntervalJoinRule.java     |  6 ++++
 .../plan/hints/stream/EarlyFireJoinHintTest.java   | 21 ++++++++++++++
 .../plan/hints/stream/EarlyFireJoinHintTest.xml    | 32 ++++++++++++++++++++++
 5 files changed, 80 insertions(+)

diff --git 
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/EarlyFireJoinHintOptions.java
 
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/EarlyFireJoinHintOptions.java
index 6c919f318f2..792d1b423bc 100644
--- 
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/EarlyFireJoinHintOptions.java
+++ 
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/EarlyFireJoinHintOptions.java
@@ -36,6 +36,18 @@ import static 
org.apache.flink.configuration.ConfigOptions.key;
 @PublicEvolving
 public class EarlyFireJoinHintOptions {
 
+    /** The only operator kind the EARLY_FIRE hint currently supports. */
+    public static final String INTERVAL_JOIN = "interval_join";
+
+    public static final ConfigOption<String> TARGET =
+            key("target")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "The operator kind that the EARLY_FIRE hint 
applies to. Currently only"
+                                    + " 'interval_join' is supported. When 
omitted, the hint applies"
+                                    + " to the interval join.");
+
     public static final ConfigOption<Duration> DELAY =
             key("delay")
                     .durationType()
@@ -60,6 +72,7 @@ public class EarlyFireJoinHintOptions {
     static {
         requiredKeys.add(DELAY);
 
+        supportedKeys.add(TARGET);
         supportedKeys.add(DELAY);
         supportedKeys.add(TIME_MODE);
     }
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/FlinkHintStrategies.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/FlinkHintStrategies.java
index 204c23a97d8..55018dec75f 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/FlinkHintStrategies.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/FlinkHintStrategies.java
@@ -308,6 +308,14 @@ public abstract class FlinkHintStrategies {
                         "Invalid EARLY_FIRE hint option: {} value should be at 
least 1 millisecond but was {}",
                         EarlyFireJoinHintOptions.DELAY.key(),
                         delay);
+
+                String target = conf.get(EarlyFireJoinHintOptions.TARGET);
+                litmus.check(
+                        null == target || 
EarlyFireJoinHintOptions.INTERVAL_JOIN.equals(target),
+                        "Invalid EARLY_FIRE hint option: {} value '{}' is not 
supported, only '{}' is supported currently",
+                        EarlyFireJoinHintOptions.TARGET.key(),
+                        target,
+                        EarlyFireJoinHintOptions.INTERVAL_JOIN);
                 return true;
             };
 
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/physical/stream/StreamPhysicalIntervalJoinRule.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/physical/stream/StreamPhysicalIntervalJoinRule.java
index d7cbad27d0c..5b848e4f20c 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/physical/stream/StreamPhysicalIntervalJoinRule.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/physical/stream/StreamPhysicalIntervalJoinRule.java
@@ -170,6 +170,12 @@ public class StreamPhysicalIntervalJoinRule
         }
 
         Configuration conf = Configuration.fromMap(earlyFireHint.kvOptions);
+        // target scopes the hint to one operator kind: this rule applies it 
only when it targets
+        // the interval join, and leaves a hint aimed at any other operator 
kind untouched.
+        String target = conf.get(EarlyFireJoinHintOptions.TARGET);
+        if (target != null && 
!EarlyFireJoinHintOptions.INTERVAL_JOIN.equals(target)) {
+            return new EarlyFire(null, null);
+        }
         Duration delay = conf.get(EarlyFireJoinHintOptions.DELAY);
         TimeMode timeMode = conf.get(EarlyFireJoinHintOptions.TIME_MODE);
         if (timeMode == null) {
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 88aac40078f..1db3dc6e8c2 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
@@ -140,6 +140,17 @@ class EarlyFireJoinHintTest extends TableTestBase {
                 .hasMessageContaining("only support key-value options");
     }
 
+    @Test
+    void testEarlyFireUnsupportedTarget() {
+        String sql =
+                "SELECT /*+ EARLY_FIRE('target'='window_join', '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(() -> verify(sql))
+                .hasMessageContaining("target value 'window_join' is not 
supported");
+    }
+
     @Test
     void testEarlyFireLowerCaseHintNamePreservesOptions() {
         String sql =
@@ -160,6 +171,16 @@ class EarlyFireJoinHintTest extends TableTestBase {
         verify(sql);
     }
 
+    @Test
+    void testEarlyFireExplicitTargetIntervalJoin() {
+        String sql =
+                "SELECT /*+ EARLY_FIRE('target'='interval_join', '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";
+        verify(sql);
+    }
+
     @Test
     void testEarlyFireRowTimeOnProcTimeJoin() {
         String sql =
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 3df60f3cb8d..8daae577559 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
@@ -16,6 +16,38 @@ See the License for the specific language governing 
permissions and
 limitations under the License.
 -->
 <Root>
+  <TestCase name="testEarlyFireExplicitTargetIntervalJoin">
+    <Resource name="sql">
+      <![CDATA[SELECT /*+ EARLY_FIRE('target'='interval_join', '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, 
target=interval_join}]]])
+   :- 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 exec plan">
+      <![CDATA[
+Calc(select=[a, b])
++- IntervalJoin(joinType=[LeftOuterJoin], windowBounds=[isRowTime=true, 
leftLowerBound=-10000, leftUpperBound=3600000, leftTimeIndex=1, 
rightTimeIndex=2], where=[((a = a0) AND (rowtime >= (rowtime0 - 10000:INTERVAL 
SECOND)) AND (rowtime <= (rowtime0 + 3600000:INTERVAL HOUR)))], select=[a, 
rowtime, a0, b, rowtime0], earlyFireDelay=[5000], earlyFireTimeMode=[ROWTIME])
+   :- Exchange(distribution=[hash[a]])
+   :  +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime])
+   :     +- TableSourceScan(table=[[default_catalog, default_database, 
MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime])
+   +- 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="testEarlyFireLowerCaseHintNamePreservesOptions">
     <Resource name="sql">
       <![CDATA[SELECT /*+ early_fire('delay'='5s', 'time-mode'='rowtime') */ 
t1.a, t2.b

Reply via email to