From d53ff92c8db26f16ebd1dffc88520e23a528aeb5 Mon Sep 17 00:00:00 2001 From: mingfeng Date: Tue, 8 Sep 2026 22:32:37 -0700 Subject: [PATCH 1/5] [common] Preserve specialized iterators for full identity mappings --- .../data/columnar/ColumnarRowIterator.java | 17 +++++ .../columnar/ColumnarRowIteratorTest.java | 67 +++++++++++++++++++ 2 files changed, 84 insertions(+) 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 49a0fe5e714f..2ba712cc8228 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; + } + public ColumnarRowIterator assignRowTracking( Long firstRowId, Long snapshotId, Map meta) { VectorizedColumnBatch vectorizedColumnBatch = row.batch(); 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 8ab926a193cd..bf6655e2a9b7 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,8 +18,13 @@ 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.fs.Path; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; import org.apache.paimon.utils.LongIterator; import org.junit.jupiter.api.Test; @@ -62,4 +67,66 @@ public void testRowIterator() { assertThat(rowIterator.returnedPosition()).isEqualTo(positions[rowIterator.index - 1]); } } + + @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 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); + } + } } From 151f38b6a91830770caf14f55e3a1b8bf690e1a5 Mon Sep 17 00:00:00 2001 From: mingfeng Date: Wed, 9 Sep 2026 21:20:41 -0700 Subject: [PATCH 2/5] [core] Isolate row tracking from reusable column batches --- .../paimon/io/DataFileRecordReader.java | 8 +- .../paimon/io/DataFileRecordReaderTest.java | 115 ++++++++++++++++++ 2 files changed, 122 insertions(+), 1 deletion(-) create mode 100644 paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java 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..138a2407f14d 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 @@ -178,8 +178,14 @@ private FileRecordIterator readBatchInternal() throws IOException { } if (iterator instanceof ColumnarRowIterator) { - iterator = ((ColumnarRowIterator) iterator).mapping(partitionInfo, indexMapping); + ColumnarRowIterator sourceIterator = (ColumnarRowIterator) iterator; + iterator = sourceIterator.mapping(partitionInfo, indexMapping); if (rowTrackingEnabled) { + if (iterator == sourceIterator) { + // Row tracking replaces columns in place, so isolate reusable reader batches. + ColumnarRowIterator columnarIterator = (ColumnarRowIterator) iterator; + iterator = columnarIterator.copy(columnarIterator.batch().columns.clone()); + } 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..862362143180 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java @@ -0,0 +1,115 @@ +/* + * 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.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.RowType; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +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(); + Map systemFields = new HashMap<>(); + systemFields.put(SpecialFields.ROW_ID.name(), 0); + systemFields.put(SpecialFields.SEQUENCE_NUMBER.name(), 1); + DataFileRecordReader reader = + new DataFileRecordReader( + new RowType( + Arrays.asList(SpecialFields.ROW_ID, SpecialFields.SEQUENCE_NUMBER)), + delegate, + false, + false, + new int[] {0, 1}, + null, + null, + true, + 100L, + 7L, + systemFields, + null, + new Path("test")); + + for (int batchIndex = 0; batchIndex < 3; batchIndex++) { + FileRecordIterator iterator = reader.readBatch(); + + assertThat(iterator).isNotSameAs(delegate.iterator); + assertThat(iterator).isInstanceOf(VectorizedRowIterator.class); + assertThat(delegate.batch.columns) + .containsExactly(delegate.rowIdVector, delegate.sequenceNumberVector); + + InternalRow row = iterator.next(); + assertThat(row.getLong(0)).isEqualTo(100L + batchIndex); + assertThat(row.getLong(1)).isEqualTo(7L); + + iterator.releaseBatch(); + assertThat(delegate.recycleCount).isEqualTo(batchIndex + 1); + assertThat(delegate.batch.columns) + .containsExactly(delegate.rowIdVector, delegate.sequenceNumberVector); + } + reader.close(); + } + + private static class ReusingColumnarReader implements FileRecordReader { + + private final HeapLongVector rowIdVector = new HeapLongVector(1); + private final HeapLongVector sequenceNumberVector = new HeapLongVector(1); + private final VectorizedColumnBatch batch = + new VectorizedColumnBatch(new ColumnVector[] {rowIdVector, sequenceNumberVector}); + private final VectorizedRowIterator iterator; + private int nextPosition; + private int recycleCount; + + private ReusingColumnarReader() { + 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() {} + } +} From 3c2c993e6e7fab3505e5f4643ab200911fff39d3 Mon Sep 17 00:00:00 2001 From: mingfeng Date: Thu, 10 Sep 2026 00:30:36 -0700 Subject: [PATCH 3/5] [core] Avoid covariant array row-tracking failure --- .../org/apache/paimon/io/DataFileRecordReader.java | 11 +++++++++-- .../apache/paimon/io/DataFileRecordReaderTest.java | 13 +++++++++---- 2 files changed, 18 insertions(+), 6 deletions(-) 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 138a2407f14d..0330e88f6b39 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; @@ -182,9 +183,15 @@ private FileRecordIterator readBatchInternal() throws IOException { iterator = sourceIterator.mapping(partitionInfo, indexMapping); if (rowTrackingEnabled) { if (iterator == sourceIterator) { - // Row tracking replaces columns in place, so isolate reusable reader batches. + // 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(columnarIterator.batch().columns.clone()); + iterator = + columnarIterator.copy( + Arrays.copyOf( + columnarIterator.batch().columns, + columnarIterator.batch().columns.length, + ColumnVector[].class)); } iterator = ((ColumnarRowIterator) iterator) 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 index 862362143180..9a10f3ae92b1 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java @@ -44,6 +44,9 @@ 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(), 0); systemFields.put(SpecialFields.SEQUENCE_NUMBER.name(), 1); @@ -69,16 +72,17 @@ public void testRowTrackingIsolatedFromReusedIdentityBatch() throws Exception { assertThat(iterator).isNotSameAs(delegate.iterator); assertThat(iterator).isInstanceOf(VectorizedRowIterator.class); - assertThat(delegate.batch.columns) + assertThat(delegate.readerOwnedColumns) .containsExactly(delegate.rowIdVector, delegate.sequenceNumberVector); InternalRow row = iterator.next(); assertThat(row.getLong(0)).isEqualTo(100L + batchIndex); assertThat(row.getLong(1)).isEqualTo(7L); + assertThat(delegate.recycleCount).isEqualTo(batchIndex); iterator.releaseBatch(); assertThat(delegate.recycleCount).isEqualTo(batchIndex + 1); - assertThat(delegate.batch.columns) + assertThat(delegate.readerOwnedColumns) .containsExactly(delegate.rowIdVector, delegate.sequenceNumberVector); } reader.close(); @@ -88,8 +92,9 @@ private static class ReusingColumnarReader implements FileRecordReader Date: Thu, 10 Sep 2026 02:57:55 -0700 Subject: [PATCH 4/5] [core] Preserve identity iterators without tracking fields --- .../paimon/io/DataFileRecordReader.java | 2 +- .../paimon/io/DataFileRecordReaderTest.java | 43 +++++++++++++++++++ 2 files changed, 44 insertions(+), 1 deletion(-) 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 0330e88f6b39..22c6faf1baf0 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 @@ -181,7 +181,7 @@ private FileRecordIterator readBatchInternal() throws IOException { if (iterator instanceof ColumnarRowIterator) { ColumnarRowIterator sourceIterator = (ColumnarRowIterator) iterator; iterator = sourceIterator.mapping(partitionInfo, indexMapping); - if (rowTrackingEnabled) { + if (rowTrackingEnabled && !systemFields.isEmpty()) { if (iterator == sourceIterator) { // Copy to a ColumnVector[] because cloning a subtype array preserves its // runtime type and cannot accept row-tracking wrapper vectors. 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 index 9a10f3ae92b1..74e265cfcd04 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java @@ -21,6 +21,7 @@ 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; @@ -28,11 +29,13 @@ 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.Arrays; +import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -88,6 +91,46 @@ public void testRowTrackingIsolatedFromReusedIdentityBatch() throws Exception { 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(); + } + private static class ReusingColumnarReader implements FileRecordReader { private final HeapLongVector rowIdVector = new HeapLongVector(1); From 53b5ad988f23e6b26b74b06bb16b8ad137c938b6 Mon Sep 17 00:00:00 2001 From: mingfeng Date: Thu, 10 Sep 2026 05:15:06 -0700 Subject: [PATCH 5/5] [common][core] Strengthen identity row-tracking tests --- .../columnar/ColumnarRowIteratorTest.java | 4 + .../paimon/io/DataFileRecordReaderTest.java | 95 ++++++++++++++----- 2 files changed, 75 insertions(+), 24 deletions(-) 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 bf6655e2a9b7..7b14a80f1224 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 @@ -95,6 +95,10 @@ public void testNonIdentityMappingCopiesIterator() { 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); 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 index 74e265cfcd04..6b655c96fc5e 100644 --- a/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/io/DataFileRecordReaderTest.java @@ -34,7 +34,6 @@ import org.junit.jupiter.api.Test; -import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -51,46 +50,64 @@ public void testRowTrackingIsolatedFromReusedIdentityBatch() throws Exception { assertThat(delegate.readerOwnedColumns.getClass()).isEqualTo(HeapLongVector[].class); Map systemFields = new HashMap<>(); - systemFields.put(SpecialFields.ROW_ID.name(), 0); - systemFields.put(SpecialFields.SEQUENCE_NUMBER.name(), 1); + systemFields.put(SpecialFields.ROW_ID.name(), 1); + systemFields.put(SpecialFields.SEQUENCE_NUMBER.name(), 2); DataFileRecordReader reader = - new DataFileRecordReader( - new RowType( - Arrays.asList(SpecialFields.ROW_ID, SpecialFields.SEQUENCE_NUMBER)), + createRowTrackingReader( + SpecialFields.rowTypeWithRowTracking(RowType.of(DataTypes.BIGINT())), delegate, - false, - false, - new int[] {0, 1}, - null, - null, - true, - 100L, - 7L, - systemFields, - null, - new Path("test")); + 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); - assertThat(delegate.readerOwnedColumns) - .containsExactly(delegate.rowIdVector, delegate.sequenceNumberVector); + assertReaderOwnedColumns(delegate); InternalRow row = iterator.next(); - assertThat(row.getLong(0)).isEqualTo(100L + batchIndex); - assertThat(row.getLong(1)).isEqualTo(7L); + 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); - assertThat(delegate.readerOwnedColumns) - .containsExactly(delegate.rowIdVector, delegate.sequenceNumberVector); + 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); @@ -131,18 +148,48 @@ public void close() {} 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[] {rowIdVector, sequenceNumberVector}; + 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);