Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<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()
Expand All @@ -60,6 +72,7 @@ public class EarlyFireJoinHintOptions {
static {
requiredKeys.add(DELAY);

supportedKeys.add(TARGET);
supportedKeys.add(DELAY);
supportedKeys.add(TIME_MODE);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,12 @@ private static EarlyFire extractEarlyFire(List<RelHint> 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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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"
+ ")");
}

Expand Down Expand Up @@ -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 =
Expand All @@ -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 =
Expand Down Expand Up @@ -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 =
Expand All @@ -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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,70 @@ 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="testEarlyFireInnerJoinStaysInsertOnly">
<Resource name="sql">
<![CDATA[SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b
FROM MyTable t1 JOIN MyTable2 t2 ON
t1.a = t2.a AND
t1.rowtime BETWEEN t2.rowtime - INTERVAL '10' SECOND AND t2.rowtime + INTERVAL '1' HOUR]]>
</Resource>
<Resource name="ast">
<![CDATA[
LogicalProject(a=[$0], b=[$6])
+- LogicalJoin(condition=[AND(=($0, $5), >=($4, -($9, 10000:INTERVAL SECOND)), <=($4, +($9, 3600000:INTERVAL HOUR)))], joinType=[inner], joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s}]]])
:- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
: +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()], rowtime=[$3])
: +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
+- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()], rowtime=[$3])
+- LogicalTableScan(table=[[default_catalog, default_database, MyTable2]])
]]>
</Resource>
<Resource name="optimized rel plan">
<![CDATA[
Calc(select=[a, b], changelogMode=[I])
+- IntervalJoin(joinType=[InnerJoin], windowBounds=[isRowTime=true, leftLowerBound=-10000, leftUpperBound=3600000, leftTimeIndex=1, rightTimeIndex=2], where=[AND(=(a, a0), >=(rowtime, -(rowtime0, 10000:INTERVAL SECOND)), <=(rowtime, +(rowtime0, 3600000:INTERVAL HOUR)))], select=[a, rowtime, a0, b, rowtime0], earlyFireDelay=[5000], earlyFireTimeMode=[ROWTIME], changelogMode=[I])
:- Exchange(distribution=[hash[a]], changelogMode=[I])
: +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime], changelogMode=[I])
: +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime], changelogMode=[I])
+- Exchange(distribution=[hash[a]], changelogMode=[I])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime], changelogMode=[I])
+- TableSourceScan(table=[[default_catalog, default_database, MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime], changelogMode=[I])
]]>
</Resource>
</TestCase>
<TestCase name="testEarlyFireLowerCaseHintNamePreservesOptions">
<Resource name="sql">
<![CDATA[SELECT /*+ early_fire('delay'='5s', 'time-mode'='rowtime') */ t1.a, t2.b
Expand Down Expand Up @@ -45,6 +109,38 @@ Calc(select=[a, b])
+- Exchange(distribution=[hash[a]])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime])
+- TableSourceScan(table=[[default_catalog, default_database, MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime])
]]>
</Resource>
</TestCase>
<TestCase name="testEarlyFireNegativeWindowStaysInsertOnly">
<Resource name="sql">
<![CDATA[SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b
FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON
t1.a = t2.a AND
t1.rowtime BETWEEN t2.rowtime + INTERVAL '10' SECOND AND t2.rowtime + INTERVAL '5' SECOND]]>
</Resource>
<Resource name="ast">
<![CDATA[
LogicalProject(a=[$0], b=[$6])
+- LogicalJoin(condition=[AND(=($0, $5), >=($4, +($9, 10000:INTERVAL SECOND)), <=($4, +($9, 5000:INTERVAL SECOND)))], joinType=[left], joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s}]]])
:- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
: +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()], rowtime=[$3])
: +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
+- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()], rowtime=[$3])
+- LogicalTableScan(table=[[default_catalog, default_database, MyTable2]])
]]>
</Resource>
<Resource name="optimized rel plan">
<![CDATA[
Calc(select=[a, b], changelogMode=[I])
+- IntervalJoin(joinType=[LeftOuterJoin], windowBounds=[isRowTime=true, leftLowerBound=10000, leftUpperBound=5000, leftTimeIndex=1, rightTimeIndex=2], where=[AND(=(a, a0), >=(rowtime, +(rowtime0, 10000:INTERVAL SECOND)), <=(rowtime, +(rowtime0, 5000:INTERVAL SECOND)))], select=[a, rowtime, a0, b, rowtime0], earlyFireDelay=[5000], earlyFireTimeMode=[ROWTIME], changelogMode=[I])
:- Exchange(distribution=[hash[a]], changelogMode=[I])
: +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime], changelogMode=[I])
: +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime], changelogMode=[I])
+- Exchange(distribution=[hash[a]], changelogMode=[I])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime], changelogMode=[I])
+- TableSourceScan(table=[[default_catalog, default_database, MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime], changelogMode=[I])
]]>
</Resource>
</TestCase>
Expand Down Expand Up @@ -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])
]]>
</Resource>
</TestCase>
<TestCase name="testEarlyFireOuterJoinProducesUpdates">
<Resource name="sql">
<![CDATA[SELECT /*+ EARLY_FIRE('delay'='5s') */ t1.a, t2.b
FROM MyTable t1 LEFT OUTER JOIN MyTable2 t2 ON
t1.a = t2.a AND
t1.rowtime BETWEEN t2.rowtime - INTERVAL '10' SECOND AND t2.rowtime + INTERVAL '1' HOUR]]>
</Resource>
<Resource name="ast">
<![CDATA[
LogicalProject(a=[$0], b=[$6])
+- LogicalJoin(condition=[AND(=($0, $5), >=($4, -($9, 10000:INTERVAL SECOND)), <=($4, +($9, 3600000:INTERVAL HOUR)))], joinType=[left], joinHints=[[[EARLY_FIRE inheritPath:[0] options:{delay=5s}]]])
:- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
: +- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()], rowtime=[$3])
: +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
+- LogicalWatermarkAssigner(rowtime=[rowtime], watermark=[$4])
+- LogicalProject(a=[$0], b=[$1], c=[$2], proctime=[PROCTIME()], rowtime=[$3])
+- LogicalTableScan(table=[[default_catalog, default_database, MyTable2]])
]]>
</Resource>
<Resource name="optimized rel plan">
<![CDATA[
Calc(select=[a, b], changelogMode=[I,UA])
+- IntervalJoin(joinType=[LeftOuterJoin], windowBounds=[isRowTime=true, leftLowerBound=-10000, leftUpperBound=3600000, leftTimeIndex=1, rightTimeIndex=2], where=[AND(=(a, a0), >=(rowtime, -(rowtime0, 10000:INTERVAL SECOND)), <=(rowtime, +(rowtime0, 3600000:INTERVAL HOUR)))], select=[a, rowtime, a0, b, rowtime0], earlyFireDelay=[5000], earlyFireTimeMode=[ROWTIME], changelogMode=[I,UA])
:- Exchange(distribution=[hash[a]], changelogMode=[I])
: +- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime], changelogMode=[I])
: +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a, rowtime], metadata=[]]], fields=[a, rowtime], changelogMode=[I])
+- Exchange(distribution=[hash[a]], changelogMode=[I])
+- WatermarkAssigner(rowtime=[rowtime], watermark=[rowtime], changelogMode=[I])
+- TableSourceScan(table=[[default_catalog, default_database, MyTable2, project=[a, b, rowtime], metadata=[]]], fields=[a, b, rowtime], changelogMode=[I])
]]>
</Resource>
</TestCase>
Expand Down
Loading