diff --git a/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarRowIterator.java b/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarRowIterator.java index 4489388e3410..0779b5b8b139 100644 --- a/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarRowIterator.java +++ b/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarRowIterator.java @@ -115,6 +115,10 @@ public ColumnarRowIterator copy(ColumnVector[] vectors) { public ColumnarRowIterator mapping( @Nullable PartitionInfo partitionInfo, @Nullable int[] indexMapping) { + if (partitionInfo == null && isIdentityMapping(indexMapping, row.batch().getArity())) { + return this; + } + if (partitionInfo != null || indexMapping != null) { VectorizedColumnBatch vectorizedColumnBatch = row.batch(); ColumnVector[] vectors = vectorizedColumnBatch.columns; @@ -129,6 +133,19 @@ public ColumnarRowIterator mapping( return this; } + private static boolean isIdentityMapping(@Nullable int[] indexMapping, int arity) { + if (indexMapping == null || indexMapping.length != arity) { + return false; + } + + for (int i = 0; i < indexMapping.length; i++) { + if (indexMapping[i] != i) { + return false; + } + } + return true; + } + /** * Strips a row-tracking wrapper previously installed by {@link #assignRowTracking}, so repeated * assignment re-wraps the base vector instead of nesting. diff --git a/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowIteratorTest.java b/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowIteratorTest.java index 8734886d3751..5a3722443ed1 100644 --- a/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowIteratorTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowIteratorTest.java @@ -18,10 +18,15 @@ package org.apache.paimon.data.columnar; +import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.data.BinaryRowWriter; +import org.apache.paimon.data.PartitionInfo; import org.apache.paimon.data.columnar.heap.HeapIntVector; import org.apache.paimon.data.columnar.heap.HeapLongVector; import org.apache.paimon.fs.Path; import org.apache.paimon.table.SpecialFields; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; import org.apache.paimon.utils.LongIterator; import org.junit.jupiter.api.Test; @@ -115,4 +120,70 @@ public void testRepeatedAssignRowTrackingDoesNotNest() { assertThat(tracked.next()).isNotNull(); assertThat(tracked.row.getLong(0)).isEqualTo(300L); } + + @Test + public void testIdentityMappingPreservesSpecializedIterator() { + HeapIntVector firstVector = new HeapIntVector(1); + HeapIntVector secondVector = new HeapIntVector(1); + VectorizedColumnBatch batch = + new VectorizedColumnBatch(new ColumnVector[] {firstVector, secondVector}); + batch.setNumRows(1); + ColumnarRowIterator rowIterator = new TestingSpecializedIterator(batch); + rowIterator.reset(0); + + assertThat(rowIterator.mapping(null, new int[] {0, 1})).isSameAs(rowIterator); + } + + @Test + public void testNonIdentityMappingCopiesIterator() { + HeapIntVector firstVector = new HeapIntVector(1); + HeapIntVector secondVector = new HeapIntVector(1); + VectorizedColumnBatch batch = + new VectorizedColumnBatch(new ColumnVector[] {firstVector, secondVector}); + batch.setNumRows(1); + ColumnarRowIterator rowIterator = new TestingSpecializedIterator(batch); + rowIterator.reset(0); + + ColumnarRowIterator reordered = rowIterator.mapping(null, new int[] {1, 0}); + assertThat(reordered).isNotSameAs(rowIterator); + assertThat(reordered.batch().columns).containsExactly(secondVector, firstVector); + + ColumnarRowIterator duplicated = rowIterator.mapping(null, new int[] {0, 0}); + assertThat(duplicated).isNotSameAs(rowIterator); + assertThat(duplicated.batch().columns).containsExactly(firstVector, firstVector); + + ColumnarRowIterator projected = rowIterator.mapping(null, new int[] {0}); + assertThat(projected).isNotSameAs(rowIterator); + assertThat(projected.batch().columns).containsExactly(firstVector); + } + + @Test + public void testPartitionMappingCopiesIterator() { + HeapIntVector dataVector = new HeapIntVector(1); + dataVector.setInt(0, 7); + VectorizedColumnBatch batch = new VectorizedColumnBatch(new ColumnVector[] {dataVector}); + batch.setNumRows(1); + ColumnarRowIterator rowIterator = new TestingSpecializedIterator(batch); + rowIterator.reset(0); + + BinaryRow partition = new BinaryRow(1); + BinaryRowWriter writer = new BinaryRowWriter(partition); + writer.writeInt(0, 42); + writer.complete(); + PartitionInfo partitionInfo = + new PartitionInfo(new int[] {1, -1, 0}, RowType.of(DataTypes.INT()), partition); + + ColumnarRowIterator mapped = rowIterator.mapping(partitionInfo, new int[] {0, 1}); + assertThat(mapped).isNotSameAs(rowIterator); + assertThat(mapped.batch().getArity()).isEqualTo(2); + assertThat(mapped.batch().getInt(0, 0)).isEqualTo(7); + assertThat(mapped.batch().getInt(0, 1)).isEqualTo(42); + } + + private static class TestingSpecializedIterator extends ColumnarRowIterator { + + private TestingSpecializedIterator(VectorizedColumnBatch batch) { + super(new Path("test"), new ColumnarRow(batch), null); + } + } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileRecordReader.java b/paimon-core/src/main/java/org/apache/paimon/io/DataFileRecordReader.java index 70650fbf4b8c..038ac57662d8 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileRecordReader.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileRecordReader.java @@ -25,6 +25,7 @@ import org.apache.paimon.data.GenericRow; import org.apache.paimon.data.InternalRow; import org.apache.paimon.data.PartitionInfo; +import org.apache.paimon.data.columnar.ColumnVector; import org.apache.paimon.data.columnar.ColumnarRowIterator; import org.apache.paimon.format.FormatReaderFactory; import org.apache.paimon.fs.Path; @@ -178,8 +179,26 @@ private FileRecordIterator readBatchInternal() throws IOException { } if (iterator instanceof ColumnarRowIterator) { - iterator = ((ColumnarRowIterator) iterator).mapping(partitionInfo, indexMapping); - if (rowTrackingEnabled) { + ColumnarRowIterator sourceIterator = (ColumnarRowIterator) iterator; + iterator = sourceIterator.mapping(partitionInfo, indexMapping); + boolean assignRowTracking = + rowTrackingEnabled + && (systemFields.containsKey(SpecialFields.SEQUENCE_NUMBER.name()) + || (firstRowId != null + && systemFields.containsKey( + SpecialFields.ROW_ID.name()))); + if (assignRowTracking) { + if (iterator == sourceIterator) { + // Copy to a ColumnVector[] because cloning a subtype array preserves its + // runtime type and cannot accept row-tracking wrapper vectors. + ColumnarRowIterator columnarIterator = (ColumnarRowIterator) iterator; + iterator = + columnarIterator.copy( + Arrays.copyOf( + columnarIterator.batch().columns, + columnarIterator.batch().columns.length, + ColumnVector[].class)); + } iterator = ((ColumnarRowIterator) iterator) .assignRowTracking(firstRowId, maxSequenceNumber, systemFields); diff --git a/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java b/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java new file mode 100644 index 000000000000..ed5be946d0df --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java @@ -0,0 +1,262 @@ +/* + * 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.paimon.io; + +import org.apache.paimon.data.InternalRow; +import org.apache.paimon.data.columnar.ColumnVector; +import org.apache.paimon.data.columnar.ColumnarRow; +import org.apache.paimon.data.columnar.ColumnarRowIterator; +import org.apache.paimon.data.columnar.VectorizedColumnBatch; +import org.apache.paimon.data.columnar.VectorizedRowIterator; +import org.apache.paimon.data.columnar.heap.HeapLongVector; +import org.apache.paimon.fs.Path; +import org.apache.paimon.reader.FileRecordIterator; +import org.apache.paimon.reader.FileRecordReader; +import org.apache.paimon.table.SpecialFields; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; + +import org.junit.jupiter.api.Test; + +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link DataFileRecordReader}. */ +public class DataFileRecordReaderTest { + + @Test + public void testRowTrackingIsolatedFromReusedIdentityBatch() throws Exception { + ReusingColumnarReader delegate = new ReusingColumnarReader(); + assertThat(delegate.batch.columns).isSameAs(delegate.readerOwnedColumns); + assertThat(delegate.readerOwnedColumns.getClass()).isEqualTo(HeapLongVector[].class); + + Map systemFields = new HashMap<>(); + systemFields.put(SpecialFields.ROW_ID.name(), 1); + systemFields.put(SpecialFields.SEQUENCE_NUMBER.name(), 2); + DataFileRecordReader reader = + createRowTrackingReader( + SpecialFields.rowTypeWithRowTracking(RowType.of(DataTypes.BIGINT())), + delegate, + new int[] {0, 1, 2}, + systemFields); + + for (int batchIndex = 0; batchIndex < 3; batchIndex++) { + FileRecordIterator iterator = reader.readBatch(); + + assertThat(iterator).isNotSameAs(delegate.iterator); + assertThat(iterator).isInstanceOf(VectorizedRowIterator.class); + assertReaderOwnedColumns(delegate); + + InternalRow row = iterator.next(); + assertThat(row.getFieldCount()).isEqualTo(3); + assertThat(row.getLong(0)).isEqualTo(42L); + assertThat(row.getLong(1)).isEqualTo(100L + batchIndex); + assertThat(row.getLong(2)).isEqualTo(7L); + assertThat(delegate.recycleCount).isEqualTo(batchIndex); + + iterator.releaseBatch(); + assertThat(delegate.recycleCount).isEqualTo(batchIndex + 1); + assertReaderOwnedColumns(delegate); + } + reader.close(); + } + + @Test + public void testSingletonRowTrackingFieldAssigned() throws Exception { + ReusingColumnarReader delegate = new ReusingColumnarReader(); + DataFileRecordReader reader = + createRowTrackingReader( + SpecialFields.rowTypeWithRowId(RowType.of(DataTypes.BIGINT())), + delegate, + new int[] {0, 1}, + Collections.singletonMap(SpecialFields.ROW_ID.name(), 1)); + + FileRecordIterator iterator = reader.readBatch(); + + assertThat(iterator).isNotSameAs(delegate.iterator); + assertThat(iterator).isInstanceOf(VectorizedRowIterator.class); + assertReaderOwnedColumns(delegate); + + InternalRow row = iterator.next(); + assertThat(row.getFieldCount()).isEqualTo(2); + assertThat(row.getLong(0)).isEqualTo(42L); + assertThat(row.getLong(1)).isEqualTo(100L); + assertThat(delegate.recycleCount).isZero(); + + iterator.releaseBatch(); + assertThat(delegate.recycleCount).isEqualTo(1); + assertReaderOwnedColumns(delegate); + reader.close(); + } + + @Test + public void testEmptyRowTrackingFieldsPreserveSpecializedIdentityIterator() throws Exception { + HeapLongVector dataVector = new HeapLongVector(1); + dataVector.setLong(0, 42L); + VectorizedColumnBatch batch = new VectorizedColumnBatch(new ColumnVector[] {dataVector}); + batch.setNumRows(1); + ColumnarRowIterator specializedIterator = + new ColumnarRowIterator(new Path("test"), new ColumnarRow(batch), null) {}; + specializedIterator.reset(0); + + FileRecordReader delegate = + new FileRecordReader() { + @Override + public FileRecordIterator readBatch() { + return specializedIterator; + } + + @Override + public void close() {} + }; + DataFileRecordReader reader = + new DataFileRecordReader( + RowType.of(DataTypes.BIGINT()), + delegate, + false, + false, + new int[] {0}, + null, + null, + true, + 100L, + 7L, + Collections.emptyMap(), + null, + new Path("test")); + + assertThat(reader.readBatch()).isSameAs(specializedIterator); + reader.close(); + } + + @Test + public void testStoredRowIdPreservesSpecializedIdentityIterator() throws Exception { + HeapLongVector dataVector = new HeapLongVector(1); + dataVector.setLong(0, 42L); + HeapLongVector rowIdVector = new HeapLongVector(1); + rowIdVector.setLong(0, 123L); + VectorizedColumnBatch batch = + new VectorizedColumnBatch(new ColumnVector[] {dataVector, rowIdVector}); + batch.setNumRows(1); + int[] recycleCount = {0}; + ColumnarRowIterator specializedIterator = + new ColumnarRowIterator( + new Path("test"), new ColumnarRow(batch), () -> recycleCount[0]++) {}; + specializedIterator.reset(0); + + FileRecordReader delegate = + new FileRecordReader() { + @Override + public FileRecordIterator readBatch() { + return specializedIterator; + } + + @Override + public void close() {} + }; + DataFileRecordReader reader = + new DataFileRecordReader( + SpecialFields.rowTypeWithRowId(RowType.of(DataTypes.BIGINT())), + delegate, + false, + false, + new int[] {0, 1}, + null, + null, + true, + null, + 7L, + Collections.singletonMap(SpecialFields.ROW_ID.name(), 1), + null, + new Path("test")); + + FileRecordIterator iterator = reader.readBatch(); + assertThat(iterator).isSameAs(specializedIterator); + InternalRow row = iterator.next(); + assertThat(row.getLong(0)).isEqualTo(42L); + assertThat(row.getLong(1)).isEqualTo(123L); + + iterator.releaseBatch(); + assertThat(recycleCount[0]).isEqualTo(1); + reader.close(); + } + + private static DataFileRecordReader createRowTrackingReader( + RowType rowType, + ReusingColumnarReader delegate, + int[] indexMapping, + Map systemFields) { + return new DataFileRecordReader( + rowType, + delegate, + false, + false, + indexMapping, + null, + null, + true, + 100L, + 7L, + systemFields, + null, + new Path("test")); + } + + private static void assertReaderOwnedColumns(ReusingColumnarReader delegate) { + assertThat(delegate.batch.columns).isSameAs(delegate.readerOwnedColumns); + assertThat(delegate.readerOwnedColumns) + .containsExactly( + delegate.dataVector, delegate.rowIdVector, delegate.sequenceNumberVector); + } + + private static class ReusingColumnarReader implements FileRecordReader { + + private final HeapLongVector dataVector = new HeapLongVector(1); + private final HeapLongVector rowIdVector = new HeapLongVector(1); + private final HeapLongVector sequenceNumberVector = new HeapLongVector(1); + private final ColumnVector[] readerOwnedColumns = + new HeapLongVector[] {dataVector, rowIdVector, sequenceNumberVector}; + private final VectorizedColumnBatch batch = new VectorizedColumnBatch(readerOwnedColumns); + private final VectorizedRowIterator iterator; + private int nextPosition; + private int recycleCount; + + private ReusingColumnarReader() { + dataVector.setLong(0, 42L); + rowIdVector.fillWithNulls(); + sequenceNumberVector.fillWithNulls(); + batch.setNumRows(1); + iterator = + new VectorizedRowIterator( + new Path("test"), new ColumnarRow(batch), () -> recycleCount++); + } + + @Override + public FileRecordIterator readBatch() { + iterator.reset(nextPosition++); + return iterator; + } + + @Override + public void close() {} + } +}