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