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]

Reply via email to