cshuo opened a new pull request, #20124:
URL: https://github.com/apache/hudi/pull/20124

   ### Describe the issue this Pull Request addresses
   
   Closes #20123.
   
   Flink LSM write buckets are already sorted by record key before pre-combine, 
but deduplication still collects the entire batch into a 
`LinkedHashMap<RecordKey, List<Record>>`. Stream adjacent equal-key records to 
remove this extra batch retention.
   
   ### Summary and Changelog
   
   - Add `FlinkWriteHelper.deduplicateSortedRecords(..., isSortedByRecordKey, 
...)`: reduce sorted input lazily with one lookahead record, and delegate 
unsorted input to the existing implementation.
   - Reuse `reduceRecords()` to preserve the merger, left-to-right reduction 
order, delete handling, and empty-merge-result behavior.
   - Pass the sorted-input flag from `StreamWriteFunction` after bucket 
sorting; callers that have not sorted their input retain the grouped path.
   - Add coverage for lazy iteration, empty input, unsorted non-adjacent 
duplicates, event-time/commit-time ordering, ties, deletes/reinserts, partial 
updates, custom merger order, and COW/MOR writes with default and LSM layouts.
   
   Validation: 13 targeted tests passed (9 helper tests and 4 COW/MOR write 
cases), with Checkstyle and `git diff --check` passing.
   
   ```bash
   mvn -o -Pflink2.2 -pl hudi-flink-datasource/hudi-flink -am \
     
'-Dtest=TestFlinkWriteHelper,TestWriteCopyOnWrite#testDeduplicationWithInterleavedKeys,TestWriteMergeOnRead#testDeduplicationWithInterleavedKeys'
 \
     -Dsurefire.failIfNoSpecifiedTests=false \
     -DskipITs -DskipSparkTests -DskipScalaTests test
   ```
   
   ### Impact
   
   For sorted LSM batches with pre-combine enabled, deduplication retains a 
constant number of record states instead of O(N) batch references and grouping 
containers. The sorting buffer remains allocated, and the number of merge 
operations is unchanged. No new configuration or storage-format changes. 
Throughput and peak-memory improvements have not been benchmarked.
   
   ### Risk Level
   
   Low. The streaming path requires contiguous equal record keys and records 
that remain valid as the iterator advances. It is selected only after the 
existing LSM bucket sort; unsorted callers use the old implementation. 
Differential tests cover merge semantics, and COW/MOR tests cover write-path 
integration.
   
   ### Documentation Update
   
   Added method Javadoc documenting sorted-input preconditions and unsorted 
fallback. No user-facing documentation changes are required.
   
   ### Contributor's checklist
   
   - [ ] Read through [contributor's 
guide](https://hudi.apache.org/contribute/how-to-contribute)
   - [x] Enough context is provided in the sections above
   - [x] Adequate tests were added if applicable
   


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