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]