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 9c71ad1f42 [vortex] Read the primary-key format's physical 
_SEQUENCE_NUMBER column (#8876) (#8877)
9c71ad1f42 is described below

commit 9c71ad1f42b39fb3479cb48bca7fd2b3eb060fae
Author: Jordan Epstein <[email protected]>
AuthorDate: Tue Jul 28 08:10:25 2026 -0500

    [vortex] Read the primary-key format's physical _SEQUENCE_NUMBER column 
(#8876) (#8877)
---
 .../paimon/format/vortex/VortexRecordsReader.java  | 37 +++++++++++++++-------
 .../format/vortex/VortexRecordsReaderTest.java     | 28 ++++++++++++++--
 2 files changed, 52 insertions(+), 13 deletions(-)

diff --git 
a/paimon-vortex/paimon-vortex-format/src/main/java/org/apache/paimon/format/vortex/VortexRecordsReader.java
 
b/paimon-vortex/paimon-vortex-format/src/main/java/org/apache/paimon/format/vortex/VortexRecordsReader.java
index 571c4edf04..2d3cfcbbc8 100644
--- 
a/paimon-vortex/paimon-vortex-format/src/main/java/org/apache/paimon/format/vortex/VortexRecordsReader.java
+++ 
b/paimon-vortex/paimon-vortex-format/src/main/java/org/apache/paimon/format/vortex/VortexRecordsReader.java
@@ -76,7 +76,8 @@ public class VortexRecordsReader implements 
FileRecordReader<InternalRow> {
             Map<String, String> storageOptions) {
         this.filePath = path;
         RowType physicalReadRowType = physicalReadRowType(dataSchemaRowType, 
projectedRowType);
-        this.physicalFieldMapping = physicalFieldMapping(physicalReadRowType, 
projectedRowType);
+        this.physicalFieldMapping =
+                physicalFieldMapping(dataSchemaRowType, physicalReadRowType, 
projectedRowType);
         this.allocator =
                 ArrowAllocation.rootAllocator()
                         .newChildAllocator("vortex-reader", 0, Long.MAX_VALUE);
@@ -207,7 +208,7 @@ public class VortexRecordsReader implements 
FileRecordReader<InternalRow> {
 
     @VisibleForTesting
     static RowType physicalReadRowType(RowType dataSchemaRowType, RowType 
projectedRowType) {
-        if (!hasRowTrackingField(projectedRowType)) {
+        if (!hasSynthesizedRowTrackingField(dataSchemaRowType, 
projectedRowType)) {
             return projectedRowType;
         }
 
@@ -215,7 +216,7 @@ public class VortexRecordsReader implements 
FileRecordReader<InternalRow> {
         Set<Integer> selectedFieldIds = new HashSet<>();
         Set<String> selectedFieldNames = new HashSet<>();
         for (DataField projectedField : projectedRowType.getFields()) {
-            if (isRowTrackingField(projectedField)) {
+            if (isSynthesizedRowTrackingField(dataSchemaRowType, 
projectedField)) {
                 continue;
             }
 
@@ -234,15 +235,16 @@ public class VortexRecordsReader implements 
FileRecordReader<InternalRow> {
 
     @Nullable
     @VisibleForTesting
-    static int[] physicalFieldMapping(RowType physicalReadRowType, RowType 
projectedRowType) {
-        if (!hasRowTrackingField(projectedRowType)) {
+    static int[] physicalFieldMapping(
+            RowType dataSchemaRowType, RowType physicalReadRowType, RowType 
projectedRowType) {
+        if (!hasSynthesizedRowTrackingField(dataSchemaRowType, 
projectedRowType)) {
             return null;
         }
 
         int[] mapping = new int[projectedRowType.getFieldCount()];
         for (int i = 0; i < projectedRowType.getFieldCount(); i++) {
             DataField field = projectedRowType.getFields().get(i);
-            if (isRowTrackingField(field)) {
+            if (isSynthesizedRowTrackingField(dataSchemaRowType, field)) {
                 mapping[i] = -1;
             } else {
                 if (physicalReadRowType.containsField(field.id())) {
@@ -273,17 +275,30 @@ public class VortexRecordsReader implements 
FileRecordReader<InternalRow> {
         return null;
     }
 
-    private static boolean hasRowTrackingField(RowType rowType) {
+    private static boolean hasSynthesizedRowTrackingField(
+            RowType dataSchemaRowType, RowType rowType) {
         for (DataField field : rowType.getFields()) {
-            if (isRowTrackingField(field)) {
+            if (isSynthesizedRowTrackingField(dataSchemaRowType, field)) {
                 return true;
             }
         }
         return false;
     }
 
-    private static boolean isRowTrackingField(DataField field) {
-        return SpecialFields.ROW_ID.name().equals(field.name())
-                || SpecialFields.SEQUENCE_NUMBER.name().equals(field.name());
+    /**
+     * A row-tracking field is synthesized only when the file does not store 
it physically. {@code
+     * _SEQUENCE_NUMBER} doubles as the primary-key file format's physical 
sequence column: for
+     * key-value reads it must be read like any other data column, not mapped 
away for synthesis —
+     * doing so left its projection index at -1 and crashed every vortex read 
of a primary-key
+     * table.
+     */
+    private static boolean isSynthesizedRowTrackingField(
+            RowType dataSchemaRowType, DataField field) {
+        if (!SpecialFields.ROW_ID.name().equals(field.name())
+                && !SpecialFields.SEQUENCE_NUMBER.name().equals(field.name())) 
{
+            return false;
+        }
+        return !dataSchemaRowType.containsField(field.id())
+                && dataSchemaRowType.getFieldIndex(field.name()) < 0;
     }
 }
diff --git 
a/paimon-vortex/paimon-vortex-format/src/test/java/org/apache/paimon/format/vortex/VortexRecordsReaderTest.java
 
b/paimon-vortex/paimon-vortex-format/src/test/java/org/apache/paimon/format/vortex/VortexRecordsReaderTest.java
index abd0a2ece5..f9ceff705a 100644
--- 
a/paimon-vortex/paimon-vortex-format/src/test/java/org/apache/paimon/format/vortex/VortexRecordsReaderTest.java
+++ 
b/paimon-vortex/paimon-vortex-format/src/test/java/org/apache/paimon/format/vortex/VortexRecordsReaderTest.java
@@ -48,13 +48,36 @@ public class VortexRecordsReaderTest {
         RowType physicalReadRowType =
                 VortexRecordsReader.physicalReadRowType(dataSchemaRowType, 
projectedRowType);
         int[] physicalFieldMapping =
-                VortexRecordsReader.physicalFieldMapping(physicalReadRowType, 
projectedRowType);
+                VortexRecordsReader.physicalFieldMapping(
+                        dataSchemaRowType, physicalReadRowType, 
projectedRowType);
 
         assertEquals(1, physicalReadRowType.getFieldCount());
         assertEquals("f_string", 
physicalReadRowType.getFields().get(0).name());
         assertArrayEquals(new int[] {0, -1}, physicalFieldMapping);
     }
 
+    @Test
+    public void testPhysicalSequenceColumnIsReadNotSynthesized() {
+        // The primary-key file format stores _SEQUENCE_NUMBER as a physical 
column; it must map
+        // to its physical index, not be reserved for row-tracking synthesis.
+        RowType dataSchemaRowType =
+                RowType.of(
+                        new DataField(0, "k", DataTypes.INT()),
+                        SpecialFields.SEQUENCE_NUMBER,
+                        SpecialFields.VALUE_KIND,
+                        new DataField(1, "v", DataTypes.BIGINT()));
+        RowType projectedRowType = dataSchemaRowType;
+
+        RowType physicalReadRowType =
+                VortexRecordsReader.physicalReadRowType(dataSchemaRowType, 
projectedRowType);
+        int[] physicalFieldMapping =
+                VortexRecordsReader.physicalFieldMapping(
+                        dataSchemaRowType, physicalReadRowType, 
projectedRowType);
+
+        assertEquals(dataSchemaRowType, physicalReadRowType);
+        assertEquals(null, physicalFieldMapping);
+    }
+
     @Test
     public void testPhysicalFieldMappingPrefersFieldIdOverName() {
         RowType dataSchemaRowType =
@@ -67,7 +90,8 @@ public class VortexRecordsReaderTest {
         RowType physicalReadRowType =
                 VortexRecordsReader.physicalReadRowType(dataSchemaRowType, 
projectedRowType);
         int[] physicalFieldMapping =
-                VortexRecordsReader.physicalFieldMapping(physicalReadRowType, 
projectedRowType);
+                VortexRecordsReader.physicalFieldMapping(
+                        dataSchemaRowType, physicalReadRowType, 
projectedRowType);
 
         assertEquals(1, physicalReadRowType.getFieldCount());
         assertEquals("old_name", 
physicalReadRowType.getFields().get(0).name());

Reply via email to