hudi-agent commented on code in PR #20138: URL: https://github.com/apache/hudi/pull/20138#discussion_r4148232416
########## hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/buffer/StableSortBuffer.java: ########## @@ -0,0 +1,232 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hudi.sink.buffer; + +import org.apache.flink.core.memory.MemorySegment; +import org.apache.flink.runtime.io.disk.SimpleCollectingOutputView; +import org.apache.flink.table.data.RowData; +import org.apache.flink.table.data.binary.BinaryRowData; +import org.apache.flink.table.runtime.generated.NormalizedKeyComputer; +import org.apache.flink.table.runtime.generated.RecordComparator; +import org.apache.flink.table.runtime.operators.sort.BinaryIndexedSortable; +import org.apache.flink.table.runtime.typeutils.AbstractRowDataSerializer; +import org.apache.flink.table.runtime.typeutils.BinaryRowDataSerializer; +import org.apache.flink.table.runtime.util.MemorySegmentPool; +import org.apache.flink.util.MutableObjectIterator; + +import java.io.EOFException; +import java.io.IOException; +import java.util.ArrayList; + +import static org.apache.flink.util.Preconditions.checkArgument; + +/** + * Copied from Flink's {@code BinaryInMemorySortBuffer}. + * + * <p>Provides stable sorting by comparing record offsets when keys are equal, + * preserving the original insertion order of records with equal keys. + */ +public final class StableSortBuffer extends BinaryIndexedSortable { + + private static final int MIN_REQUIRED_BUFFERS = 3; + + private final AbstractRowDataSerializer<RowData> inputSerializer; + private final ArrayList<MemorySegment> recordBufferSegments; + private final SimpleCollectingOutputView recordCollector; + private final int totalNumBuffers; + + private long currentDataBufferOffset; + private long sortIndexBytes; + + /** Create a memory sorter in `insert` way. */ + public static StableSortBuffer createBuffer( + NormalizedKeyComputer normalizedKeyComputer, + AbstractRowDataSerializer<RowData> inputSerializer, + BinaryRowDataSerializer serializer, + RecordComparator comparator, + MemorySegmentPool memoryPool) { + checkArgument(memoryPool.freePages() >= MIN_REQUIRED_BUFFERS); + int totalNumBuffers = memoryPool.freePages(); + ArrayList<MemorySegment> recordBufferSegments = new ArrayList<>(16); + return new StableSortBuffer( + normalizedKeyComputer, + inputSerializer, + serializer, + comparator, + recordBufferSegments, + new SimpleCollectingOutputView( + recordBufferSegments, memoryPool, memoryPool.pageSize()), + memoryPool, + totalNumBuffers); + } + + private StableSortBuffer( + NormalizedKeyComputer normalizedKeyComputer, + AbstractRowDataSerializer<RowData> inputSerializer, + BinaryRowDataSerializer serializer, + RecordComparator comparator, + ArrayList<MemorySegment> recordBufferSegments, + SimpleCollectingOutputView recordCollector, + MemorySegmentPool pool, + int totalNumBuffers) { + super(normalizedKeyComputer, serializer, comparator, recordBufferSegments, pool); + this.inputSerializer = inputSerializer; + this.recordBufferSegments = recordBufferSegments; + this.recordCollector = recordCollector; + this.totalNumBuffers = totalNumBuffers; + } + + // Preserve arrival order for equal keys using the existing record offsets. + @Override + public int compare(int segmentNumberI, int segmentOffsetI, int segmentNumberJ, int segmentOffsetJ) { + int result = super.compare(segmentNumberI, segmentOffsetI, segmentNumberJ, segmentOffsetJ); + // ------------------ hoodie update --------------------- Review Comment: 🤖 nit: the `// --- hoodie update ---` marker sits in the middle of `compare`, and nothing marks where the Hudi change ends. Could you drop it and rely on the comment above the method, or wrap the whole override in matching begin/end markers? <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag quality.</i></sub> ########## hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/buffer/StableSortBuffer.java: ########## @@ -0,0 +1,232 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hudi.sink.buffer; + +import org.apache.flink.core.memory.MemorySegment; +import org.apache.flink.runtime.io.disk.SimpleCollectingOutputView; +import org.apache.flink.table.data.RowData; +import org.apache.flink.table.data.binary.BinaryRowData; +import org.apache.flink.table.runtime.generated.NormalizedKeyComputer; +import org.apache.flink.table.runtime.generated.RecordComparator; +import org.apache.flink.table.runtime.operators.sort.BinaryIndexedSortable; +import org.apache.flink.table.runtime.typeutils.AbstractRowDataSerializer; +import org.apache.flink.table.runtime.typeutils.BinaryRowDataSerializer; +import org.apache.flink.table.runtime.util.MemorySegmentPool; +import org.apache.flink.util.MutableObjectIterator; + +import java.io.EOFException; +import java.io.IOException; +import java.util.ArrayList; + +import static org.apache.flink.util.Preconditions.checkArgument; + +/** + * Copied from Flink's {@code BinaryInMemorySortBuffer}. Review Comment: 🤖 nit: could you say which Flink version `BinaryInMemorySortBuffer` was copied from, e.g. "Copied from Flink 1.20's ..."? It also seems worth noting that `compare` is the only intended change. <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]
