stivsh opened a new pull request, #29382:
URL: https://github.com/apache/flink/pull/29382

   ## What is the purpose of the change
   
   I hit a performance problem in `TemporalRowTimeJoinOperator`.
   
   During startup, when both sides produce a lot of data, millions of rows 
accumulate in `leftState`. On every watermark timer, Flink scans **all of 
them** to find rows that are ready to emit. Rows that are not ready are scanned 
again on the next watermark.
   
   In my case join stucked with out 100s of CPU.
   
   I propose replacing this full scan with a binary min-heap ordered by event 
time.
   
   The heap is implemented as `MapStateHeap` on top of `MapState`. `MapState` 
is used as a dynamic array with integer keys `0..N-1`, which is enough to 
implement a standard binary heap.
   
   With this structure, `emitResultAndCleanUpState` reads and removes only rows 
that are ready for the current watermark. **It no longer scans rows that should 
stay in the state**.
   
   The timer cost changes from `O(N)` per watermark to `O(K log N)`, where `K` 
is the number of rows emitted and `N` is the number of buffered rows.
   
   The idea that any structure could be impemented on top of the map could also 
be used for the right side, but this change only updates the left side.
   
   There is currently no Flink JIRA issue for this change, so I could not name 
the branch or PR with a `FLINK-XXXX` prefix.
   
   ## Brief change log
   
   - Added `MapStateHeap`, a binary min-heap backed by `MapState`.
   - Replaced the full scan of `leftState` with heap-based processing.
   - `emitResultAndCleanUpState` now processes only rows ready for the current 
watermark.
   


-- 
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