This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 55063fab6c [common] Stop nesting row-tracking vector wrappers per 
batch (#9634)
55063fab6c is described below

commit 55063fab6c7e0ed0d4f091c633c1712387c22887
Author: YangJie <[email protected]>
AuthorDate: Thu Sep 10 02:52:06 2026 -0400

    [common] Stop nesting row-tracking vector wrappers per batch (#9634)
---
 .../paimon/data/columnar/ColumnarRowIterator.java  | 34 +++++++++++---
 .../data/columnar/ColumnarRowIteratorTest.java     | 53 ++++++++++++++++++++++
 2 files changed, 81 insertions(+), 6 deletions(-)

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 49a0fe5e71..4489388e34 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
@@ -129,16 +129,38 @@ public class ColumnarRowIterator extends 
RecyclableIterator<InternalRow>
         return this;
     }
 
+    /**
+     * Strips a row-tracking wrapper previously installed by {@link 
#assignRowTracking}, so repeated
+     * assignment re-wraps the base vector instead of nesting.
+     */
+    private static ColumnVector unwrapLong(ColumnVector vector) {
+        return vector instanceof TrackingLongColumnVector
+                ? ((TrackingLongColumnVector) vector).base
+                : vector;
+    }
+
+    /** Row-tracking wrapper installed by {@link #assignRowTracking}. */
+    private abstract static class TrackingLongColumnVector implements 
LongColumnVector {
+        final ColumnVector base;
+
+        TrackingLongColumnVector(ColumnVector base) {
+            this.base = base;
+        }
+    }
+
     public ColumnarRowIterator assignRowTracking(
             Long firstRowId, Long snapshotId, Map<String, Integer> meta) {
         VectorizedColumnBatch vectorizedColumnBatch = row.batch();
         ColumnVector[] vectors = vectorizedColumnBatch.columns;
 
-        if (meta.containsKey(SpecialFields.ROW_ID.name())) {
+        if (meta.containsKey(SpecialFields.ROW_ID.name()) && firstRowId != 
null) {
             Integer index = meta.get(SpecialFields.ROW_ID.name());
-            final ColumnVector rowIdVector = vectors[index];
+            // Wrap the base vector once: assignRowTracking runs per batch and 
the
+            // wrapped vector persists, so re-wrapping would nest one 
delegation
+            // level per batch (O(batches) depth, every read O(depth)).
+            final ColumnVector rowIdVector = unwrapLong(vectors[index]);
             vectors[index] =
-                    new LongColumnVector() {
+                    new TrackingLongColumnVector(rowIdVector) {
                         @Override
                         public long getLong(int i) {
                             if (rowIdVector.isNullAt(i)) {
@@ -155,11 +177,11 @@ public class ColumnarRowIterator extends 
RecyclableIterator<InternalRow>
                     };
         }
 
-        if (meta.containsKey(SpecialFields.SEQUENCE_NUMBER.name())) {
+        if (meta.containsKey(SpecialFields.SEQUENCE_NUMBER.name()) && 
snapshotId != null) {
             Integer index = meta.get(SpecialFields.SEQUENCE_NUMBER.name());
-            final ColumnVector versionVector = vectors[index];
+            final ColumnVector versionVector = unwrapLong(vectors[index]);
             vectors[index] =
-                    new LongColumnVector() {
+                    new TrackingLongColumnVector(versionVector) {
                         @Override
                         public long getLong(int i) {
                             if (versionVector.isNullAt(i)) {
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 8ab926a193..8734886d37 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
@@ -19,11 +19,15 @@
 package org.apache.paimon.data.columnar;
 
 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.utils.LongIterator;
 
 import org.junit.jupiter.api.Test;
 
+import java.util.HashMap;
+import java.util.Map;
 import java.util.Random;
 
 import static org.assertj.core.api.Assertions.assertThat;
@@ -62,4 +66,53 @@ public class ColumnarRowIteratorTest {
             
assertThat(rowIterator.returnedPosition()).isEqualTo(positions[rowIterator.index
 - 1]);
         }
     }
+
+    @Test
+    public void testAssignRowTrackingSkipsNullFirstRowId() {
+        Map<String, Integer> meta = new HashMap<>();
+        meta.put(SpecialFields.ROW_ID.name(), 0);
+
+        // A file without stored _rowid gives a null firstRowId: the vector 
must be left
+        // alone rather than wrapped around an unboxing of null.
+        HeapLongVector base = new HeapLongVector(4);
+        VectorizedColumnBatch batch = new VectorizedColumnBatch(new 
ColumnVector[] {base});
+        batch.setNumRows(4);
+        ColumnarRowIterator nullFirst =
+                new ColumnarRowIterator(new Path("t"), new ColumnarRow(batch), 
null);
+        nullFirst.reset(LongIterator.fromArray(new long[] {0, 1, 2, 3}));
+        ColumnarRowIterator untouched = nullFirst.assignRowTracking(null, 7L, 
meta);
+        assertThat(untouched.next()).isNotNull();
+        assertThat(untouched.batch().columns[0]).isSameAs(base);
+    }
+
+    @Test
+    public void testRepeatedAssignRowTrackingDoesNotNest() {
+        Map<String, Integer> meta = new HashMap<>();
+        meta.put(SpecialFields.ROW_ID.name(), 0);
+
+        HeapLongVector withIds = new HeapLongVector(4);
+        withIds.setNullAt(0); // row 0 falls back to firstRowId + position
+        withIds.setLong(1, 100);
+        withIds.setLong(2, 200);
+        withIds.setLong(3, 300);
+        VectorizedColumnBatch batch2 = new VectorizedColumnBatch(new 
ColumnVector[] {withIds});
+        batch2.setNumRows(4);
+        ColumnarRowIterator tracked =
+                new ColumnarRowIterator(new Path("t"), new 
ColumnarRow(batch2), null);
+        tracked.reset(LongIterator.fromArray(new long[] {0, 1, 2, 3}));
+        // One assignment per batch, and the wrapped vector outlives the 
batch. Nesting a
+        // wrapper each time makes every read recurse once per assignment, so 
this many
+        // rounds overflows the stack rather than merely being slow.
+        for (int i = 0; i < 100_000; i++) {
+            tracked = tracked.assignRowTracking(10L, 7L, meta);
+        }
+        assertThat(tracked.next()).isNotNull();
+        assertThat(tracked.row.getLong(0)).isEqualTo(10L);
+        assertThat(tracked.next()).isNotNull();
+        assertThat(tracked.row.getLong(0)).isEqualTo(100L);
+        assertThat(tracked.next()).isNotNull();
+        assertThat(tracked.row.getLong(0)).isEqualTo(200L);
+        assertThat(tracked.next()).isNotNull();
+        assertThat(tracked.row.getLong(0)).isEqualTo(300L);
+    }
 }

Reply via email to