[FLINK-40172][table-runtime] Support processing-time early fire on a row-time interval join - #28953
Draft
weiqingy wants to merge 2 commits into
Draft
[FLINK-40172][table-runtime] Support processing-time early fire on a row-time interval join#28953weiqingy wants to merge 2 commits into
weiqingy wants to merge 2 commits into
Conversation
…he interval join operator Wire the EARLY_FIRE delay into the interval join operator so an outer join speculatively emits its padded unmatched row after the delay and corrects it when a real match arrives. Covers the natural timer pairings: a row-time join fires on event time, a processing-time join fires on processing time. Processing-time triggering on a row-time join stays rejected at planning. When an unmatched outer row is cached, the operator registers an early-fire timer at rowTime + delay. On that timer it emits the padded row as an INSERT and records that it fired. When the row later matches, it retracts the padded row as UPDATE_BEFORE and emits the matched row as UPDATE_AFTER, matching the update-producing changelog mode inferred for the node. The retraction is tied to the one-time matched-and-emitted flip, so a row that matches several times emits a single correction followed by ordinary inserts. The already-fired marker is a new per-side MapState<Long, List<Boolean>> kept positionally aligned with the existing row cache, rather than widening the cache tuple, so the cache serializer is unchanged and old savepoints restore the new state empty. The marker is the single gate that keeps a row padded exactly once when the delay is at or beyond the window span. All early-fire work is gated on the hint being set, an outer join, and a non-negative window, so a plain interval join is unchanged and allocates nothing new. EmitAwareCollector carries the changelog stamping so IntervalJoinFunction stays changelog-agnostic, and every padded or matched emit stamps its RowKind explicitly to avoid leaking a kind onto a reused row.
…row-time interval join
Add the cross-domain timer combination the previous commit left out: an
event-time interval join with EARLY_FIRE('time_mode'='proctime') now fires its
speculative pads on the wall clock while keeping its event-time cleanup. The
temporary "not yet supported" rejection in the planner rule is removed; the
row-time-on-processing-time rejection is retained.
onTimer distinguishes the two timer kinds by OnTimerContext.timeDomain(): in
the cross-domain case early-fire timers are processing-time and cleanup timers
are event-time, so a processing-time firing runs early fire and returns while
an event-time firing runs cleanup only. The discrimination is gated on a new
cross-domain flag, so the natural pairings keep the previous timestamp - delay
recovery where early fire and cleanup share a domain.
A processing-time firing timestamp cannot be mapped back to an event-time cache
bucket arithmetically, so a per-side MapState<Long, List<Long>> keyed by firing
processing-time records the event-time bucket keys due to fire then. It is
allocated only in the cross-domain case and reuses the existing per-bucket emit
and positional fired bit, so the retract-and-correct path is shared. Every
scheduled firing time fires and removes its own entry, and a bucket already
cleaned by event-time expiry makes the firing a no-op, so nothing accumulates.
The schedule is value-typed and order-preserving and processing-time timers are
checkpointed, so a timer pending at snapshot fires after restore against the
restored schedule and fired bits and emits at most the not-yet-emitted pad.
Harness tests cover the wall-clock trigger without watermark advance, a snapshot
before the timer fires, and a snapshot after the pad is emitted.
Collaborator
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Part of the FLIP-497 implementation stack under umbrella FLINK-36953. Landing order:
targetoption (#28827, merged)Opened as a draft because it is stacked on #28952, which is in review. Until that merges, the commit list and diff here also carry PR-4's commit. Once #28952 merges I will rebase onto master, leaving only this PR's change, and take it out of draft.
What is the purpose of the change
Allows a processing-time early-fire delay on an event-time interval join. That pairing was previously rejected at planning with a "not yet supported" error. The join still matches on event time; only the early-fire trigger runs on processing time, so a speculative row can be emitted on a wall-clock delay even when watermarks are sparse or stalled.
Brief change log
currentProcessingTime() + delay. The join's matching and its state cleanup stay on event time.onTimerserves both domains and dispatches onctx.timeDomain().StreamExecIntervalJoinpasses the resolved time mode through to the operator.Verifying this change
This change added tests and can be verified as follows:
EarlyFireJoinHintTest.testEarlyFireProcTimeOnRowTimeJoinflips from asserting the old rejection to a golden plan, which recordsearlyFireTimeMode=[PROCTIME]on a join whose bounds areisRowTime=true. The opposite pairing stays rejected, covered bytestEarlyFireRowTimeOnProcTimeJoin.RowTimeIntervalJoinTestdrive the two clocks independently. One case fires the pad on a wall-clock advance with the watermark untouched. Together they cover the speculative emit, its-U/+Ucorrection, an inner join correctly ignoring the hint, and snapshot/restore on both sides of the firing.Does this pull request potentially affect one of the following parts:
@Public(Evolving): noMapStateintroduced in PR-4. Restore is covered by tests taking a snapshot both before and after the speculative row fires.Documentation
Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Anthropic)