hudi-agent commented on code in PR #20124:
URL: https://github.com/apache/hudi/pull/20124#discussion_r4127111823
##########
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java:
##########
@@ -560,7 +560,7 @@ protected List<WriteStatus> writeRecords(
rowItr, rowData -> recordConverter.convert(rowData,
rowDataBucket.getBucketInfo()));
List<WriteStatus> statuses = writeFunction.write(
- deduplicateRecordsIfNeeded(recordItr), rowDataBucket.getBucketInfo(),
instant);
+ deduplicateRecordsIfNeeded(recordItr, recordKeyComputer != null),
rowDataBucket.getBucketInfo(), instant);
Review Comment:
🤖 Not introduced here, but relevant to the order this path preserves:
`RowDataBucket.sort()` uses Flink's `QuickSort`, which keeps arrival order only
below 13 records, and neither LSM comparator breaks ties by arrival order. In
larger buckets, equal-key records can be reordered before dedup, so
COMMIT_TIME_ORDERING or event-time ties could keep an older record. Would a
follow-up with an arrival-order tie-breaker, plus a >13-record interleaved
test, make sense? @danny0405 FYI.
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
--
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]