This is an automated email from the ASF dual-hosted git repository.
RocMarshal pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/master by this push:
new 2e224426850 [FLINK-40171][table-runtime] Emit and retract early-fire
results in the interval join operator (#28952)
2e224426850 is described below
commit 2e224426850258c9a342db9ea87fd625c6692744
Author: Weiqing Yang <[email protected]>
AuthorDate: Wed Aug 26 21:48:17 2026 -0700
[FLINK-40171][table-runtime] Emit and retract early-fire results in the
interval join operator (#28952)
---
.../nodes/exec/stream/StreamExecIntervalJoin.java | 6 +-
.../join/interval/EmitAwareCollector.java | 44 ++-
.../join/interval/ProcTimeIntervalJoin.java | 6 +-
.../join/interval/RowTimeIntervalJoin.java | 6 +-
.../operators/join/interval/TimeIntervalJoin.java | 310 ++++++++++++++++++---
.../join/interval/ProcTimeIntervalJoinTest.java | 84 +++++-
.../join/interval/RowTimeIntervalJoinTest.java | 298 +++++++++++++++++++-
7 files changed, 694 insertions(+), 60 deletions(-)
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 20b676af785..3d8d4c7b01b 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
@@ -391,7 +391,8 @@ public class StreamExecIntervalJoin extends
ExecNodeBase<RowData>
minCleanUpIntervalMillis,
leftTypeInfo,
rightTypeInfo,
- joinFunction);
+ joinFunction,
+ earlyFireDelay == null ? -1L : earlyFireDelay);
// TODO: add async version procJoinFunc to use
AsyncKeyedCoProcessOperator
return ExecNodeUtil.createTwoInputTransformation(
leftInputTransform,
@@ -428,7 +429,8 @@ public class StreamExecIntervalJoin extends
ExecNodeBase<RowData>
rightTypeInfo,
joinFunction,
windowBounds.getLeftTimeIdx(),
- windowBounds.getRightTimeIdx());
+ windowBounds.getRightTimeIdx(),
+ earlyFireDelay == null ? -1L : earlyFireDelay);
// TODO: add async version rowJoinFunc to use
AsyncKeyedCoProcessOperator
return ExecNodeUtil.createTwoInputTransformation(
leftInputTransform,
diff --git
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/EmitAwareCollector.java
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/EmitAwareCollector.java
index e9fe4447f55..a4e9c2f4652 100644
---
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/EmitAwareCollector.java
+++
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/EmitAwareCollector.java
@@ -19,17 +19,30 @@
package org.apache.flink.table.runtime.operators.join.interval;
import org.apache.flink.table.data.RowData;
+import org.apache.flink.types.RowKind;
import org.apache.flink.util.Collector;
/**
- * Collector to wrap a [[org.apache.flink.table.dataformat.RowData]] and to
track whether a row has
- * been emitted by the inner collector.
+ * Collector to wrap a {@link RowData} and to track whether a row has been
emitted by the inner
+ * collector.
+ *
+ * <p>The collector can be armed with a correction before a single matched row
is collected. When
+ * armed, the next collected row is treated as the corrected result of a
previously emitted
+ * speculative outer-join pad: the pending pad is emitted first stamped {@link
+ * RowKind#UPDATE_BEFORE}, then the matched row is stamped {@link
RowKind#UPDATE_AFTER}. This turns
+ * the join function's single {@code INSERT} emit into the {@code -U}/{@code
+U} pair without the
+ * join function knowing about changelogs. When not armed, the collected row
is stamped {@link
+ * RowKind#INSERT}, because the join function reuses a single row instance
whose kind may have been
+ * left at {@link RowKind#UPDATE_AFTER} by an earlier correction.
*/
class EmitAwareCollector implements Collector<RowData> {
private boolean emitted = false;
private Collector<RowData> innerCollector;
+ // The pad to retract before the next matched row, or null when no
correction is armed.
+ private RowData pendingRetraction;
+
void reset() {
emitted = false;
}
@@ -42,10 +55,35 @@ class EmitAwareCollector implements Collector<RowData> {
this.innerCollector = innerCollector;
}
+ /**
+ * Arms the collector so the next collected matched row is corrected into
a {@code -U}/{@code
+ * +U} pair against the given padded row.
+ */
+ void armRetraction(RowData retractionPad) {
+ retractionPad.setRowKind(RowKind.UPDATE_BEFORE);
+ this.pendingRetraction = retractionPad;
+ }
+
+ /** Clears an armed correction that was never consumed (the join condition
did not match). */
+ void disarm() {
+ this.pendingRetraction = null;
+ }
+
@Override
public void collect(RowData record) {
emitted = true;
- innerCollector.collect(record);
+ if (pendingRetraction != null) {
+ innerCollector.collect(pendingRetraction);
+ pendingRetraction = null;
+ record.setRowKind(RowKind.UPDATE_AFTER);
+ innerCollector.collect(record);
+ } else {
+ // The matched row reuses a single instance whose kind may have
been left as
+ // UPDATE_AFTER by a previous correction; force INSERT so a later
ordinary match is not
+ // mis-emitted as an update.
+ record.setRowKind(RowKind.INSERT);
+ innerCollector.collect(record);
+ }
}
@Override
diff --git
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoin.java
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoin.java
index 1fefc759fd4..84ad4526289 100644
---
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoin.java
+++
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoin.java
@@ -34,7 +34,8 @@ public final class ProcTimeIntervalJoin extends
TimeIntervalJoin {
long minCleanUpInterval,
InternalTypeInfo<RowData> leftType,
InternalTypeInfo<RowData> rightType,
- IntervalJoinFunction genJoinFunc) {
+ IntervalJoinFunction genJoinFunc,
+ long earlyFireDelay) {
super(
joinType,
leftLowerBound,
@@ -43,7 +44,8 @@ public final class ProcTimeIntervalJoin extends
TimeIntervalJoin {
minCleanUpInterval,
leftType,
rightType,
- genJoinFunc);
+ genJoinFunc,
+ earlyFireDelay);
}
@Override
diff --git
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoin.java
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoin.java
index 5d972104a66..57972aff227 100644
---
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoin.java
+++
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoin.java
@@ -40,7 +40,8 @@ public final class RowTimeIntervalJoin extends
TimeIntervalJoin {
InternalTypeInfo<RowData> rightType,
IntervalJoinFunction joinFunc,
int leftTimeIdx,
- int rightTimeIdx) {
+ int rightTimeIdx,
+ long earlyFireDelay) {
super(
joinType,
leftLowerBound,
@@ -49,7 +50,8 @@ public final class RowTimeIntervalJoin extends
TimeIntervalJoin {
minCleanUpInterval,
leftType,
rightType,
- joinFunc);
+ joinFunc,
+ earlyFireDelay);
this.leftTimeIdx = leftTimeIdx;
this.rightTimeIdx = rightTimeIdx;
}
diff --git
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/TimeIntervalJoin.java
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/TimeIntervalJoin.java
index 4dbf1250dac..f2abfeb00f2 100644
---
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/TimeIntervalJoin.java
+++
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/interval/TimeIntervalJoin.java
@@ -34,6 +34,7 @@ import org.apache.flink.table.data.RowData;
import org.apache.flink.table.runtime.operators.join.FlinkJoinType;
import org.apache.flink.table.runtime.operators.join.OuterJoinPaddingUtil;
import org.apache.flink.table.runtime.typeutils.InternalTypeInfo;
+import org.apache.flink.types.RowKind;
import org.apache.flink.util.Collector;
import org.slf4j.Logger;
@@ -64,6 +65,13 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
private final IntervalJoinFunction joinFunction;
private transient OuterJoinPaddingUtil paddingUtil;
+ // Delay after a row's time at which an unmatched outer row is
speculatively padded and emitted.
+ // A negative value disables early firing, in which case the operator
behaves as a plain
+ // interval join.
+ private final long earlyFireDelay;
+ // True only for an outer join with a non-negative window span and a
non-negative delay.
+ private final boolean earlyFireEnabled;
+
private transient EmitAwareCollector joinCollector;
// cache to store rows form the left stream
@@ -71,6 +79,15 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
// cache to store rows from the right stream
private transient MapState<Long, List<Tuple2<RowData, Boolean>>>
rightCache;
+ // For each cached outer row, whether its speculative early-fire pad has
already been emitted.
+ // The list is positionally aligned 1:1 with the row-time bucket in
leftCache / rightCache, so
+ // firedState.get(t).get(i) corresponds to cache.get(t).get(i). It is kept
as a parallel list
+ // rather than a third tuple field so the existing cache serializer stays
unchanged. The bit
+ // gates both the unmatched window-close pad (it must not be emitted
twice) and the retraction
+ // on a later match (only a row that was speculatively padded needs
correcting).
+ private transient MapState<Long, List<Boolean>> leftFiredState;
+ private transient MapState<Long, List<Boolean>> rightFiredState;
+
// state to record the timer on the left stream. 0 means no timer set
private transient ValueState<Long> leftTimerState;
// state to record the timer on the right stream. 0 means no timer set
@@ -92,7 +109,8 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
long minCleanUpInterval,
InternalTypeInfo<RowData> leftType,
InternalTypeInfo<RowData> rightType,
- IntervalJoinFunction joinFunc) {
+ IntervalJoinFunction joinFunc,
+ long earlyFireDelay) {
this.joinType = joinType;
this.leftRelativeSize = -leftLowerBound;
this.rightRelativeSize = leftUpperBound;
@@ -104,6 +122,14 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
this.leftType = leftType;
this.rightType = rightType;
this.joinFunction = joinFunc;
+ this.earlyFireDelay = earlyFireDelay;
+ // leftRelativeSize + rightRelativeSize equals the window span
(leftUpperBound minus
+ // leftLowerBound), matching the planner gate that enables
update-producing early fire for
+ // outer joins.
+ this.earlyFireEnabled =
+ earlyFireDelay >= 0
+ && joinType.isOuter()
+ && (leftRelativeSize + rightRelativeSize) >= 0;
}
@Override
@@ -129,6 +155,27 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
rightRowListTypeInfo);
rightCache = getRuntimeContext().getMapState(rightMapStateDescriptor);
+ // Early-fire bookkeeping, aligned with the caches above. New
descriptor names restore as
+ // empty state from savepoints taken before early firing existed.
+ if (earlyFireEnabled) {
+ ListTypeInfo<Boolean> firedListTypeInfo =
+ new ListTypeInfo<>(BasicTypeInfo.BOOLEAN_TYPE_INFO);
+ leftFiredState =
+ getRuntimeContext()
+ .getMapState(
+ new MapStateDescriptor<>(
+ "IntervalJoinLeftFired",
+ BasicTypeInfo.LONG_TYPE_INFO,
+ firedListTypeInfo));
+ rightFiredState =
+ getRuntimeContext()
+ .getMapState(
+ new MapStateDescriptor<>(
+ "IntervalJoinRightFired",
+ BasicTypeInfo.LONG_TYPE_INFO,
+ firedListTypeInfo));
+ }
+
// Initialize the timer states.
ValueStateDescriptor<Long> leftValueStateDescriptor =
new ValueStateDescriptor<>("IntervalJoinLeftTimerState",
Long.class);
@@ -178,10 +225,28 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
if (rightTime >= rightQualifiedLowerBound
&& rightTime <= rightQualifiedUpperBound) {
List<Tuple2<RowData, Boolean>> rightRows =
rightEntry.getValue();
+ List<Boolean> rightFired =
+ earlyFireEnabled && joinType.isRightOuter()
+ ? firedBits(rightFiredState, rightTime,
rightRows)
+ : null;
boolean entryUpdated = false;
- for (Tuple2<RowData, Boolean> tuple : rightRows) {
+ for (int i = 0; i < rightRows.size(); i++) {
+ Tuple2<RowData, Boolean> tuple = rightRows.get(i);
joinCollector.reset();
+ boolean retract =
+ rightFired != null
+ && joinType.isRightOuter()
+ && !tuple.f1
+ && rightFired.get(i);
+ if (retract) {
+ // The speculative pad for this right row was
already emitted as an
+ // insert; arm the collector so the match becomes
-U(pad)/+U(match).
+
joinCollector.armRetraction(paddingUtil.padRight(tuple.f0));
+ }
joinFunction.join(leftRow, tuple.f0, joinCollector);
+ if (retract && !joinCollector.isEmitted()) {
+ joinCollector.disarm();
+ }
emitted = emitted || joinCollector.isEmitted();
if (joinType.isRightOuter()) {
if (!tuple.f1 && joinCollector.isEmitted()) {
@@ -200,17 +265,24 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
if (rightTime <= rightExpirationTime) {
if (joinType.isRightOuter()) {
List<Tuple2<RowData, Boolean>> rightRows =
rightEntry.getValue();
- rightRows.forEach(
- (Tuple2<RowData, Boolean> tuple) -> {
- if (!tuple.f1) {
- // Emit a null padding result if the
right row has never
- // been successfully joined.
-
joinCollector.collect(paddingUtil.padRight(tuple.f0));
- }
- });
+ List<Boolean> rightFired =
+ earlyFireEnabled
+ ? firedBits(rightFiredState,
rightTime, rightRows)
+ : null;
+ for (int i = 0; i < rightRows.size(); i++) {
+ Tuple2<RowData, Boolean> tuple = rightRows.get(i);
+ // Skip a row whose speculative pad already fired:
it is correct as
+ // emitted and must not be padded a second time.
+ if (!tuple.f1 && (rightFired == null ||
!rightFired.get(i))) {
+ collectPad(paddingUtil.padRight(tuple.f0));
+ }
+ }
}
// eager remove
rightIterator.remove();
+ if (earlyFireEnabled) {
+ removeFired(rightFiredState, rightTime);
+ }
} // We could do the short-cutting optimization here once we
get a state with
// ordered keys.
}
@@ -226,13 +298,21 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
}
leftRowList.add(Tuple2.of(leftRow, emitted));
leftCache.put(timeForLeftRow, leftRowList);
+ if (earlyFireEnabled && joinType.isLeftOuter()) {
+ // The new tuple has not been speculatively padded yet, so its
bit starts false.
+ appendFired(leftFiredState, timeForLeftRow);
+ if (!emitted) {
+ // Schedule a speculative pad of this unmatched left row
after the delay.
+ registerTimer(ctx, timeForLeftRow + earlyFireDelay);
+ }
+ }
if (rightTimerState.value() == null) {
// Register a timer on the RIGHT stream to remove rows.
registerCleanUpTimer(ctx, timeForLeftRow, true);
}
} else if (!emitted && joinType.isLeftOuter()) {
// Emit a null padding result if the left row is not cached and
successfully joined.
- joinCollector.collect(paddingUtil.padLeft(leftRow));
+ collectPad(paddingUtil.padLeft(leftRow));
}
}
@@ -261,10 +341,28 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
Long leftTime = leftEntry.getKey();
if (leftTime >= leftQualifiedLowerBound && leftTime <=
leftQualifiedUpperBound) {
List<Tuple2<RowData, Boolean>> leftRows =
leftEntry.getValue();
+ List<Boolean> leftFired =
+ earlyFireEnabled && joinType.isLeftOuter()
+ ? firedBits(leftFiredState, leftTime,
leftRows)
+ : null;
boolean entryUpdated = false;
- for (Tuple2<RowData, Boolean> tuple : leftRows) {
+ for (int i = 0; i < leftRows.size(); i++) {
+ Tuple2<RowData, Boolean> tuple = leftRows.get(i);
joinCollector.reset();
+ boolean retract =
+ leftFired != null
+ && joinType.isLeftOuter()
+ && !tuple.f1
+ && leftFired.get(i);
+ if (retract) {
+ // The speculative pad for this left row was
already emitted as an
+ // insert; arm the collector so the match becomes
-U(pad)/+U(match).
+
joinCollector.armRetraction(paddingUtil.padLeft(tuple.f0));
+ }
joinFunction.join(tuple.f0, rightRow, joinCollector);
+ if (retract && !joinCollector.isEmitted()) {
+ joinCollector.disarm();
+ }
emitted = emitted || joinCollector.isEmitted();
if (joinType.isLeftOuter()) {
if (!tuple.f1 && joinCollector.isEmitted()) {
@@ -283,17 +381,24 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
if (leftTime <= leftExpirationTime) {
if (joinType.isLeftOuter()) {
List<Tuple2<RowData, Boolean>> leftRows =
leftEntry.getValue();
- leftRows.forEach(
- (Tuple2<RowData, Boolean> tuple) -> {
- if (!tuple.f1) {
- // Emit a null padding result if the
left row has never been
- // successfully joined.
-
joinCollector.collect(paddingUtil.padLeft(tuple.f0));
- }
- });
+ List<Boolean> leftFired =
+ earlyFireEnabled
+ ? firedBits(leftFiredState, leftTime,
leftRows)
+ : null;
+ for (int i = 0; i < leftRows.size(); i++) {
+ Tuple2<RowData, Boolean> tuple = leftRows.get(i);
+ // Skip a row whose speculative pad already fired:
it is correct as
+ // emitted and must not be padded a second time.
+ if (!tuple.f1 && (leftFired == null ||
!leftFired.get(i))) {
+ collectPad(paddingUtil.padLeft(tuple.f0));
+ }
+ }
}
// eager remove
leftIterator.remove();
+ if (earlyFireEnabled) {
+ removeFired(leftFiredState, leftTime);
+ }
} // We could do the short-cutting optimization here once we
get a state with
// ordered keys.
}
@@ -309,13 +414,21 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
}
rightRowList.add(Tuple2.of(rightRow, emitted));
rightCache.put(timeForRightRow, rightRowList);
+ if (earlyFireEnabled && joinType.isRightOuter()) {
+ // The new tuple has not been speculatively padded yet, so its
bit starts false.
+ appendFired(rightFiredState, timeForRightRow);
+ if (!emitted) {
+ // Schedule a speculative pad of this unmatched right row
after the delay.
+ registerTimer(ctx, timeForRightRow + earlyFireDelay);
+ }
+ }
if (leftTimerState.value() == null) {
// Register a timer on the LEFT stream to remove rows.
registerCleanUpTimer(ctx, timeForRightRow, false);
}
} else if (!emitted && joinType.isRightOuter()) {
// Emit a null padding result if the right row is not cached and
successfully joined.
- joinCollector.collect(paddingUtil.padRight(rightRow));
+ collectPad(paddingUtil.padRight(rightRow));
}
}
@@ -325,6 +438,22 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
joinFunction.setJoinKey(ctx.getCurrentKey());
joinCollector.setInnerCollector(out);
updateOperatorTime(ctx);
+
+ // Early fire runs before cleanup at a shared timestamp so a row that
is both due to fire
+ // and
+ // due to expire emits its speculative pad here; the cleanup branch's
fired-bit gate then
+ // suppresses a second pad. A cleanup-only timestamp finds no live
unfired-unmatched row at
+ // timestamp - earlyFireDelay and is a cheap no-op.
+ if (earlyFireEnabled) {
+ long rowTime = timestamp - earlyFireDelay;
+ if (joinType.isLeftOuter()) {
+ earlyFire(leftCache, leftFiredState, rowTime, true);
+ }
+ if (joinType.isRightOuter()) {
+ earlyFire(rightCache, rightFiredState, rowTime, false);
+ }
+ }
+
// In the future, we should separate the left and right watermarks.
Otherwise, the
// registered timer of the faster stream will be delayed, even if the
watermarks have
// already been emitted by the source.
@@ -332,14 +461,57 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
if (leftCleanUpTime != null && timestamp == leftCleanUpTime) {
rightExpirationTime = calExpirationTime(leftOperatorTime,
rightRelativeSize);
removeExpiredRows(
- joinCollector, rightExpirationTime, rightCache,
leftTimerState, ctx, false);
+ joinCollector,
+ rightExpirationTime,
+ rightCache,
+ rightFiredState,
+ leftTimerState,
+ ctx,
+ false);
}
Long rightCleanUpTime = rightTimerState.value();
if (rightCleanUpTime != null && timestamp == rightCleanUpTime) {
leftExpirationTime = calExpirationTime(rightOperatorTime,
leftRelativeSize);
removeExpiredRows(
- joinCollector, leftExpirationTime, leftCache,
rightTimerState, ctx, true);
+ joinCollector,
+ leftExpirationTime,
+ leftCache,
+ leftFiredState,
+ rightTimerState,
+ ctx,
+ true);
+ }
+ }
+
+ /**
+ * Emit the speculative null-padding result for every cached outer row at
the given row time
+ * that is still unmatched and has not yet had its pad emitted, flipping
its fired bit so
+ * neither this path nor the later window-close pad emits it again.
+ */
+ private void earlyFire(
+ MapState<Long, List<Tuple2<RowData, Boolean>>> rowCache,
+ MapState<Long, List<Boolean>> firedState,
+ long rowTime,
+ boolean padLeft)
+ throws Exception {
+ List<Tuple2<RowData, Boolean>> rows = rowCache.get(rowTime);
+ if (rows == null) {
+ return;
+ }
+ List<Boolean> fired = firedBits(firedState, rowTime, rows);
+ boolean changed = false;
+ for (int i = 0; i < rows.size(); i++) {
+ Tuple2<RowData, Boolean> tuple = rows.get(i);
+ if (!tuple.f1 && !fired.get(i)) {
+ collectPad(
+ padLeft ? paddingUtil.padLeft(tuple.f0) :
paddingUtil.padRight(tuple.f0));
+ fired.set(i, true);
+ changed = true;
+ }
+ }
+ if (changed) {
+ firedState.put(rowTime, fired);
}
}
@@ -396,6 +568,7 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
Collector<RowData> collector,
long expirationTime,
MapState<Long, List<Tuple2<RowData, Boolean>>> rowCache,
+ MapState<Long, List<Boolean>> firedState,
ValueState<Long> timerState,
OnTimerContext ctx,
boolean removeLeft)
@@ -410,28 +583,29 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
Map.Entry<Long, List<Tuple2<RowData, Boolean>>> entry =
iterator.next();
Long rowTime = entry.getKey();
if (rowTime <= expirationTime) {
- if (removeLeft && joinType.isLeftOuter()) {
+ boolean removeOuter =
+ (removeLeft && joinType.isLeftOuter())
+ || (!removeLeft && joinType.isRightOuter());
+ if (removeOuter) {
List<Tuple2<RowData, Boolean>> rows = entry.getValue();
- rows.forEach(
- (Tuple2<RowData, Boolean> tuple) -> {
- if (!tuple.f1) {
- // Emit a null padding result if the row
has never been
- // successfully joined.
-
collector.collect(paddingUtil.padLeft(tuple.f0));
- }
- });
- } else if (!removeLeft && joinType.isRightOuter()) {
- List<Tuple2<RowData, Boolean>> rows = entry.getValue();
- rows.forEach(
- (Tuple2<RowData, Boolean> tuple) -> {
- if (!tuple.f1) {
- // Emit a null padding result if the row
has never been
- // successfully joined.
-
collector.collect(paddingUtil.padRight(tuple.f0));
- }
- });
+ List<Boolean> fired =
+ earlyFireEnabled ? firedBits(firedState, rowTime,
rows) : null;
+ for (int i = 0; i < rows.size(); i++) {
+ Tuple2<RowData, Boolean> tuple = rows.get(i);
+ // Emit a null padding result only if the row was
never matched and its
+ // speculative pad has not already been emitted.
+ if (!tuple.f1 && (fired == null || !fired.get(i))) {
+ collectPad(
+ removeLeft
+ ? paddingUtil.padLeft(tuple.f0)
+ : paddingUtil.padRight(tuple.f0));
+ }
+ }
}
iterator.remove();
+ if (earlyFireEnabled) {
+ removeFired(firedState, rowTime);
+ }
} else {
// We find the earliest timestamp that is still valid.
if (rowTime < earliestTimestamp || earliestTimestamp < 0) {
@@ -447,6 +621,58 @@ abstract class TimeIntervalJoin extends
KeyedCoProcessFunction<RowData, RowData,
// No rows left in the cache. Clear the states and the timerState
will be 0.
timerState.clear();
rowCache.clear();
+ if (earlyFireEnabled && firedState != null) {
+ firedState.clear();
+ }
+ }
+ }
+
+ /**
+ * Emit a padded outer-join row as an insert, overriding any leaked row
kind on the reused row.
+ */
+ private void collectPad(RowData paddedRow) {
+ paddedRow.setRowKind(RowKind.INSERT);
+ joinCollector.collect(paddedRow);
+ }
+
+ /**
+ * Return the fired-bit list aligned with the given cache bucket. Only
called when early firing
+ * is enabled. When the stored list is absent or its length no longer
matches the bucket (e.g.
+ * after a restore), a fresh all-false list of the right length is rebuilt
so no row is ever
+ * treated as already fired.
+ */
+ private List<Boolean> firedBits(
+ MapState<Long, List<Boolean>> firedState,
+ long rowTime,
+ List<Tuple2<RowData, Boolean>> rows)
+ throws Exception {
+ if (firedState != null) {
+ List<Boolean> fired = firedState.get(rowTime);
+ if (fired != null && fired.size() == rows.size()) {
+ return fired;
+ }
+ }
+ List<Boolean> fired = new ArrayList<>(rows.size());
+ for (int i = 0; i < rows.size(); i++) {
+ fired.add(Boolean.FALSE);
+ }
+ return fired;
+ }
+
+ private void appendFired(MapState<Long, List<Boolean>> firedState, long
rowTime)
+ throws Exception {
+ List<Boolean> fired = firedState.get(rowTime);
+ if (fired == null) {
+ fired = new ArrayList<>(1);
+ }
+ fired.add(Boolean.FALSE);
+ firedState.put(rowTime, fired);
+ }
+
+ private void removeFired(MapState<Long, List<Boolean>> firedState, long
rowTime)
+ throws Exception {
+ if (firedState != null) {
+ firedState.remove(rowTime);
}
}
diff --git
a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoinTest.java
b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoinTest.java
index 071a5669afe..42a2e1dacbd 100644
---
a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoinTest.java
+++
b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/ProcTimeIntervalJoinTest.java
@@ -33,6 +33,8 @@ import java.util.ArrayList;
import java.util.List;
import static
org.apache.flink.table.runtime.util.StreamRecordUtils.insertRecord;
+import static
org.apache.flink.table.runtime.util.StreamRecordUtils.updateAfterRecord;
+import static
org.apache.flink.table.runtime.util.StreamRecordUtils.updateBeforeRecord;
import static org.assertj.core.api.Assertions.assertThat;
/** Test for {@link ProcTimeIntervalJoin}. */
@@ -49,7 +51,7 @@ class ProcTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testProcTimeInnerJoinWithCommonBounds() throws Exception {
ProcTimeIntervalJoin joinProcessFunc =
new ProcTimeIntervalJoin(
- FlinkJoinType.INNER, -10, 20, 15, rowType, rowType,
joinFunction);
+ FlinkJoinType.INNER, -10, 20, 15, rowType, rowType,
joinFunction, -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
testHarness.open();
@@ -108,7 +110,7 @@ class ProcTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testProcTimeInnerJoinWithNegativeBounds() throws Exception {
ProcTimeIntervalJoin joinProcessFunc =
new ProcTimeIntervalJoin(
- FlinkJoinType.INNER, -10, -5, 2, rowType, rowType,
joinFunction);
+ FlinkJoinType.INNER, -10, -5, 2, rowType, rowType,
joinFunction, -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -168,6 +170,84 @@ class ProcTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
testHarness.close();
}
+ /** Early fire on processing time, then a match retracts the speculative
pad. */
+ @Test
+ void testProcTimeLeftOuterEarlyFireThenMatch() throws Exception {
+ ProcTimeIntervalJoin joinProcessFunc =
+ new ProcTimeIntervalJoin(
+ FlinkJoinType.LEFT, -5, 9, 0, rowType, rowType,
joinFunction, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.setProcessingTime(10);
+ testHarness.processElement1(insertRecord(1L, "a"));
+ // One cleanup timer plus one early-fire timer at 10 + 3 = 13.
+ assertThat(testHarness.numProcessingTimeTimers()).isEqualTo(2);
+
+ // Fire the early-fire timer: the unmatched left row is speculatively
padded.
+ testHarness.setProcessingTime(13);
+
+ // A right row matches the early-fired left row.
+ testHarness.setProcessingTime(14);
+ testHarness.processElement2(insertRecord(1L, "b"));
+
+ // Advance past cleanup: no further pad.
+ testHarness.setProcessingTime(40);
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(1L, "a", null, null));
+ expectedOutput.add(updateBeforeRecord(1L, "a", null, null));
+ expectedOutput.add(updateAfterRecord(1L, "a", 1L, "b"));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** With early fire disabled the processing-time inner join behaves
exactly as before. */
+ @Test
+ void testProcTimeInnerJoinIgnoresEarlyFire() throws Exception {
+ ProcTimeIntervalJoin joinProcessFunc =
+ new ProcTimeIntervalJoin(
+ FlinkJoinType.INNER, -5, 9, 0, rowType, rowType,
joinFunction, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.setProcessingTime(10);
+ testHarness.processElement1(insertRecord(1L, "a"));
+ // No early-fire timer for an inner join.
+ assertThat(testHarness.numProcessingTimeTimers()).isEqualTo(1);
+
+ testHarness.setProcessingTime(13);
+ testHarness.setProcessingTime(40);
+
+ List<Object> expectedOutput = new ArrayList<>();
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** Delay larger than the window span still pads an unmatched row exactly
once. */
+ @Test
+ void testProcTimeLeftOuterEarlyFireDelayExceedsSpan() throws Exception {
+ // Window span is 5 + 9 = 14; the delay exceeds it so cleanup may
reach the row first.
+ ProcTimeIntervalJoin joinProcessFunc =
+ new ProcTimeIntervalJoin(
+ FlinkJoinType.LEFT, -5, 9, 0, rowType, rowType,
joinFunction, 20L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.setProcessingTime(10);
+ testHarness.processElement1(insertRecord(1L, "a"));
+ // Cleanup at 16, early fire at 30: advancing past both must still
emit a single pad.
+ testHarness.setProcessingTime(35);
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(1L, "a", null, null));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
private KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData>
createTestHarness(ProcTimeIntervalJoin intervalJoinFunc) throws
Exception {
KeyedCoProcessOperator<RowData, RowData, RowData, RowData> operator =
diff --git
a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoinTest.java
b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoinTest.java
index d0e6530c190..a3748a61563 100644
---
a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoinTest.java
+++
b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/join/interval/RowTimeIntervalJoinTest.java
@@ -44,6 +44,8 @@ import static
org.apache.flink.configuration.CheckpointingOptions.CHECKPOINTING_
import static
org.apache.flink.configuration.CheckpointingOptions.ENABLE_UNALIGNED;
import static
org.apache.flink.configuration.CheckpointingOptions.ENABLE_UNALIGNED_INTERRUPTIBLE_TIMERS;
import static
org.apache.flink.table.runtime.util.StreamRecordUtils.insertRecord;
+import static
org.apache.flink.table.runtime.util.StreamRecordUtils.updateAfterRecord;
+import static
org.apache.flink.table.runtime.util.StreamRecordUtils.updateBeforeRecord;
import static org.assertj.core.api.Assertions.assertThat;
/** Test for {@link RowTimeIntervalJoin}. */
@@ -60,7 +62,17 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testRowTimeInnerJoinWithCommonBounds() throws Exception {
RowTimeIntervalJoin joinProcessFunc =
new RowTimeIntervalJoin(
- FlinkJoinType.INNER, -10, 20, 0, 15, rowType, rowType,
joinFunction, 0, 0);
+ FlinkJoinType.INNER,
+ -10,
+ 20,
+ 0,
+ 15,
+ rowType,
+ rowType,
+ joinFunction,
+ 0,
+ 0,
+ -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -125,7 +137,17 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testRowTimeInnerJoinWithNegativeBounds() throws Exception {
RowTimeIntervalJoin joinProcessFunc =
new RowTimeIntervalJoin(
- FlinkJoinType.INNER, -10, -7, 0, 0, rowType, rowType,
joinFunction, 0, 0);
+ FlinkJoinType.INNER,
+ -10,
+ -7,
+ 0,
+ 0,
+ rowType,
+ rowType,
+ joinFunction,
+ 0,
+ 0,
+ -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -180,7 +202,7 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testRowTimeInnerJoinRealtimeCleanUp() throws Exception {
RowTimeIntervalJoin joinProcessFunc =
new RowTimeIntervalJoin(
- FlinkJoinType.LEFT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0);
+ FlinkJoinType.LEFT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -209,7 +231,7 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testRowTimeLeftOuterJoin() throws Exception {
RowTimeIntervalJoin joinProcessFunc =
new RowTimeIntervalJoin(
- FlinkJoinType.LEFT, -5, 9, 0, 7, rowType, rowType,
joinFunction, 0, 0);
+ FlinkJoinType.LEFT, -5, 9, 0, 7, rowType, rowType,
joinFunction, 0, 0, -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -279,7 +301,17 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testRowTimeRightOuterJoin() throws Exception {
RowTimeIntervalJoin joinProcessFunc =
new RowTimeIntervalJoin(
- FlinkJoinType.RIGHT, -5, 9, 0, 7, rowType, rowType,
joinFunction, 0, 0);
+ FlinkJoinType.RIGHT,
+ -5,
+ 9,
+ 0,
+ 7,
+ rowType,
+ rowType,
+ joinFunction,
+ 0,
+ 0,
+ -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -350,7 +382,7 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
void testRowTimeFullOuterJoin() throws Exception {
RowTimeIntervalJoin joinProcessFunc =
new RowTimeIntervalJoin(
- FlinkJoinType.FULL, -5, 9, 0, 7, rowType, rowType,
joinFunction, 0, 0);
+ FlinkJoinType.FULL, -5, 9, 0, 7, rowType, rowType,
joinFunction, 0, 0, -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -439,7 +471,8 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
rowType,
joinFunction,
0,
- 0);
+ 0,
+ -1L);
KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
createTestHarness(joinProcessFunc);
@@ -512,6 +545,257 @@ class RowTimeIntervalJoinTest extends
TimeIntervalStreamJoinTestBase {
testHarness.close();
}
+ /** Early fire: an unmatched left outer row is speculatively padded once
the delay elapses. */
+ @Test
+ void testRowTimeLeftOuterEarlyFire() throws Exception {
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.LEFT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.processElement1(insertRecord(10L, "k1"));
+ // One cleanup timer plus one early-fire timer at 10 + 3 = 13.
+ assertThat(testHarness.numEventTimeTimers()).isEqualTo(2);
+
+ // Cross the early-fire time but not the cleanup time (16): the
speculative pad is emitted.
+ testHarness.processWatermark1(new Watermark(13));
+ testHarness.processWatermark2(new Watermark(13));
+
+ // Cross the cleanup time: the already-fired row must not be padded
again.
+ testHarness.processWatermark1(new Watermark(20));
+ testHarness.processWatermark2(new Watermark(20));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(10L, "k1", null, null));
+ expectedOutput.add(new Watermark(13 - 9));
+ expectedOutput.add(new Watermark(20 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** Early fire then a match: the speculative pad is retracted and replaced
by the joined row. */
+ @Test
+ void testRowTimeLeftOuterEarlyFireThenMatch() throws Exception {
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.LEFT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.processElement1(insertRecord(10L, "k1"));
+ testHarness.processWatermark1(new Watermark(13));
+ testHarness.processWatermark2(new Watermark(13));
+
+ // A right row arrives in window (10 in [12 - 5, 12 + 9]) and matches
the early-fired left
+ // row.
+ testHarness.processElement2(insertRecord(12L, "k1"));
+
+ // Cross cleanup: no further pad, the row already matched.
+ testHarness.processWatermark1(new Watermark(30));
+ testHarness.processWatermark2(new Watermark(30));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(10L, "k1", null, null));
+ expectedOutput.add(new Watermark(13 - 9));
+ expectedOutput.add(updateBeforeRecord(10L, "k1", null, null));
+ expectedOutput.add(updateAfterRecord(10L, "k1", 12L, "k1"));
+ expectedOutput.add(new Watermark(30 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** Symmetric retraction for a right outer join. */
+ @Test
+ void testRowTimeRightOuterEarlyFireThenMatch() throws Exception {
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.RIGHT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.processElement2(insertRecord(10L, "k1"));
+ testHarness.processWatermark1(new Watermark(13));
+ testHarness.processWatermark2(new Watermark(13));
+
+ // A left row in window matches the early-fired right row.
+ testHarness.processElement1(insertRecord(12L, "k1"));
+
+ testHarness.processWatermark1(new Watermark(30));
+ testHarness.processWatermark2(new Watermark(30));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(null, null, 10L, "k1"));
+ expectedOutput.add(new Watermark(13 - 9));
+ expectedOutput.add(updateBeforeRecord(null, null, 10L, "k1"));
+ expectedOutput.add(updateAfterRecord(12L, "k1", 10L, "k1"));
+ expectedOutput.add(new Watermark(30 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** Full outer: both sides early-fire and both are retracted once their
match arrives. */
+ @Test
+ void testRowTimeFullOuterEarlyFireOneMatches() throws Exception {
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.FULL, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ // Left row at 10 and right row at 40, each matched later by a row
from the other side.
+ testHarness.processElement1(insertRecord(10L, "k1"));
+ testHarness.processElement2(insertRecord(40L, "k1"));
+ testHarness.processWatermark1(new Watermark(13));
+ testHarness.processWatermark2(new Watermark(13));
+
+ // Match the left row.
+ testHarness.processElement2(insertRecord(12L, "k1"));
+
+ // Fire the right row's early-fire timer (43).
+ testHarness.processWatermark1(new Watermark(43));
+ testHarness.processWatermark2(new Watermark(43));
+
+ // A left row in window (40 in [45 - 9, 45 + 5]) matches the
early-fired right row.
+ testHarness.processElement1(insertRecord(45L, "k1"));
+
+ // Close everything.
+ testHarness.processWatermark1(new Watermark(60));
+ testHarness.processWatermark2(new Watermark(60));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(10L, "k1", null, null));
+ expectedOutput.add(new Watermark(13 - 9));
+ expectedOutput.add(updateBeforeRecord(10L, "k1", null, null));
+ expectedOutput.add(updateAfterRecord(10L, "k1", 12L, "k1"));
+ expectedOutput.add(insertRecord(null, null, 40L, "k1"));
+ expectedOutput.add(new Watermark(43 - 9));
+ expectedOutput.add(updateBeforeRecord(null, null, 40L, "k1"));
+ expectedOutput.add(updateAfterRecord(45L, "k1", 40L, "k1"));
+ expectedOutput.add(new Watermark(60 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** With early fire disabled the operator output is identical to a plain
interval join. */
+ @Test
+ void testRowTimeInnerJoinIgnoresEarlyFire() throws Exception {
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.INNER, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.processElement1(insertRecord(10L, "k1"));
+ // No early-fire timer for an inner join: only the cleanup timer is
registered.
+ assertThat(testHarness.numEventTimeTimers()).isEqualTo(1);
+
+ testHarness.processWatermark1(new Watermark(13));
+ testHarness.processWatermark2(new Watermark(13));
+ testHarness.processWatermark1(new Watermark(30));
+ testHarness.processWatermark2(new Watermark(30));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(new Watermark(13 - 9));
+ expectedOutput.add(new Watermark(30 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** Delay larger than the window span still pads an unmatched row exactly
once. */
+ @Test
+ void testRowTimeLeftOuterEarlyFireDelayExceedsSpan() throws Exception {
+ // Window span is 5 + 9 = 14; the delay exceeds it so cleanup may
reach the row first.
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.LEFT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 20L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.processElement1(insertRecord(10L, "k1"));
+ // Cleanup at 16, early fire at 30: advancing past both must still
emit a single pad.
+ testHarness.processWatermark1(new Watermark(35));
+ testHarness.processWatermark2(new Watermark(35));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(10L, "k1", null, null));
+ expectedOutput.add(new Watermark(35 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** A normal pad emitted after a retraction must be an insert, not a
leaked update-before. */
+ @Test
+ void testRowTimeEarlyFireRowKindIsolation() throws Exception {
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.LEFT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ // Row A early-fires then matches, producing a retraction that leaves
the reused pad row at
+ // UPDATE_BEFORE.
+ testHarness.processElement1(insertRecord(10L, "k1"));
+ testHarness.processWatermark1(new Watermark(13));
+ testHarness.processWatermark2(new Watermark(13));
+ testHarness.processElement2(insertRecord(12L, "k1"));
+
+ // Row B early-fires and never matches; its window-close pad must be
an insert.
+ testHarness.processElement1(insertRecord(40L, "k2"));
+ testHarness.processWatermark1(new Watermark(60));
+ testHarness.processWatermark2(new Watermark(60));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(10L, "k1", null, null));
+ expectedOutput.add(new Watermark(13 - 9));
+ expectedOutput.add(updateBeforeRecord(10L, "k1", null, null));
+ expectedOutput.add(updateAfterRecord(10L, "k1", 12L, "k1"));
+ expectedOutput.add(insertRecord(40L, "k2", null, null));
+ expectedOutput.add(new Watermark(60 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
+ /** Multiple matches of an early-fired row produce exactly one retraction.
*/
+ @Test
+ void testRowTimeLeftOuterEarlyFireMultiMatch() throws Exception {
+ RowTimeIntervalJoin joinProcessFunc =
+ new RowTimeIntervalJoin(
+ FlinkJoinType.LEFT, -5, 9, 0, 0, rowType, rowType,
joinFunction, 0, 0, 3L);
+ KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData> testHarness =
+ createTestHarness(joinProcessFunc);
+ testHarness.open();
+
+ testHarness.processElement1(insertRecord(10L, "k1"));
+ testHarness.processWatermark1(new Watermark(13));
+ testHarness.processWatermark2(new Watermark(13));
+
+ // First match: corrected via -U/+U.
+ testHarness.processElement2(insertRecord(12L, "k1"));
+ // Second match of the same left row: an ordinary insert, no second
retraction.
+ testHarness.processElement2(insertRecord(14L, "k1"));
+
+ testHarness.processWatermark1(new Watermark(30));
+ testHarness.processWatermark2(new Watermark(30));
+
+ List<Object> expectedOutput = new ArrayList<>();
+ expectedOutput.add(insertRecord(10L, "k1", null, null));
+ expectedOutput.add(new Watermark(13 - 9));
+ expectedOutput.add(updateBeforeRecord(10L, "k1", null, null));
+ expectedOutput.add(updateAfterRecord(10L, "k1", 12L, "k1"));
+ expectedOutput.add(insertRecord(10L, "k1", 14L, "k1"));
+ expectedOutput.add(new Watermark(30 - 9));
+ assertor.assertOutputEquals("output wrong.", expectedOutput,
testHarness.getOutput());
+ testHarness.close();
+ }
+
private KeyedTwoInputStreamOperatorTestHarness<RowData, RowData, RowData,
RowData>
createTestHarness(RowTimeIntervalJoin intervalJoinFunc) throws
Exception {
KeyedCoProcessOperator<RowData, RowData, RowData, RowData> operator =