From 4215571399ad995ffb929306acb142cc10669ec2 Mon Sep 17 00:00:00 2001 From: weiqingy Date: Sat, 18 Jul 2026 15:09:44 -0700 Subject: [PATCH 1/6] [FLINK-36953][table] Add EARLY_FIRE join hint surface and option validation Register the EARLY_FIRE join hint and its typed options (EarlyFireJoinHintOptions: delay, time_mode) with a key-value option checker, and wire the query-hint propagation touchpoints. Exclude EARLY_FIRE from the generic join-hint test coverage and config-docs generation, since it is a key-value hint validated on its own. The hint is recognized and validated here but not yet consumed by any rule; threading it into the interval join follows in a separate change. --- .../docs/util/ConfigurationOptionLocator.java | 3 +- .../api/config/EarlyFireJoinHintOptions.java | 81 +++++++++++++++++++ .../hint/CapitalizeQueryHintsShuttle.java | 3 +- .../planner/hint/FlinkHintStrategies.java | 58 +++++++++++++ .../table/planner/hint/JoinStrategy.java | 13 +++ .../plan/optimize/QueryHintsResolver.java | 7 ++ .../plan/hints/batch/JoinHintTestBase.java | 5 +- .../stream/sql/join/IntervalJoinTest.scala | 57 +++++++++++++ 8 files changed, 223 insertions(+), 4 deletions(-) create mode 100644 flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/EarlyFireJoinHintOptions.java diff --git a/flink-docs/src/main/java/org/apache/flink/docs/util/ConfigurationOptionLocator.java b/flink-docs/src/main/java/org/apache/flink/docs/util/ConfigurationOptionLocator.java index 95dae776d7d837..ed2d987be504a9 100644 --- a/flink-docs/src/main/java/org/apache/flink/docs/util/ConfigurationOptionLocator.java +++ b/flink-docs/src/main/java/org/apache/flink/docs/util/ConfigurationOptionLocator.java @@ -105,7 +105,8 @@ public class ConfigurationOptionLocator { "org.apache.flink.state.rocksdb.PredefinedOptions", "org.apache.flink.python.PythonConfig", "org.apache.flink.cep.configuration.SharedBufferCacheConfig", - "org.apache.flink.table.api.config.LookupJoinHintOptions")); + "org.apache.flink.table.api.config.LookupJoinHintOptions", + "org.apache.flink.table.api.config.EarlyFireJoinHintOptions")); private static final String DEFAULT_PATH_PREFIX = "src/main/java"; 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 new file mode 100644 index 00000000000000..2e633eef1d7a28 --- /dev/null +++ b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/EarlyFireJoinHintOptions.java @@ -0,0 +1,81 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.api.config; + +import org.apache.flink.annotation.PublicEvolving; +import org.apache.flink.configuration.ConfigOption; + +import org.apache.flink.shaded.guava33.com.google.common.collect.ImmutableSet; + +import java.time.Duration; +import java.util.HashSet; +import java.util.Set; + +import static org.apache.flink.configuration.ConfigOptions.key; + +/** + * This class holds hint option name definitions for EARLY_FIRE join hints based on {@link + * org.apache.flink.configuration.ConfigOption}. + */ +@PublicEvolving +public class EarlyFireJoinHintOptions { + + public static final ConfigOption DELAY = + key("delay") + .durationType() + .noDefaultValue() + .withDescription( + "The delay between the time an unmatched outer row becomes eligible to" + + " be emitted with null padding and the time it is actually" + + " emitted. Must be a positive duration."); + + public static final ConfigOption TIME_MODE = + key("time_mode") + .enumType(TimeMode.class) + .noDefaultValue() + .withDescription( + "The time domain that drives the early-fire delay, can be 'rowtime' or" + + " 'proctime'. If not set, it defaults to the time domain of" + + " the interval join."); + + private static final Set> requiredKeys = new HashSet<>(); + private static final Set> supportedKeys = new HashSet<>(); + + static { + requiredKeys.add(DELAY); + + supportedKeys.add(DELAY); + supportedKeys.add(TIME_MODE); + } + + public static ImmutableSet getRequiredOptions() { + return ImmutableSet.copyOf(requiredKeys); + } + + public static ImmutableSet getSupportedOptions() { + return ImmutableSet.copyOf(supportedKeys); + } + + /** The time domain that drives the early-fire delay. */ + @PublicEvolving + public enum TimeMode { + ROWTIME, + PROCTIME + } +} diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/CapitalizeQueryHintsShuttle.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/CapitalizeQueryHintsShuttle.java index 006dab4b832a7f..6da4452010d380 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/CapitalizeQueryHintsShuttle.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/CapitalizeQueryHintsShuttle.java @@ -46,7 +46,8 @@ protected RelNode doVisit(RelNode node) { changed.set(true); if (JoinStrategy.isJoinStrategy(capitalHintName)) { - if (JoinStrategy.isLookupHint(hint.hintName)) { + if (JoinStrategy.isLookupHint(hint.hintName) + || JoinStrategy.isEarlyFireHint(hint.hintName)) { return RelHint.builder(capitalHintName) .hintOptions(hint.kvOptions) .inheritPath(hint.inheritPath) 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 5978f98017d91d..aa2ea0fd1372e9 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 @@ -20,6 +20,7 @@ import org.apache.flink.configuration.ConfigOption; import org.apache.flink.configuration.Configuration; +import org.apache.flink.table.api.config.EarlyFireJoinHintOptions; import org.apache.flink.table.api.config.LookupJoinHintOptions; import org.apache.flink.table.factories.FactoryUtil; import org.apache.flink.table.planner.plan.rules.logical.WrapJsonAggFunctionArgumentsRule; @@ -36,6 +37,9 @@ import java.time.Duration; import java.util.Collections; import java.util.Optional; +import java.util.Set; +import java.util.TreeSet; +import java.util.stream.Collectors; /** * A collection of Flink style {@link HintStrategy}s. @@ -135,6 +139,11 @@ public static HintStrategyTable createHintStrategyTable() { HintPredicates.JOIN, HintPredicates.AGGREGATE)) .optionChecker(STATE_TTL_NON_EMPTY_KV_OPTION_CHECKER) .build()) + .hintStrategy( + JoinStrategy.EARLY_FIRE.getJoinHintName(), + HintStrategy.builder(HintPredicates.JOIN) + .optionChecker(EARLY_FIRE_KV_OPTION_CHECKER) + .build()) .build(); } @@ -253,6 +262,55 @@ private static HintOptionChecker fixedSizeListOptionChecker(int size) { return true; }; + private static final HintOptionChecker EARLY_FIRE_KV_OPTION_CHECKER = + (earlyFireHint, litmus) -> { + litmus.check( + earlyFireHint.listOptions.size() == 0, + "Invalid list options in EARLY_FIRE hint, only support key-value options."); + + Configuration conf = Configuration.fromMap(earlyFireHint.kvOptions); + ImmutableSet requiredKeys = + EarlyFireJoinHintOptions.getRequiredOptions(); + litmus.check( + requiredKeys.stream().allMatch(conf::contains), + "Invalid EARLY_FIRE hint: incomplete required option(s): {}", + requiredKeys); + + ImmutableSet supportedKeys = + EarlyFireJoinHintOptions.getSupportedOptions(); + Set supportedKeyNames = + supportedKeys.stream().map(ConfigOption::key).collect(Collectors.toSet()); + Set unknownKeys = + earlyFireHint.kvOptions.keySet().stream() + .filter(key -> !supportedKeyNames.contains(key)) + .collect(Collectors.toCollection(TreeSet::new)); + litmus.check( + unknownKeys.isEmpty(), + "Unsupported EARLY_FIRE hint option(s) {}, supported options are {}.", + unknownKeys, + new TreeSet<>(supportedKeyNames)); + litmus.check( + earlyFireHint.kvOptions.size() <= supportedKeys.size(), + "Too many EARLY_FIRE hint options {} beyond max number of supported options {}", + earlyFireHint.kvOptions.size(), + supportedKeys.size()); + + try { + // try to validate all hint options by parsing them + supportedKeys.forEach(conf::get); + } catch (IllegalArgumentException e) { + litmus.fail("Invalid EARLY_FIRE hint options: {}", e.getMessage()); + } + + Duration delay = conf.get(EarlyFireJoinHintOptions.DELAY); + litmus.check( + null != delay && delay.toMillis() > 0, + "Invalid EARLY_FIRE hint option: {} value should be a positive duration but was {}", + EarlyFireJoinHintOptions.DELAY.key(), + delay); + return true; + }; + private static final HintOptionChecker STATE_TTL_NON_EMPTY_KV_OPTION_CHECKER = (ttlHint, litmus) -> { litmus.check( diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/JoinStrategy.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/JoinStrategy.java index 0ded47cb22230e..d9e3acff7a248b 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/JoinStrategy.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/JoinStrategy.java @@ -50,6 +50,12 @@ public enum JoinStrategy { /** Instructs the optimizer to use lookup join strategy. Only accept key-value hint options. */ LOOKUP("LOOKUP"), + /** + * Instructs an outer interval join to emit unmatched outer rows with null padding after a + * configurable delay. Only accept key-value hint options. + */ + EARLY_FIRE("EARLY_FIRE"), + /** * Instructs the optimizer to use multi-way join strategy for streaming queries. This hint * allows specifying multiple tables to be joined together in a single {@link @@ -89,6 +95,7 @@ public static boolean validOptions(String hintName, List options) { case NEST_LOOP: return options.size() > 0; case LOOKUP: + case EARLY_FIRE: return null == options || options.size() == 0; case MULTI_JOIN: return options.size() > 0; @@ -101,4 +108,10 @@ public static boolean isLookupHint(String hintName) { return isJoinStrategy(formalizedHintName) && JoinStrategy.valueOf(formalizedHintName) == LOOKUP; } + + public static boolean isEarlyFireHint(String hintName) { + String formalizedHintName = hintName.toUpperCase(Locale.ROOT); + return isJoinStrategy(formalizedHintName) + && JoinStrategy.valueOf(formalizedHintName) == EARLY_FIRE; + } } diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/optimize/QueryHintsResolver.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/optimize/QueryHintsResolver.java index b288f111e30cc4..b153ef2b885514 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/optimize/QueryHintsResolver.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/optimize/QueryHintsResolver.java @@ -146,6 +146,13 @@ private List validateAndGetNewHints( updateInfoForOptionCheck(hint.hintName, rightName); newHints.add(hint); } + } else if (JoinStrategy.isEarlyFireHint(hint.hintName)) { + // EARLY_FIRE carries only key-value options and is not bound to a specific input + // side, so it is passed through unchanged once its options are validated by the + // hint option checker. + allHints.add(trimInheritPath(hint)); + validHints.add(trimInheritPath(hint)); + newHints.add(hint); } else if (JoinStrategy.isJoinStrategy(hint.hintName)) { allHints.add(trimInheritPath(hint)); // add options about this hint for finally checking diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java index 24ae0e4025d08d..217e3c2b8af7e1 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java @@ -61,8 +61,9 @@ public abstract class JoinHintTestBase extends TableTestBase { private final List allJoinHintNames = Lists.newArrayList(JoinStrategy.values()).stream() - // LOOKUP hint has different kv-options against other join hints - .filter(hint -> hint != JoinStrategy.LOOKUP) + // LOOKUP and EARLY_FIRE hints only support key-value options, unlike the + // list-option join hints exercised here + .filter(hint -> hint != JoinStrategy.LOOKUP && hint != JoinStrategy.EARLY_FIRE) .map(JoinStrategy::getJoinHintName) .collect(Collectors.toList()); diff --git a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala index bf96e146bb74ba..c29b9756c16357 100644 --- a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala +++ b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala @@ -395,6 +395,63 @@ class IntervalJoinTest extends TableTestBase { util.verifyExecPlan(sqlQuery) } + // Tests for the EARLY_FIRE join hint + @Test + def testEarlyFireMissingDelay(): Unit = { + val sqlQuery = + """ + |SELECT /*+ EARLY_FIRE('time_mode'='rowtime') */ 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 + """.stripMargin + + assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) + .hasMessageContaining("incomplete required option(s)") + } + + @Test + def testEarlyFireNonPositiveDelay(): Unit = { + val sqlQuery = + """ + |SELECT /*+ EARLY_FIRE('delay'='0s') */ 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 + """.stripMargin + + assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) + .hasMessageContaining("value should be a positive duration") + } + + @Test + def testEarlyFireInvalidTimeMode(): Unit = { + val sqlQuery = + """ + |SELECT /*+ EARLY_FIRE('delay'='5s', 'time_mode'='unknown') */ 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 + """.stripMargin + + assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) + .hasMessageContaining("Invalid EARLY_FIRE hint options") + } + + @Test + def testEarlyFireUnknownOption(): Unit = { + val sqlQuery = + """ + |SELECT /*+ EARLY_FIRE('delay'='5s', 'timemode'='proctime') */ 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 + """.stripMargin + + assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) + .hasMessageContaining("Unsupported EARLY_FIRE hint option(s) [timemode]") + } + // Other tests @Test def testJoinTimeBoundary(): Unit = { From d64c58a70f6a548b564a5623bd6b03745fdae37f Mon Sep 17 00:00:00 2001 From: weiqingy Date: Mon, 20 Jul 2026 18:46:10 -0700 Subject: [PATCH 2/6] [FLINK-40167][table] Address review: rename time_mode hint option to time-mode Flink option keys use hyphens (e.g. output-mode, fixed-delay); rename the EARLY_FIRE hint's time_mode option key to time-mode to match, and update the affected tests. --- .../flink/table/api/config/EarlyFireJoinHintOptions.java | 2 +- .../table/planner/plan/stream/sql/join/IntervalJoinTest.scala | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) 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 2e633eef1d7a28..df77cc3167779d 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 @@ -46,7 +46,7 @@ public class EarlyFireJoinHintOptions { + " emitted. Must be a positive duration."); public static final ConfigOption TIME_MODE = - key("time_mode") + key("time-mode") .enumType(TimeMode.class) .noDefaultValue() .withDescription( diff --git a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala index c29b9756c16357..95522ba9857658 100644 --- a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala +++ b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.scala @@ -400,7 +400,7 @@ class IntervalJoinTest extends TableTestBase { def testEarlyFireMissingDelay(): Unit = { val sqlQuery = """ - |SELECT /*+ EARLY_FIRE('time_mode'='rowtime') */ t1.a, t2.b + |SELECT /*+ EARLY_FIRE('time-mode'='rowtime') */ 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 @@ -428,7 +428,7 @@ class IntervalJoinTest extends TableTestBase { def testEarlyFireInvalidTimeMode(): Unit = { val sqlQuery = """ - |SELECT /*+ EARLY_FIRE('delay'='5s', 'time_mode'='unknown') */ t1.a, t2.b + |SELECT /*+ EARLY_FIRE('delay'='5s', 'time-mode'='unknown') */ 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 From a0d69c9ae7280ad694b01595142b7a2f803579fa Mon Sep 17 00:00:00 2001 From: weiqingy Date: Mon, 20 Jul 2026 21:25:10 -0700 Subject: [PATCH 3/6] [FLINK-40167][table] Address review: require delay >= 1ms and add validation tests The runtime uses delay.toMillis(), so a sub-millisecond delay truncates to zero. Make the contract explicit: require at least 1 millisecond in the DELAY description and the checker error message. Add coverage for a sub-millisecond delay, list-style options (only key-value is supported), and case-insensitive hint-name capitalization that preserves the key-value options. --- .../api/config/EarlyFireJoinHintOptions.java | 2 +- .../planner/hint/FlinkHintStrategies.java | 2 +- .../plan/stream/sql/join/IntervalJoinTest.xml | 30 +++++++++++++ .../stream/sql/join/IntervalJoinTest.scala | 43 ++++++++++++++++++- 4 files changed, 74 insertions(+), 3 deletions(-) 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 df77cc3167779d..6c919f318f2469 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 @@ -43,7 +43,7 @@ public class EarlyFireJoinHintOptions { .withDescription( "The delay between the time an unmatched outer row becomes eligible to" + " be emitted with null padding and the time it is actually" - + " emitted. Must be a positive duration."); + + " emitted. Must be at least 1 millisecond."); public static final ConfigOption TIME_MODE = key("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 aa2ea0fd1372e9..204c23a97d84a5 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 @@ -305,7 +305,7 @@ private static HintOptionChecker fixedSizeListOptionChecker(int size) { Duration delay = conf.get(EarlyFireJoinHintOptions.DELAY); litmus.check( null != delay && delay.toMillis() > 0, - "Invalid EARLY_FIRE hint option: {} value should be a positive duration but was {}", + "Invalid EARLY_FIRE hint option: {} value should be at least 1 millisecond but was {}", EarlyFireJoinHintOptions.DELAY.key(), delay); return true; diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml index 680bc4a3bbb8bd..11770b8596f180 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml @@ -16,6 +16,36 @@ See the License for the specific language governing permissions and limitations under the License. --> + + + + + + =($4, -($9, 10000:INTERVAL SECOND)), <=($4, +($9, 3600000:INTERVAL HOUR)))], joinType=[left], joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s, time-mode=rowtime}]]]) + :- LogicalTableScan(table=[[default_catalog, default_database, MyTable]], hints=[[[ALIAS inheritPath:[] options:[t1]]]]) + +- LogicalTableScan(table=[[default_catalog, default_database, MyTable2]], hints=[[[ALIAS inheritPath:[] options:[t2]]]]) +]]> + + + = (rowtime0 - 10000:INTERVAL SECOND)) AND (rowtime <= (rowtime0 + 3600000:INTERVAL HOUR)))], select=[a, rowtime, a0, b, rowtime0]) + :- Exchange(distribution=[hash[a]]) + : +- Calc(select=[a, rowtime]) + : +- DataStreamScan(table=[[default_catalog, default_database, MyTable]], fields=[a, b, c, proctime, rowtime]) + +- Exchange(distribution=[hash[a]]) + +- Calc(select=[a, b, rowtime]) + +- DataStreamScan(table=[[default_catalog, default_database, MyTable2]], fields=[a, b, c, proctime, rowtime]) +]]> + + util.verifyExecPlan(sqlQuery)) - .hasMessageContaining("value should be a positive duration") + .hasMessageContaining("value should be at least 1 millisecond") } @Test @@ -452,6 +452,47 @@ class IntervalJoinTest extends TableTestBase { .hasMessageContaining("Unsupported EARLY_FIRE hint option(s) [timemode]") } + @Test + def testEarlyFireSubMillisecondDelay(): Unit = { + val sqlQuery = + """ + |SELECT /*+ EARLY_FIRE('delay'='1ns') */ 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 + """.stripMargin + + assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) + .hasMessageContaining("value should be at least 1 millisecond") + } + + @Test + def testEarlyFireListOptionsRejected(): Unit = { + val sqlQuery = + """ + |SELECT /*+ EARLY_FIRE('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 + """.stripMargin + + assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) + .hasMessageContaining("only support key-value options") + } + + @Test + def testEarlyFireLowerCaseHintNamePreservesOptions(): Unit = { + val sqlQuery = + """ + |SELECT /*+ early_fire('delay'='5s', 'time-mode'='rowtime') */ 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 + """.stripMargin + + util.verifyExecPlan(sqlQuery) + } + // Other tests @Test def testJoinTimeBoundary(): Unit = { From f6af6637194dacc7a5694e0f4c625c87fc0c1ad3 Mon Sep 17 00:00:00 2001 From: weiqingy Date: Tue, 21 Jul 2026 11:07:30 -0700 Subject: [PATCH 4/6] [FLINK-40167][table] Address review: move EARLY_FIRE tests to a Java test class Move the EARLY_FIRE hint surface and option-validation tests out of the Scala IntervalJoinTest into a dedicated Java EarlyFireJoinHintTest (under plan/hints/stream, mirroring StateTtlHintTest), and relocate their plan golden accordingly. IntervalJoinTest keeps only its original tests. --- .../hints/stream/EarlyFireJoinHintTest.java | 152 ++++++++++++++++++ .../hints/stream/EarlyFireJoinHintTest.xml | 51 ++++++ .../plan/stream/sql/join/IntervalJoinTest.xml | 30 ---- .../stream/sql/join/IntervalJoinTest.scala | 98 ----------- 4 files changed, 203 insertions(+), 128 deletions(-) create mode 100644 flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.java create mode 100644 flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.xml 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 new file mode 100644 index 00000000000000..47685e380d4f75 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.java @@ -0,0 +1,152 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.planner.plan.hints.stream; + +import org.apache.flink.table.api.ExplainDetail; +import org.apache.flink.table.api.TableConfig; +import org.apache.flink.table.planner.utils.PlanKind; +import org.apache.flink.table.planner.utils.StreamTableTestUtil; +import org.apache.flink.table.planner.utils.TableTestBase; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import scala.Enumeration; + +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Test for the EARLY_FIRE join hint surface and option validation. */ +class EarlyFireJoinHintTest extends TableTestBase { + + protected StreamTableTestUtil util; + + @BeforeEach + void before() { + util = streamTestUtil(TableConfig.getDefault()); + util.tableEnv() + .executeSql( + "CREATE TABLE MyTable (\n" + + " a INT,\n" + + " b VARCHAR,\n" + + " c BIGINT,\n" + + " proctime AS PROCTIME(),\n" + + " rowtime TIMESTAMP(3),\n" + + " WATERMARK FOR rowtime AS rowtime\n" + + ") WITH (\n" + + " 'connector' = 'values',\n" + + " 'bounded' = 'false'\n" + + ")"); + util.tableEnv() + .executeSql( + "CREATE TABLE MyTable2 (\n" + + " a INT,\n" + + " b VARCHAR,\n" + + " c BIGINT,\n" + + " proctime AS PROCTIME(),\n" + + " rowtime TIMESTAMP(3),\n" + + " WATERMARK FOR rowtime AS rowtime\n" + + ") WITH (\n" + + " 'connector' = 'values',\n" + + " 'bounded' = 'false'\n" + + ")"); + } + + @Test + void testEarlyFireMissingDelay() { + String sql = + "SELECT /*+ EARLY_FIRE('time-mode'='rowtime') */ 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("incomplete required option(s)"); + } + + @Test + void testEarlyFireNonPositiveDelay() { + String sql = + "SELECT /*+ EARLY_FIRE('delay'='0s') */ 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("value should be at least 1 millisecond"); + } + + @Test + void testEarlyFireSubMillisecondDelay() { + String sql = + "SELECT /*+ EARLY_FIRE('delay'='1ns') */ 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("value should be at least 1 millisecond"); + } + + @Test + void testEarlyFireInvalidTimeMode() { + String sql = + "SELECT /*+ EARLY_FIRE('delay'='5s', 'time-mode'='unknown') */ 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("Invalid EARLY_FIRE hint options"); + } + + @Test + void testEarlyFireUnknownOption() { + String sql = + "SELECT /*+ EARLY_FIRE('delay'='5s', 'timemode'='proctime') */ 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("Unsupported EARLY_FIRE hint option(s) [timemode]"); + } + + @Test + void testEarlyFireListOptionsRejected() { + String sql = + "SELECT /*+ EARLY_FIRE('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("only support key-value options"); + } + + @Test + void testEarlyFireLowerCaseHintNamePreservesOptions() { + String sql = + "SELECT /*+ early_fire('delay'='5s', 'time-mode'='rowtime') */ 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); + } + + private void verify(String sql) { + util.doVerifyPlan( + sql, + new ExplainDetail[] {}, + false, + new Enumeration.Value[] {PlanKind.AST(), PlanKind.OPT_EXEC()}, + false); + } +} 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 new file mode 100644 index 00000000000000..df5bc8675e72a6 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest.xml @@ -0,0 +1,51 @@ + + + + + + + + + =($4, -($9, 10000:INTERVAL SECOND)), <=($4, +($9, 3600000:INTERVAL HOUR)))], joinType=[left], joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s, time-mode=rowtime}]]]) + :- 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]]) +]]> + + + = (rowtime0 - 10000:INTERVAL SECOND)) AND (rowtime <= (rowtime0 + 3600000:INTERVAL HOUR)))], select=[a, rowtime, a0, b, rowtime0]) + :- 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]) +]]> + + + diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml index 11770b8596f180..680bc4a3bbb8bd 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/join/IntervalJoinTest.xml @@ -16,36 +16,6 @@ See the License for the specific language governing permissions and limitations under the License. --> - - - - - - =($4, -($9, 10000:INTERVAL SECOND)), <=($4, +($9, 3600000:INTERVAL HOUR)))], joinType=[left], joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s, time-mode=rowtime}]]]) - :- LogicalTableScan(table=[[default_catalog, default_database, MyTable]], hints=[[[ALIAS inheritPath:[] options:[t1]]]]) - +- LogicalTableScan(table=[[default_catalog, default_database, MyTable2]], hints=[[[ALIAS inheritPath:[] options:[t2]]]]) -]]> - - - = (rowtime0 - 10000:INTERVAL SECOND)) AND (rowtime <= (rowtime0 + 3600000:INTERVAL HOUR)))], select=[a, rowtime, a0, b, rowtime0]) - :- Exchange(distribution=[hash[a]]) - : +- Calc(select=[a, rowtime]) - : +- DataStreamScan(table=[[default_catalog, default_database, MyTable]], fields=[a, b, c, proctime, rowtime]) - +- Exchange(distribution=[hash[a]]) - +- Calc(select=[a, b, rowtime]) - +- DataStreamScan(table=[[default_catalog, default_database, MyTable2]], fields=[a, b, c, proctime, rowtime]) -]]> - - util.verifyExecPlan(sqlQuery)) - .hasMessageContaining("incomplete required option(s)") - } - - @Test - def testEarlyFireNonPositiveDelay(): Unit = { - val sqlQuery = - """ - |SELECT /*+ EARLY_FIRE('delay'='0s') */ 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 - """.stripMargin - - assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) - .hasMessageContaining("value should be at least 1 millisecond") - } - - @Test - def testEarlyFireInvalidTimeMode(): Unit = { - val sqlQuery = - """ - |SELECT /*+ EARLY_FIRE('delay'='5s', 'time-mode'='unknown') */ 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 - """.stripMargin - - assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) - .hasMessageContaining("Invalid EARLY_FIRE hint options") - } - - @Test - def testEarlyFireUnknownOption(): Unit = { - val sqlQuery = - """ - |SELECT /*+ EARLY_FIRE('delay'='5s', 'timemode'='proctime') */ 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 - """.stripMargin - - assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) - .hasMessageContaining("Unsupported EARLY_FIRE hint option(s) [timemode]") - } - - @Test - def testEarlyFireSubMillisecondDelay(): Unit = { - val sqlQuery = - """ - |SELECT /*+ EARLY_FIRE('delay'='1ns') */ 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 - """.stripMargin - - assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) - .hasMessageContaining("value should be at least 1 millisecond") - } - - @Test - def testEarlyFireListOptionsRejected(): Unit = { - val sqlQuery = - """ - |SELECT /*+ EARLY_FIRE('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 - """.stripMargin - - assertThatThrownBy(() => util.verifyExecPlan(sqlQuery)) - .hasMessageContaining("only support key-value options") - } - - @Test - def testEarlyFireLowerCaseHintNamePreservesOptions(): Unit = { - val sqlQuery = - """ - |SELECT /*+ early_fire('delay'='5s', 'time-mode'='rowtime') */ 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 - """.stripMargin - - util.verifyExecPlan(sqlQuery) - } - // Other tests @Test def testJoinTimeBoundary(): Unit = { From 229dcb514d4efcb5647b74eb87db65cd8f9f2682 Mon Sep 17 00:00:00 2001 From: weiqingy Date: Tue, 21 Jul 2026 13:28:01 -0700 Subject: [PATCH 5/6] [FLINK-40167][table] Fix Spotless line wrapping in EarlyFireJoinHintTest Wrap the assertThatThrownBy chain in testEarlyFireListOptionsRejected to satisfy Spotless. --- .../table/planner/plan/hints/stream/EarlyFireJoinHintTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 47685e380d4f75..447c198eef22ed 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 @@ -128,7 +128,8 @@ void testEarlyFireListOptionsRejected() { + "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("only support key-value options"); + assertThatThrownBy(() -> verify(sql)) + .hasMessageContaining("only support key-value options"); } @Test From 647116f42b1cc933c0d3a792292df77585cb3761 Mon Sep 17 00:00:00 2001 From: weiqingy Date: Sat, 18 Jul 2026 15:11:00 -0700 Subject: [PATCH 6/6] [FLINK-40168][table] Thread the EARLY_FIRE hint into the interval join --- .../exec/stream/StreamExecIntervalJoin.java | 25 + .../StreamPhysicalIntervalJoinRule.java | 59 ++- .../stream/StreamPhysicalIntervalJoin.scala | 13 +- .../hints/stream/EarlyFireJoinHintTest.java | 60 +++ .../hints/stream/EarlyFireJoinHintTest.xml | 70 ++- .../testEarlyFireJsonPlanRoundTrip.out | 446 ++++++++++++++++++ 6 files changed, 669 insertions(+), 4 deletions(-) create mode 100644 flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest_jsonplan/testEarlyFireJsonPlanRoundTrip.out diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecIntervalJoin.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecIntervalJoin.java index 2ac0af781e0c66..20b676af78562e 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecIntervalJoin.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecIntervalJoin.java @@ -28,6 +28,7 @@ import org.apache.flink.streaming.api.transformations.TwoInputTransformation; import org.apache.flink.streaming.api.transformations.UnionTransformation; import org.apache.flink.table.api.TableException; +import org.apache.flink.table.api.config.EarlyFireJoinHintOptions; import org.apache.flink.table.api.config.ExecutionConfigOptions; import org.apache.flink.table.data.RowData; import org.apache.flink.table.planner.delegation.PlannerBase; @@ -59,11 +60,14 @@ import org.apache.flink.shaded.guava33.com.google.common.collect.Lists; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonCreator; +import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonInclude; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import javax.annotation.Nullable; + import java.util.List; /** {@link StreamExecNode} for a time interval stream join. */ @@ -91,13 +95,27 @@ public class StreamExecIntervalJoin extends ExecNodeBase public static final String INTERVAL_JOIN_TRANSFORMATION = "interval-join"; public static final String FIELD_NAME_INTERVAL_JOIN_SPEC = "intervalJoinSpec"; + public static final String FIELD_NAME_EARLY_FIRE_DELAY = "earlyFireDelay"; + public static final String FIELD_NAME_EARLY_FIRE_TIME_MODE = "earlyFireTimeMode"; @JsonProperty(FIELD_NAME_INTERVAL_JOIN_SPEC) private final IntervalJoinSpec intervalJoinSpec; + @Nullable + @JsonProperty(FIELD_NAME_EARLY_FIRE_DELAY) + @JsonInclude(JsonInclude.Include.NON_NULL) + private final Long earlyFireDelay; + + @Nullable + @JsonProperty(FIELD_NAME_EARLY_FIRE_TIME_MODE) + @JsonInclude(JsonInclude.Include.NON_NULL) + private final EarlyFireJoinHintOptions.TimeMode earlyFireTimeMode; + public StreamExecIntervalJoin( ReadableConfig tableConfig, IntervalJoinSpec intervalJoinSpec, + @Nullable Long earlyFireDelay, + @Nullable EarlyFireJoinHintOptions.TimeMode earlyFireTimeMode, InputProperty leftInputProperty, InputProperty rightInputProperty, RowType outputType, @@ -107,6 +125,8 @@ public StreamExecIntervalJoin( ExecNodeContext.newContext(StreamExecIntervalJoin.class), ExecNodeContext.newPersistedConfig(StreamExecIntervalJoin.class, tableConfig), intervalJoinSpec, + earlyFireDelay, + earlyFireTimeMode, Lists.newArrayList(leftInputProperty, rightInputProperty), outputType, description); @@ -118,12 +138,17 @@ public StreamExecIntervalJoin( @JsonProperty(FIELD_NAME_TYPE) ExecNodeContext context, @JsonProperty(FIELD_NAME_CONFIGURATION) ReadableConfig persistedConfig, @JsonProperty(FIELD_NAME_INTERVAL_JOIN_SPEC) IntervalJoinSpec intervalJoinSpec, + @Nullable @JsonProperty(FIELD_NAME_EARLY_FIRE_DELAY) Long earlyFireDelay, + @Nullable @JsonProperty(FIELD_NAME_EARLY_FIRE_TIME_MODE) + EarlyFireJoinHintOptions.TimeMode earlyFireTimeMode, @JsonProperty(FIELD_NAME_INPUT_PROPERTIES) List inputProperties, @JsonProperty(FIELD_NAME_OUTPUT_TYPE) RowType outputType, @JsonProperty(FIELD_NAME_DESCRIPTION) String description) { super(id, context, persistedConfig, inputProperties, outputType, description); Preconditions.checkArgument(inputProperties.size() == 2); this.intervalJoinSpec = Preconditions.checkNotNull(intervalJoinSpec); + this.earlyFireDelay = earlyFireDelay; + this.earlyFireTimeMode = earlyFireTimeMode; } @Override 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 6354f04de07485..d7cbad27d0c5aa 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 @@ -19,9 +19,13 @@ package org.apache.flink.table.planner.plan.rules.physical.stream; import org.apache.flink.api.java.tuple.Tuple2; +import org.apache.flink.configuration.Configuration; import org.apache.flink.table.api.TableException; import org.apache.flink.table.api.ValidationException; +import org.apache.flink.table.api.config.EarlyFireJoinHintOptions; +import org.apache.flink.table.api.config.EarlyFireJoinHintOptions.TimeMode; import org.apache.flink.table.planner.calcite.FlinkTypeFactory; +import org.apache.flink.table.planner.hint.JoinStrategy; import org.apache.flink.table.planner.plan.nodes.FlinkRelNode; import org.apache.flink.table.planner.plan.nodes.exec.spec.IntervalJoinSpec; import org.apache.flink.table.planner.plan.nodes.logical.FlinkLogicalJoin; @@ -32,11 +36,16 @@ import org.apache.calcite.plan.RelOptRuleCall; import org.apache.calcite.plan.RelTraitSet; import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.hint.RelHint; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rex.RexNode; import org.immutables.value.Value; +import javax.annotation.Nullable; + +import java.time.Duration; import java.util.Collection; +import java.util.List; import java.util.function.Function; import java.util.stream.Collectors; @@ -133,6 +142,8 @@ public FlinkRelNode transform( RelTraitSet providedTraitSet) { Tuple2, Option> tuple2 = extractWindowBounds(join); + boolean isEventTime = tuple2.f0.get().isEventTime(); + EarlyFire earlyFire = extractEarlyFire(join.getHints(), isEventTime); return new StreamPhysicalIntervalJoin( join.getCluster(), providedTraitSet, @@ -141,7 +152,53 @@ public FlinkRelNode transform( join.getJoinType(), join.getCondition(), tuple2.f1.getOrElse(() -> join.getCluster().getRexBuilder().makeLiteral(true)), - tuple2.f0.get()); + tuple2.f0.get(), + earlyFire.delay, + earlyFire.timeMode); + } + + private static EarlyFire extractEarlyFire(List hints, boolean isEventTime) { + RelHint earlyFireHint = null; + for (RelHint hint : hints) { + if (JoinStrategy.isEarlyFireHint(hint.hintName)) { + earlyFireHint = hint; + break; + } + } + if (earlyFireHint == null) { + return new EarlyFire(null, null); + } + + Configuration conf = Configuration.fromMap(earlyFireHint.kvOptions); + Duration delay = conf.get(EarlyFireJoinHintOptions.DELAY); + TimeMode timeMode = conf.get(EarlyFireJoinHintOptions.TIME_MODE); + if (timeMode == null) { + timeMode = isEventTime ? TimeMode.ROWTIME : TimeMode.PROCTIME; + } + + if (!isEventTime && timeMode == TimeMode.ROWTIME) { + throw new ValidationException( + "EARLY_FIRE hint requested row-time triggering on a processing-time interval" + + " join. Row-time triggering requires a row-time interval join."); + } + if (isEventTime && timeMode == TimeMode.PROCTIME) { + // Processing-time triggering on an event-time interval join is not supported. + throw new TableException( + "EARLY_FIRE hint requested processing-time triggering on a row-time interval" + + " join, which is not yet supported."); + } + + return new EarlyFire(delay == null ? null : delay.toMillis(), timeMode); + } + + private static final class EarlyFire { + @Nullable private final Long delay; + @Nullable private final TimeMode timeMode; + + EarlyFire(@Nullable Long delay, @Nullable TimeMode timeMode) { + this.delay = delay; + this.timeMode = timeMode; + } } /** Configuration for {@link StreamPhysicalIntervalJoinRule}. */ 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 4916d653b2e236..d3401783989d9f 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 @@ -18,6 +18,7 @@ package org.apache.flink.table.planner.plan.nodes.physical.stream import org.apache.flink.table.api.TableException +import org.apache.flink.table.api.config.EarlyFireJoinHintOptions.TimeMode import org.apache.flink.table.planner.calcite.FlinkTypeFactory import org.apache.flink.table.planner.plan.nodes.exec.{ExecNode, InputProperty} import org.apache.flink.table.planner.plan.nodes.exec.spec.IntervalJoinSpec @@ -45,7 +46,9 @@ class StreamPhysicalIntervalJoin( val originalCondition: RexNode, // remaining join condition contains all of join condition except window bounds remainingCondition: RexNode, - windowBounds: WindowBounds) + windowBounds: WindowBounds, + earlyFireDelay: java.lang.Long, + earlyFireTimeMode: TimeMode) extends CommonPhysicalJoin(cluster, traitSet, leftRel, rightRel, remainingCondition, joinType) with StreamPhysicalRel { @@ -76,7 +79,9 @@ class StreamPhysicalIntervalJoin( joinType, originalCondition, conditionExpr, - windowBounds) + windowBounds, + earlyFireDelay, + earlyFireTimeMode) } override def explainTerms(pw: RelWriter): RelWriter = { @@ -98,12 +103,16 @@ class StreamPhysicalIntervalJoin( preferExpressionFormat(pw), pw.getDetailLevel)) .item("select", getRowType.getFieldNames.mkString(", ")) + .itemIf("earlyFireDelay", earlyFireDelay, earlyFireDelay != null) + .itemIf("earlyFireTimeMode", earlyFireTimeMode, earlyFireTimeMode != null) } override def translateToExecNode(): ExecNode[_] = { new StreamExecIntervalJoin( unwrapTableConfig(this), new IntervalJoinSpec(joinSpec, windowBounds), + earlyFireDelay, + earlyFireTimeMode, InputProperty.DEFAULT, InputProperty.DEFAULT, FlinkTypeFactory.toLogicalRowType(getRowType), 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 447c198eef22ed..88aac40078f8e2 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 @@ -65,6 +65,14 @@ void before() { + " 'connector' = 'values',\n" + " 'bounded' = 'false'\n" + ")"); + util.tableEnv() + .executeSql( + "CREATE TABLE MySink (\n" + + " a INT,\n" + + " b VARCHAR\n" + + ") WITH (\n" + + " 'connector' = 'values'\n" + + ")"); } @Test @@ -142,6 +150,58 @@ void testEarlyFireLowerCaseHintNamePreservesOptions() { verify(sql); } + @Test + void testEarlyFireOnRowTimeLeftOuterJoin() { + 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"; + verify(sql); + } + + @Test + void testEarlyFireRowTimeOnProcTimeJoin() { + String sql = + "SELECT /*+ EARLY_FIRE('delay'='5s', 'time-mode'='rowtime') */ t1.a, t2.b\n" + + "FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON\n" + + " t1.a = t2.a AND\n" + + " t1.proctime BETWEEN t2.proctime - INTERVAL '1' HOUR AND t2.proctime + INTERVAL '1' HOUR"; + assertThatThrownBy(() -> verify(sql)) + .hasStackTraceContaining("requires a row-time interval join"); + } + + @Test + void testEarlyFireProcTimeOnRowTimeJoin() { + String sql = + "SELECT /*+ EARLY_FIRE('delay'='5s', 'time-mode'='proctime') */ 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)).hasStackTraceContaining("not yet supported"); + } + + @Test + void testEarlyFireOnProcTimeLeftOuterJoin() { + 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.proctime BETWEEN t2.proctime - INTERVAL '1' HOUR AND t2.proctime + INTERVAL '1' HOUR"; + verify(sql); + } + + @Test + void testEarlyFireJsonPlanRoundTrip() { + String insert = + "INSERT INTO MySink\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"; + util.verifyJsonPlan(insert); + } + private void verify(String sql) { util.doVerifyPlan( 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 df5bc8675e72a6..3df60f3cb8d756 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 @@ -38,7 +38,75 @@ LogicalProject(a=[$0], b=[$6]) = (rowtime0 - 10000:INTERVAL SECOND)) AND (rowtime <= (rowtime0 + 3600000:INTERVAL HOUR)))], select=[a, rowtime, a0, b, rowtime0]) ++- 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]) +]]> + + + + + + + + =($3, -($8, 3600000:INTERVAL HOUR)), <=($3, +($8, 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]]) +]]> + + + = (proctime0 - 3600000:INTERVAL HOUR)) AND (proctime <= (proctime0 + 3600000:INTERVAL HOUR)))], select=[a, proctime, a0, b, proctime0], earlyFireDelay=[5000], earlyFireTimeMode=[PROCTIME]) + :- Exchange(distribution=[hash[a]]) + : +- Calc(select=[a, proctime]) + : +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime]) + : +- Calc(select=[a, PROCTIME() AS proctime, rowtime]) + : +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime]) + +- Exchange(distribution=[hash[a]]) + +- Calc(select=[a, b, proctime]) + +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime]) + +- Calc(select=[a, b, PROCTIME() AS proctime, rowtime]) + +- TableSourceScan(table=[[default_catalog, default_database, MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime]) +]]> + + + + + + + + =($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]]) +]]> + + + = (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]) 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 new file mode 100644 index 00000000000000..2b98b28751eb25 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/stream/EarlyFireJoinHintTest_jsonplan/testEarlyFireJsonPlanRoundTrip.out @@ -0,0 +1,446 @@ +{ + "flinkVersion" : "", + "nodes" : [ { + "id" : 1, + "type" : "stream-exec-table-source-scan_2", + "scanTableSource" : { + "table" : { + "identifier" : "`default_catalog`.`default_database`.`MyTable`", + "resolvedTable" : { + "schema" : { + "columns" : [ { + "name" : "a", + "dataType" : "INT" + }, { + "name" : "b", + "dataType" : "VARCHAR(2147483647)" + }, { + "name" : "c", + "dataType" : "BIGINT" + }, { + "name" : "proctime", + "kind" : "COMPUTED", + "expression" : { + "rexNode" : { + "kind" : "CALL", + "internalName" : "$PROCTIME$1", + "type" : { + "type" : "TIMESTAMP_WITH_LOCAL_TIME_ZONE", + "nullable" : false, + "precision" : 3, + "kind" : "PROCTIME" + } + }, + "serializableString" : "PROCTIME()" + } + }, { + "name" : "rowtime", + "dataType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + } ], + "watermarkSpecs" : [ { + "rowtimeAttribute" : "rowtime", + "expression" : { + "rexNode" : { + "kind" : "INPUT_REF", + "inputIndex" : 4, + "type" : "TIMESTAMP(3)" + }, + "serializableString" : "`rowtime`" + } + } ] + }, + "options" : { + "bounded" : "false", + "connector" : "values" + } + } + }, + "abilities" : [ { + "type" : "ProjectPushDown", + "projectedFields" : [ [ 0 ], [ 3 ] ], + "producedType" : "ROW<`a` INT, `rowtime` TIMESTAMP(3)> NOT NULL" + }, { + "type" : "ReadingMetadata", + "metadataKeys" : [ ], + "producedType" : "ROW<`a` INT, `rowtime` TIMESTAMP(3)> NOT NULL" + } ] + }, + "outputType" : "ROW<`a` INT, `rowtime` TIMESTAMP(3)>", + "description" : "TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime])" + }, { + "id" : 2, + "type" : "stream-exec-watermark-assigner_1", + "watermarkExpr" : { + "kind" : "INPUT_REF", + "inputIndex" : 1, + "type" : "TIMESTAMP(3)" + }, + "rowtimeFieldIndex" : 1, + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : { + "type" : "ROW", + "fields" : [ { + "name" : "a", + "fieldType" : "INT" + }, { + "name" : "rowtime", + "fieldType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + } ] + }, + "description" : "WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime])" + }, { + "id" : 3, + "type" : "stream-exec-exchange_1", + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "HASH", + "keys" : [ 0 ] + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : { + "type" : "ROW", + "fields" : [ { + "name" : "a", + "fieldType" : "INT" + }, { + "name" : "rowtime", + "fieldType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + } ] + }, + "description" : "Exchange(distribution=[hash[a]])" + }, { + "id" : 4, + "type" : "stream-exec-table-source-scan_2", + "scanTableSource" : { + "table" : { + "identifier" : "`default_catalog`.`default_database`.`MyTable2`", + "resolvedTable" : { + "schema" : { + "columns" : [ { + "name" : "a", + "dataType" : "INT" + }, { + "name" : "b", + "dataType" : "VARCHAR(2147483647)" + }, { + "name" : "c", + "dataType" : "BIGINT" + }, { + "name" : "proctime", + "kind" : "COMPUTED", + "expression" : { + "rexNode" : { + "kind" : "CALL", + "internalName" : "$PROCTIME$1", + "type" : { + "type" : "TIMESTAMP_WITH_LOCAL_TIME_ZONE", + "nullable" : false, + "precision" : 3, + "kind" : "PROCTIME" + } + }, + "serializableString" : "PROCTIME()" + } + }, { + "name" : "rowtime", + "dataType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + } ], + "watermarkSpecs" : [ { + "rowtimeAttribute" : "rowtime", + "expression" : { + "rexNode" : { + "kind" : "INPUT_REF", + "inputIndex" : 4, + "type" : "TIMESTAMP(3)" + }, + "serializableString" : "`rowtime`" + } + } ] + }, + "options" : { + "bounded" : "false", + "connector" : "values" + } + } + }, + "abilities" : [ { + "type" : "ProjectPushDown", + "projectedFields" : [ [ 0 ], [ 1 ], [ 3 ] ], + "producedType" : "ROW<`a` INT, `b` VARCHAR(2147483647), `rowtime` TIMESTAMP(3)> NOT NULL" + }, { + "type" : "ReadingMetadata", + "metadataKeys" : [ ], + "producedType" : "ROW<`a` INT, `b` VARCHAR(2147483647), `rowtime` TIMESTAMP(3)> NOT NULL" + } ] + }, + "outputType" : "ROW<`a` INT, `b` VARCHAR(2147483647), `rowtime` TIMESTAMP(3)>", + "description" : "TableSourceScan(table=[[default_catalog, default_database, MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime])" + }, { + "id" : 5, + "type" : "stream-exec-watermark-assigner_1", + "watermarkExpr" : { + "kind" : "INPUT_REF", + "inputIndex" : 2, + "type" : "TIMESTAMP(3)" + }, + "rowtimeFieldIndex" : 2, + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : { + "type" : "ROW", + "fields" : [ { + "name" : "a", + "fieldType" : "INT" + }, { + "name" : "b", + "fieldType" : "VARCHAR(2147483647)" + }, { + "name" : "rowtime", + "fieldType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + } ] + }, + "description" : "WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime])" + }, { + "id" : 6, + "type" : "stream-exec-exchange_1", + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "HASH", + "keys" : [ 0 ] + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : { + "type" : "ROW", + "fields" : [ { + "name" : "a", + "fieldType" : "INT" + }, { + "name" : "b", + "fieldType" : "VARCHAR(2147483647)" + }, { + "name" : "rowtime", + "fieldType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + } ] + }, + "description" : "Exchange(distribution=[hash[a]])" + }, { + "id" : 7, + "type" : "stream-exec-interval-join_1", + "intervalJoinSpec" : { + "joinSpec" : { + "joinType" : "LEFT", + "leftKeys" : [ 0 ], + "rightKeys" : [ 0 ], + "filterNulls" : [ true ], + "nonEquiCondition" : null + }, + "windowBounds" : { + "isEventTime" : true, + "leftLowerBound" : -10000, + "leftUpperBound" : 3600000, + "leftTimeIndex" : 1, + "rightTimeIndex" : 2 + } + }, + "earlyFireDelay" : 5000, + "earlyFireTimeMode" : "ROWTIME", + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + }, { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : { + "type" : "ROW", + "fields" : [ { + "name" : "a", + "fieldType" : "INT" + }, { + "name" : "rowtime", + "fieldType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + }, { + "name" : "a0", + "fieldType" : "INT" + }, { + "name" : "b", + "fieldType" : "VARCHAR(2147483647)" + }, { + "name" : "rowtime0", + "fieldType" : { + "type" : "TIMESTAMP_WITHOUT_TIME_ZONE", + "precision" : 3, + "kind" : "ROWTIME" + } + } ] + }, + "description" : "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])" + }, { + "id" : 8, + "type" : "stream-exec-calc_1", + "projection" : [ { + "kind" : "INPUT_REF", + "inputIndex" : 0, + "type" : "INT" + }, { + "kind" : "INPUT_REF", + "inputIndex" : 3, + "type" : "VARCHAR(2147483647)" + } ], + "condition" : null, + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : "ROW<`a` INT, `b` VARCHAR(2147483647)>", + "description" : "Calc(select=[a, b])" + }, { + "id" : 9, + "type" : "stream-exec-sink_2", + "configuration" : { + "table.exec.sink.keyed-shuffle" : "AUTO", + "table.exec.sink.not-null-enforcer" : "ERROR", + "table.exec.sink.rowtime-inserter" : "ENABLED", + "table.exec.sink.type-length-enforcer" : "IGNORE", + "table.exec.sink.upsert-materialize" : "AUTO" + }, + "dynamicTableSink" : { + "table" : { + "identifier" : "`default_catalog`.`default_database`.`MySink`", + "resolvedTable" : { + "schema" : { + "columns" : [ { + "name" : "a", + "dataType" : "INT" + }, { + "name" : "b", + "dataType" : "VARCHAR(2147483647)" + } ] + }, + "options" : { + "connector" : "values" + } + } + } + }, + "inputChangelogMode" : [ "INSERT" ], + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : "ROW<`a` INT, `b` VARCHAR(2147483647)>", + "description" : "Sink(table=[default_catalog.default_database.MySink], fields=[a, b])" + } ], + "edges" : [ { + "source" : 1, + "target" : 2, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 2, + "target" : 3, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 4, + "target" : 5, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 5, + "target" : 6, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 3, + "target" : 7, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 6, + "target" : 7, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 7, + "target" : 8, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 8, + "target" : 9, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + } ] +} \ No newline at end of file