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 913f3d16a4 [Spark] Support reading nested null fields for Parquet 
struct (#9364)
913f3d16a4 is described below

commit 913f3d16a4eef7593b15e33998f271a61824b104
Author: Arnav Balyan <[email protected]>
AuthorDate: Sun Aug 23 20:47:24 2026 +0530

    [Spark] Support reading nested null fields for Parquet struct (#9364)
---
 .../format/parquet/ParquetReaderFactory.java       | 98 +++++++++++++++++++++-
 .../format/parquet/reader/ParquetColumnVector.java | 25 ++++--
 .../format/parquet/reader/ParquetReaderUtil.java   | 24 ++++++
 .../reader/VectorizedParquetRecordReader.java      |  5 ++
 .../format/parquet/ParquetReadWriteTest.java       | 62 ++++++++++++++
 5 files changed, 203 insertions(+), 11 deletions(-)

diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
index 5f6873fe34..2874849bb5 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
@@ -54,8 +54,10 @@ import org.apache.parquet.io.ColumnIOFactory;
 import org.apache.parquet.io.MessageColumnIO;
 import org.apache.parquet.schema.ConversionPatterns;
 import org.apache.parquet.schema.GroupType;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
 import org.apache.parquet.schema.MessageType;
 import org.apache.parquet.schema.OriginalType;
+import org.apache.parquet.schema.PrimitiveType;
 import org.apache.parquet.schema.Type;
 import org.apache.parquet.schema.Types;
 import org.slf4j.Logger;
@@ -291,16 +293,20 @@ public class ParquetReaderFactory implements 
FormatReaderFactory {
                 RowType rowType = (RowType) readType;
                 GroupType rowGroup = (GroupType) parquetType;
                 List<Type> rowGroupFields = new ArrayList<>();
+                boolean allFieldsMissing = true;
                 for (DataField field : rowType.getFields()) {
                     String fieldName = field.name();
                     Type type = matchParquetField(rowGroup, fieldName);
                     if (type != null) {
+                        allFieldsMissing = false;
                         rowGroupFields.add(clipParquetType(field.type(), 
type));
                     } else {
-                        // todo: support nested field missing
-                        throw new RuntimeException("field " + fieldName + " is 
missing");
+                        
rowGroupFields.add(ParquetSchemaConverter.convertToParquetType(field));
                     }
                 }
+                if (allFieldsMissing && rowGroup.getFieldCount() > 0) {
+                    rowGroupFields.add(findCheapestGroupField(rowGroup));
+                }
                 return rowGroup.withNewFields(rowGroupFields);
             case MAP:
                 MapType mapType = (MapType) readType;
@@ -358,6 +364,94 @@ public class ParquetReaderFactory implements 
FormatReaderFactory {
         }
     }
 
+    private Type findCheapestGroupField(GroupType groupType) {
+        return findCheapestField(groupType, 0).type.asGroupType().getType(0);
+    }
+
+    private CheapestField findCheapestField(Type type, int repetitionLevel) {
+        if (type.isPrimitive()) {
+            PrimitiveType.PrimitiveTypeName typeName =
+                    type.asPrimitiveType().getPrimitiveTypeName();
+            int cost;
+            switch (typeName) {
+                case BOOLEAN:
+                    cost = 1;
+                    break;
+                case INT32:
+                case FLOAT:
+                    cost = 4;
+                    break;
+                case INT64:
+                case DOUBLE:
+                    cost = 8;
+                    break;
+                case INT96:
+                    cost = 12;
+                    break;
+                default:
+                    cost = 32;
+            }
+            return new CheapestField(type, repetitionLevel, cost);
+        }
+
+        GroupType groupType = type.asGroupType();
+        LogicalTypeAnnotation annotation = 
groupType.getLogicalTypeAnnotation();
+        if (annotation instanceof 
LogicalTypeAnnotation.MapLogicalTypeAnnotation
+                || annotation instanceof 
LogicalTypeAnnotation.MapKeyValueTypeAnnotation) {
+            Preconditions.checkArgument(
+                    groupType.getFieldCount() == 1 && 
!groupType.getType(0).isPrimitive(),
+                    "Invalid map type: %s",
+                    groupType);
+            GroupType keyValueType = groupType.getType(0).asGroupType();
+            Preconditions.checkArgument(
+                    keyValueType.getRepetition() == Type.Repetition.REPEATED
+                            && keyValueType.getFieldCount() == 2,
+                    "Invalid map type: %s",
+                    groupType);
+            CheapestField key = findCheapestField(keyValueType.getType(0), 
repetitionLevel + 1);
+            CheapestField value = findCheapestField(keyValueType.getType(1), 
repetitionLevel + 1);
+            GroupType clippedKeyValue =
+                    keyValueType.withNewFields(Arrays.asList(key.type, 
value.type));
+            return new CheapestField(
+                    
groupType.withNewFields(Collections.singletonList(clippedKeyValue)),
+                    Math.max(key.repetitionLevel, value.repetitionLevel),
+                    key.cost + value.cost);
+        }
+
+        CheapestField cheapest = null;
+        for (Type child : groupType.getFields()) {
+            int childRepetitionLevel =
+                    repetitionLevel + (child.getRepetition() == 
Type.Repetition.REPEATED ? 1 : 0);
+            if (cheapest == null || childRepetitionLevel <= 
cheapest.repetitionLevel) {
+                CheapestField candidate = findCheapestField(child, 
childRepetitionLevel);
+                if (cheapest == null
+                        || candidate.repetitionLevel < cheapest.repetitionLevel
+                        || (candidate.repetitionLevel == 
cheapest.repetitionLevel
+                                && candidate.cost < cheapest.cost)) {
+                    cheapest = candidate;
+                }
+            }
+        }
+        Preconditions.checkNotNull(cheapest, "Parquet group must contain at 
least one field.");
+        return new CheapestField(
+                
groupType.withNewFields(Collections.singletonList(cheapest.type)),
+                cheapest.repetitionLevel,
+                cheapest.cost);
+    }
+
+    private static class CheapestField {
+
+        private final Type type;
+        private final int repetitionLevel;
+        private final int cost;
+
+        private CheapestField(Type type, int repetitionLevel, int cost) {
+            this.type = type;
+            this.repetitionLevel = repetitionLevel;
+            this.cost = cost;
+        }
+    }
+
     /**
      * Compute the batch size to use for the given file. Subclasses can 
override this to implement
      * dynamic per-file batch sizing based on footer metadata. The default 
implementation returns
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetColumnVector.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetColumnVector.java
index 667f741b54..4351976599 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetColumnVector.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetColumnVector.java
@@ -31,6 +31,8 @@ import java.util.ArrayList;
 import java.util.List;
 import java.util.Set;
 
+import static 
org.apache.paimon.data.columnar.ColumnVectorUtils.createParquetWritableColumnVector;
+
 /** Parquet Column tree. */
 public class ParquetColumnVector {
     private final ParquetField column;
@@ -48,6 +50,8 @@ public class ParquetColumnVector {
     /** Whether this column is primitive (i.e., leaf column). */
     private final boolean isPrimitive;
 
+    private final boolean isMissing;
+
     /** Reader for this column - only set if 'isPrimitive' is true. */
     private VectorizedColumnReader columnReader;
 
@@ -61,8 +65,9 @@ public class ParquetColumnVector {
         this.vector = vector;
         this.children = new ArrayList<>();
         this.isPrimitive = column.isPrimitive();
+        this.isMissing = missingColumns.contains(column);
 
-        if (missingColumns.contains(column)) {
+        if (isMissing) {
             vector.setAllNull();
             return;
         }
@@ -78,17 +83,19 @@ public class ParquetColumnVector {
         } else {
             ParquetGroupField groupField = (ParquetGroupField) column;
             Preconditions.checkArgument(
-                    groupField.getChildren().size() == 
vector.getChildren().length);
+                    groupField.getChildren().size() == 
vector.getChildren().length
+                            || groupField.getChildren().size() == 
vector.getChildren().length + 1);
             boolean allChildrenAreMissing = true;
 
             for (int i = 0; i < groupField.getChildren().size(); i++) {
+                ParquetField child = groupField.getChildren().get(i);
+                WritableColumnVector childVector =
+                        i < vector.getChildren().length
+                                ? (WritableColumnVector) 
vector.getChildren()[i]
+                                : createParquetWritableColumnVector(capacity, 
child.getType());
                 ParquetColumnVector childCv =
                         new ParquetColumnVector(
-                                groupField.getChildren().get(i),
-                                (WritableColumnVector) vector.getChildren()[i],
-                                capacity,
-                                missingColumns,
-                                false);
+                                child, childVector, capacity, missingColumns, 
false);
                 children.add(childCv);
 
                 // Only use levels from non-missing child, this can happen if 
only some but not all
@@ -165,8 +172,8 @@ public class ParquetColumnVector {
      * children.
      */
     void reset() {
-        // nothing to do if the column itself is missing
-        if (vector.isAllNull()) {
+        if (isMissing) {
+            vector.setAllNull();
             return;
         }
 
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetReaderUtil.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetReaderUtil.java
index afcdf43670..60c29f55a9 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetReaderUtil.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetReaderUtil.java
@@ -45,7 +45,9 @@ import org.apache.parquet.schema.Type;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.Objects;
+import java.util.stream.Collectors;
 
+import static 
org.apache.paimon.format.parquet.ParquetSchemaConverter.convertToPaimonField;
 import static 
org.apache.paimon.format.parquet.ParquetSchemaConverter.parquetListElementType;
 import static 
org.apache.paimon.format.parquet.ParquetSchemaConverter.parquetMapKeyValueType;
 import static org.apache.parquet.schema.Type.Repetition.REPEATED;
@@ -89,6 +91,16 @@ public class ParquetReaderUtil {
                                 lookupColumnByName(groupColumnIO, childName),
                                 getTypeIgnoreCase(parquetType.asGroupType(), 
childName)));
             }
+            GroupType parquetGroup = parquetType.asGroupType();
+            for (int i = children.size(); i < parquetGroup.getFieldCount(); 
i++) {
+                Type extraType = parquetGroup.getType(i);
+                DataField extraField = 
convertToPaimonField(addFallbackFieldIds(extraType));
+                fieldsBuilder.add(
+                        constructField(
+                                extraField,
+                                lookupColumnByName(groupColumnIO, 
extraType.getName()),
+                                extraType));
+            }
 
             return new ParquetGroupField(
                     type,
@@ -199,6 +211,18 @@ public class ParquetReaderUtil {
                 primitiveColumnIO.getFieldPath());
     }
 
+    private static Type addFallbackFieldIds(Type type) {
+        Type result = type;
+        if (!type.isPrimitive()) {
+            List<Type> children =
+                    type.asGroupType().getFields().stream()
+                            .map(ParquetReaderUtil::addFallbackFieldIds)
+                            .collect(Collectors.toList());
+            result = type.asGroupType().withNewFields(children);
+        }
+        return result.getId() == null ? result.withId(0) : result;
+    }
+
     /**
      * Parquet's column names are case insensitive. So when we look up columns 
we first check for
      * exact match, and if that can not find we look for a case-insensitive 
match.
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
index da306e4e11..cde21ec2d7 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
@@ -21,6 +21,7 @@ package org.apache.paimon.format.parquet.reader;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.data.columnar.writable.WritableColumnVector;
 import org.apache.paimon.format.parquet.type.ParquetField;
+import org.apache.paimon.format.parquet.type.ParquetGroupField;
 import org.apache.paimon.format.parquet.type.ParquetPrimitiveField;
 import org.apache.paimon.fs.FileIO;
 import org.apache.paimon.fs.Path;
@@ -162,6 +163,10 @@ public class VectorizedParquetRecordReader implements 
FileRecordReader<InternalR
                 if (!fd.equals(desc)) {
                     throw new IOException("Schema evolution not supported.");
                 }
+            } else {
+                for (ParquetField child : ((ParquetGroupField) 
field).getChildren()) {
+                    checkColumn(child);
+                }
             }
         } else {
             if (field.isRequired()) {
diff --git 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
index ccb793c6a6..9d471df92b 100644
--- 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
+++ 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
@@ -206,6 +206,68 @@ public class ParquetReadWriteTest {
         return Arrays.asList(10, 1000);
     }
 
+    @Test
+    void testMissingNestedFieldPreservesStructNullability() throws IOException 
{
+        RowType writeType =
+                RowType.builder()
+                        .field("id", new BigIntType())
+                        .field("s", RowType.builder().field("a", new 
BigIntType()).build())
+                        .field(
+                                "map_s",
+                                RowType.builder()
+                                        .field("m", new MapType(new 
BooleanType(), new IntType()))
+                                        .build())
+                        .field("partial_s", RowType.builder().field("a", new 
BigIntType()).build())
+                        .build();
+        List<InternalRow> records =
+                Arrays.asList(
+                        GenericRow.of(
+                                1L,
+                                GenericRow.of(10L),
+                                GenericRow.of(new 
GenericMap(Collections.singletonMap(true, 10))),
+                                GenericRow.of(20L)),
+                        GenericRow.of(2L, null, null, null));
+        Path path = createTempParquetFileByPaimon(folder, records, 10_000, 
writeType);
+
+        RowType readType =
+                RowType.builder()
+                        .field("id", new BigIntType())
+                        .field("s", RowType.builder().field("b", new 
BigIntType()).build())
+                        .field("map_s", RowType.builder().field("b", new 
BigIntType()).build())
+                        .field(
+                                "partial_s",
+                                RowType.builder()
+                                        .field("a", new BigIntType())
+                                        .field("b", new BigIntType())
+                                        .build())
+                        .build();
+        ParquetReaderFactory factory =
+                new ParquetReaderFactory(new Options(), readType, 1024, null);
+        LocalFileIO fileIO = new LocalFileIO();
+        InternalRowSerializer serializer = new InternalRowSerializer(readType);
+        List<InternalRow> results = new ArrayList<>();
+        try (RecordReader<InternalRow> reader =
+                factory.createReader(
+                        new FormatReaderContext(
+                                fileIO, path, fileIO.getFileSize(path), null, 
null))) {
+            reader.forEachRemaining(row -> results.add(serializer.copy(row)));
+        }
+
+        assertThat(results).hasSize(2);
+        assertThat(results.get(0).getLong(0)).isEqualTo(1L);
+        assertThat(results.get(0).isNullAt(1)).isFalse();
+        assertThat(results.get(0).getRow(1, 1).isNullAt(0)).isTrue();
+        assertThat(results.get(0).isNullAt(2)).isFalse();
+        assertThat(results.get(0).getRow(2, 1).isNullAt(0)).isTrue();
+        assertThat(results.get(0).isNullAt(3)).isFalse();
+        assertThat(results.get(0).getRow(3, 2).getLong(0)).isEqualTo(20L);
+        assertThat(results.get(0).getRow(3, 2).isNullAt(1)).isTrue();
+        assertThat(results.get(1).getLong(0)).isEqualTo(2L);
+        assertThat(results.get(1).isNullAt(1)).isTrue();
+        assertThat(results.get(1).isNullAt(2)).isTrue();
+        assertThat(results.get(1).isNullAt(3)).isTrue();
+    }
+
     @Test
     void testDynamicReadBatchSize() throws IOException {
         List<InternalRow> records = new ArrayList<>();

Reply via email to