cshuo opened a new issue, #20131:
URL: https://github.com/apache/hudi/issues/20131

   ### Bug Description
   
   **What happened:**
   
   Flink LSM writes sort buffered rows by record key before pre-combine. 
`RowDataBucket.sort()` uses Flink's `QuickSort`, and the record-key comparators 
do not break equal-key ties by arrival order. The sort can therefore reorder 
records with the same key.
   
   Under `COMMIT_TIME_ORDERING`, the batch reducer passes each subsequent 
record as `newRecord`, and `CommitTimeRecordMerger.deltaMerge()` selects that 
record. Reordering equal keys can make an earlier-arriving value overwrite a 
later-arriving value within the same write batch.
   
   This predates the streaming deduplication optimization in #20124. The old 
`LinkedHashMap` implementation also consumes the already-sorted input and has 
the same behavior.
   
   **What you expected:**
   
   Sorting an LSM bucket by record key should preserve the encounter order 
within each key so that commit-time pre-combine retains the last-arriving 
record in that bucket. This concerns ordering within a single buffered batch, 
not a global arrival order across parallel writers.
   
   **Steps to reproduce:**
   
   1. Configure Flink streaming writes with LSM storage layout, pre-combine 
enabled, and `COMMIT_TIME_ORDERING`.
   2. Buffer 14 records in one bucket without an intervening flush. Alternate 
keys `A` and `B`, with values identifying arrival sequence `0..13` (even 
sequences belong to `A`).
   3. Flush the bucket and inspect the equal-key order after sorting and the 
pre-combine winner.
   
   A standalone probe using the actual Flink 2.2.1 `QuickSort` with an 
`IndexedSortable` comparing only the alternating keys produced:
   
   ```text
   A arrival sequence before sort: 0, 2, 4, 6, 8, 10, 12
   A arrival sequence after sort:  0, 8, 4, 10, 2, 12, 6
   ```
   
   The existing commit-time reduction would select arrival `6` rather than `12` 
for this sorted sequence. The sort reordering was reproduced directly; a full 
Hudi write-path regression test still needs to be added. Flink 2.2.1 uses 
insertion sort for ranges smaller than 13 records, so small tests can miss the 
issue.
   
   **Suggested follow-up:**
   
   - Preserve equal-key arrival order, for example with an arrival-sequence 
tie-breaker or a stable sorting strategy, while retaining encoded record-key 
ordering.
   - Add a regression test with at least 13 interleaved records, preferably a 
larger batch, and assert that commit-time pre-combine keeps the last-arriving 
value per key.
   - Cover both generated string-key sorting and the encoded-key comparator 
fallback, and assess the memory overhead of the chosen solution.
   
   ### Environment
   
   **Hudi version:** 1.3.0-SNAPSHOT; code inspected at 
`77a5850b6c5a3b2e359648df1590d5b30dd94c51`.
   
   **Query engine:** Flink; standalone sort probe used Flink 2.2.1.
   
   **Relevant configuration:** LSM_TREE storage layout, PRE_COMBINE enabled, 
COMMIT_TIME_ORDERING merge mode.
   
   ### Related discussion
   
   Identified during [review of 
#20124](https://github.com/apache/hudi/pull/20124#discussion_r4127111823). 
Tracked separately from the deduplication memory optimization in #20123 / 
#20124.
   


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