stivsh commented on PR #29382:
URL: https://github.com/apache/flink/pull/29382#issuecomment-5995734514
Thanks for the review - agreed this should be tracked as a bug, not just a
performance improvement. I've applied for Flink JIRA access; while that's
pending, here is the write-up I intend to file as the JIRA issue, so the
reasoning is visible here in the meantime:
**Problem**
`TemporalRowTimeJoinOperator` (the rowtime variant of the SQL temporal /
versioned-table join, `FOR SYSTEM_TIME AS OF`) buffers probe-side (left) rows
in `leftState` until the watermark passes their event time. On every event-time
timer firing - which happens on *every* watermark advance, including ones
driven purely by the *right* (build) side - `emitResultAndCleanUpState` scans
*all* entries currently in `leftState`, not just the ones that are ready to be
emitted.
This means:
- If a startup spike (or any burst) causes a large backlog of
not-yet-ready rows to accumulate in `leftState` while the watermark is lagging
behind,
- then once the watermark starts advancing - even in small increments,
each one triggering its own timer - the operator rescans the *entire* backlog
on *every single one* of those advances, picking off only the handful of rows
that just became ready each time.
The result is O(backlogSize) work per watermark tick until the backlog fully
drains, with no way to catch up faster than one tick at a time. In practice,
under a sufficiently large initial spike combined with a steadily (even slowly)
advancing watermark, the join can look like it has stalled completely: CPU
pegged scanning the same not-yet-ready rows over and over, while making only
slow forward progress on the actual backlog. This is not a performance nit, it
is a correctness/availability problem - the join can become effectively
unusable for any workload with a sizeable backlog plus a spiky arrival pattern
(backfill, a consumer catching up after downtime, or just a traffic spike),
which is a realistic and common scenario for this operator, not an edge case.
**Proposed fix**
This PR replaces the full scan of `leftState` with a binary min-heap ordered
by event time (`MapStateHeap`), implemented directly on top of the existing
`MapState` (used as a dynamic array with integer keys `0..N-1`).
`emitResultAndCleanUpState` then only reads/removes exactly the rows that are
ready for the current watermark - O(K log N) per tick instead of O(N), where K
is the number of rows actually emitted and N is the backlog size.
As submitted, the patch changes the serialized shape of `leftState`
(`RowData` -> `Tuple2<Long, RowData>`), so it is not compatible with savepoints
taken on an older version. If that is a blocker for including this in a
maintenance release rather than waiting for a major version, a state-migration
path (reading the old format once on restore and rebuilding the heap from it)
can be added - happy to implement that if it unblocks acceptance.
Without a fix along these lines, this operator is, in my experience,
practically unusable for TemporalJoin workloads with this kind of backlog/spike
pattern - not just mine; this is inherent to how the operator buffers and
re-scans state, so I'd expect it to affect any similar temporal-join use case
under load.
I'll open the formal JIRA issue and update this PR's title/branch with the
`FLINK-XXXX` prefix as soon as my JIRA access comes through.
--
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]