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 6c919f318f2469..f8ecb172bd7f44 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 @@ @PublicEvolving public class EarlyFireJoinHintOptions { + /** The only operator kind the EARLY_FIRE hint applies to. */ + public static final String INTERVAL_JOIN = "interval_join"; + + public static final ConfigOption 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 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 204c23a97d84a5..55018dec75fc9e 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 @@ private static HintOptionChecker fixedSizeListOptionChecker(int size) { "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 d7cbad27d0c5aa..5b848e4f20c2de 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 @@ private static EarlyFire extractEarlyFire(List hints, boolean isEventTi } 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/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 d3401783989d9f..0d63c0c50fd1b4 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) >= 0 + 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 1c91e960a8a180..867066dc0176d2 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 is insert-only, 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 88aac40078f8e2..05879701b33399 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.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 @@ void before() { + " a INT,\n" + " b VARCHAR\n" + ") WITH (\n" - + " 'connector' = 'values'\n" + + " 'connector' = 'values',\n" + + " 'sink-insert-only' = 'false'\n" + ")"); } @@ -140,6 +143,17 @@ void testEarlyFireListOptionsRejected() { .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 +174,16 @@ void testEarlyFireOnRowTimeLeftOuterJoin() { 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 = @@ -191,6 +215,58 @@ void testEarlyFireOnProcTimeLeftOuterJoin() { 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 = @@ -210,4 +286,8 @@ private void verify(String sql) { 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 3df60f3cb8d756..8faa17c8bea3fa 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,70 @@ 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, 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]]) +]]> + + + = (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]) +]]> + + + + + + + + =($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]]) +]]> + + + =(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]) +]]> + + + + + + + + + + =($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]]) +]]> + + + =(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]) ]]> @@ -113,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]) +]]> + + + + + + + + =($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]]) +]]> + + + =(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]) ]]> 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 2b98b28751eb25..8b6af42a64e720 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"