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]
