weiqingy opened a new pull request, #28953: URL: https://github.com/apache/flink/pull/28953
Part of the FLIP-497 implementation stack under umbrella [FLINK-36953](https://issues.apache.org/jira/browse/FLINK-36953). Landing order: | Step | Sub-task | Scope | | --- | --- | --- | | PR-1a | [FLINK-40167](https://issues.apache.org/jira/browse/FLINK-40167) | EARLY_FIRE hint surface + option validation (#28353, merged) | | PR-1b | [FLINK-40168](https://issues.apache.org/jira/browse/FLINK-40168) | Thread the hint into the interval join (#28796, merged) | | PR-2 | [FLINK-40169](https://issues.apache.org/jira/browse/FLINK-40169) | `target` option (#28827, merged) | | PR-3 | [FLINK-40170](https://issues.apache.org/jira/browse/FLINK-40170) | Update-producing changelog mode + insert-only guard (#28877, merged) | | PR-4 | [FLINK-40171](https://issues.apache.org/jira/browse/FLINK-40171) | Runtime early-fire emit + retraction (#28952, in review) | | **PR-5 (this PR)** | [FLINK-40172](https://issues.apache.org/jira/browse/FLINK-40172) | Processing-time early fire on an event-time join | | PR-6 | [FLINK-40173](https://issues.apache.org/jira/browse/FLINK-40173) | State restore coverage | | PR-7 | [FLINK-40174](https://issues.apache.org/jira/browse/FLINK-40174) | User-facing documentation | 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 - Drop the planning guard that rejected a processing-time delay on an event-time join. - When the hint asks for it, register the early-fire timer on the processing-time timer service, computing the firing time as `currentProcessingTime() + delay`. The join's matching and its state cleanup stay on event time. - Track the row's event time as a bucket key in new per-side schedule state, since a processing-time firing timestamp cannot be mapped back to an event-time bucket. - A single `onTimer` serves both domains and dispatches on `ctx.timeDomain()`. - `StreamExecIntervalJoin` passes the resolved time mode through to the operator. ## Verifying this change This change added tests and can be verified as follows: - `EarlyFireJoinHintTest.testEarlyFireProcTimeOnRowTimeJoin` flips from asserting the old rejection to a golden plan, which records `earlyFireTimeMode=[PROCTIME]` on a join whose bounds are `isRowTime=true`. The opposite pairing stays rejected, covered by `testEarlyFireRowTimeOnProcTimeJoin`. - Harness tests in `RowTimeIntervalJoinTest` drive 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`/`+U` correction, 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: - Dependencies (does it add or upgrade a dependency): no - The public API, i.e., is any changed class annotated with `@Public(Evolving)`: no - The serializers: no - The runtime per-record code paths (performance sensitive): yes, the interval join record path. The behavior is gated on the hint and off by default. - Anything that affects deployment or recovery: yes. This adds per-side early-fire schedule state alongside the bookkeeping `MapState` introduced in PR-4. Restore is covered by tests taking a snapshot both before and after the speculative row fires. - The S3 file system connector: no ## Documentation - Does this pull request introduce a new feature? no (extends the FLIP-497 hint to a time-domain pairing that was previously rejected) - If yes, how is the feature documented? not applicable --- ##### Was generative AI tooling used to co-author this PR? - [X] Yes (please specify the tool below) Generated-by: Claude Code (Anthropic) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
