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

   ### Task Description
   
   **What needs to be done:**
   
   Use streaming pre-combine deduplication for Flink LSM write batches that 
have already been sorted by record key. Merge adjacent records with the same 
key using the existing left-to-right reduction, retaining only the current 
merged record and one lookahead record. Keep the existing grouped 
implementation for unsorted input.
   
   **Why this task is needed:**
   
   `StreamWriteFunction.writeRecords()` sorts LSM buckets before pre-combine, 
but `FlinkWriteHelper.deduplicateRecords()` still collects the entire batch 
into a `LinkedHashMap<RecordKey, List<Record>>`. This retains O(N) record 
references and grouping containers even though equal keys are already 
contiguous.
   
   The sorted path can avoid this extra batch retention while preserving the 
existing merger, per-key encounter order, deletion handling, and partial-update 
behavior. The sorting buffer remains necessary; the intended improvement is 
limited to the additional deduplication state. Throughput and peak-memory gains 
require benchmarking.
   
   Validation should cover sorted and unsorted input, iterator laziness, 
event-time and commit-time ordering, equal ordering values, deletes/reinserts, 
partial updates, custom merger reduction order, and COW/MOR writes.
   
   ### Task Type
   
   Performance optimization
   
   ### Related Issues
   
   None.
   


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