From 1c3b1c6cdbee5379dfd02df395c224f7b0648016 Mon Sep 17 00:00:00 2001 From: Xiening Dai Date: Sat, 8 Aug 2026 18:08:13 +0000 Subject: [PATCH] Parquet: Fix null counts for nested fields of a null struct When an optional struct is null, OptionWriter writes a null directly to every leaf column it contains, so the writers for those columns never see the value and cannot count it. OptionWriter dropped its own null count in that case, with a comment saying nested null stats were not used. They are used: the counts it returns become DataFile.nullValueCounts. Float, double, geometry and geography are the types whose writers report metrics, and ParquetMetrics prefers writer metrics over footer statistics, so the correct footer count was never consulted. Such a field under a nullable struct was reported as having 0 nulls even when the struct was null for some rows. Required fields are affected too, since only optional fields are wrapped in an option writer. The incorrect counting of nulls could affect query engines that relay on this stats for optimization. For example, they could simply skip the file with null_count == 0 for predicate `WHERE c.f_id IS NULL` and produces wrong result. Add the nulls counted by an option writer to the metrics of the columns it wrote them to, at any depth. And add corresponding tests. --- .../iceberg/parquet/ParquetValueWriters.java | 35 +++-- .../parquet/TestParquetValueWriters.java | 121 ++++++++++++++++++ 2 files changed, 142 insertions(+), 14 deletions(-) diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueWriters.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueWriters.java index 298ffa121585..c931207809f9 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueWriters.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueWriters.java @@ -474,17 +474,7 @@ public Stream> metrics() { // we are not tracking field metrics for this type ourselves return Stream.empty(); } else if (fieldMetricsFromWriter.size() == 1) { - FieldMetrics metrics = fieldMetricsFromWriter.get(0); - return Stream.of( - new FieldMetrics<>( - metrics.id(), - metrics.valueCount() + nullValueCount, - nullValueCount, - metrics.nanValueCount(), - metrics.lowerBound(), - metrics.upperBound(), - metrics.originalType(), - metrics.avgValueSizeInBytes())); + return Stream.of(withNullValues(fieldMetricsFromWriter.get(0))); } else { throw new IllegalStateException( String.format( @@ -494,9 +484,26 @@ public Stream> metrics() { } } - // skipping updating null stats for non-primitive types since we don't use them today, to - // avoid unnecessary work - return writer.metrics(); + // A null value here is also null for every descendant column, but those columns are written + // directly and never see it, so their writers cannot count it. Add it to their metrics. + return writer.metrics().map(this::withNullValues); + } + + /** Adds the nulls counted by this writer to metrics produced by a descendant column. */ + private FieldMetrics withNullValues(FieldMetrics metrics) { + if (nullValueCount == 0) { + return metrics; + } + + return new FieldMetrics<>( + metrics.id(), + metrics.valueCount() + nullValueCount, + metrics.nullValueCount() + nullValueCount, + metrics.nanValueCount(), + metrics.lowerBound(), + metrics.upperBound(), + metrics.originalType(), + metrics.avgValueSizeInBytes()); } } diff --git a/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetValueWriters.java b/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetValueWriters.java index 33c4044a0a7f..16a3a74688e6 100644 --- a/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetValueWriters.java +++ b/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetValueWriters.java @@ -19,13 +19,21 @@ package org.apache.iceberg.parquet; import static org.apache.iceberg.types.Types.NestedField.optional; +import static org.apache.iceberg.types.Types.NestedField.required; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; import java.nio.ByteBuffer; +import java.util.Map; +import java.util.function.Function; +import java.util.stream.Collectors; import org.apache.iceberg.FieldMetrics; import org.apache.iceberg.Schema; +import org.apache.iceberg.data.GenericRecord; +import org.apache.iceberg.data.Record; +import org.apache.iceberg.data.parquet.InternalWriter; import org.apache.iceberg.types.Types; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.ColumnWriteStore; @@ -60,4 +68,117 @@ void geospatialValueSizeMetricsExcludeNulls() { assertThat(metrics.nullValueCount()).isEqualTo(1); assertThat(metrics.avgValueSizeInBytes()).isEqualTo(31); } + + @Test + void nullStructCountsNullsForNestedFields() { + // a null struct is also null for the fields it contains, but those columns are written by the + // struct's writer and never see the value, so the struct must count the nulls for them + Types.StructType struct = + Types.StructType.of( + optional(2, "d", Types.DoubleType.get()), required(3, "f", Types.FloatType.get())); + Schema schema = new Schema(optional(1, "s", struct)); + + ParquetValueWriter writer = writerFor(schema); + Record inner = GenericRecord.create(struct); + inner.set(0, 2.0D); + inner.set(1, 1.0F); + + writer.write(0, record(schema, inner)); + writer.write(0, record(schema, null)); + writer.write(0, record(schema, null)); + + Map> metrics = metricsById(writer); + // both fields have one non-null value and two nulls from the null structs + assertThat(metrics.get(2).nullValueCount()).isEqualTo(2); + assertThat(metrics.get(2).valueCount()).isEqualTo(3); + assertThat(metrics.get(3).nullValueCount()).isEqualTo(2); + assertThat(metrics.get(3).valueCount()).isEqualTo(3); + } + + @Test + void nullStructAddsToNullsCountedByNestedField() { + Types.StructType struct = Types.StructType.of(optional(2, "d", Types.DoubleType.get())); + Schema schema = new Schema(optional(1, "s", struct)); + + ParquetValueWriter writer = writerFor(schema); + Record present = GenericRecord.create(struct); + present.set(0, 2.0D); + Record nullField = GenericRecord.create(struct); + nullField.set(0, null); + + writer.write(0, record(schema, present)); + // the field is null while the struct is present, so the field's own writer counts it + writer.write(0, record(schema, nullField)); + writer.write(0, record(schema, null)); + + Map> metrics = metricsById(writer); + assertThat(metrics.get(2).nullValueCount()).isEqualTo(2); + assertThat(metrics.get(2).valueCount()).isEqualTo(3); + } + + @Test + void nullStructCountsNullsForDeeplyNestedFields() { + Types.StructType inner = Types.StructType.of(optional(3, "d", Types.DoubleType.get())); + Types.StructType outer = Types.StructType.of(optional(2, "inner", inner)); + Schema schema = new Schema(optional(1, "s", outer)); + + ParquetValueWriter writer = writerFor(schema); + Record innerRecord = GenericRecord.create(inner); + innerRecord.set(0, 2.0D); + Record withInner = GenericRecord.create(outer); + withInner.set(0, innerRecord); + Record withoutInner = GenericRecord.create(outer); + withoutInner.set(0, null); + + writer.write(0, record(schema, withInner)); + // a null at either level is a null for the leaf + writer.write(0, record(schema, withoutInner)); + writer.write(0, record(schema, null)); + + Map> metrics = metricsById(writer); + assertThat(metrics.get(3).nullValueCount()).isEqualTo(2); + assertThat(metrics.get(3).valueCount()).isEqualTo(3); + } + + @Test + void nullStructCountsNullsForNestedGeospatialField() { + // geospatial writers also report metrics, so they are affected in the same way + Types.StructType struct = Types.StructType.of(optional(2, "g", Types.GeometryType.crs84())); + Schema schema = new Schema(optional(1, "s", struct)); + + ParquetValueWriter writer = writerFor(schema); + Record present = GenericRecord.create(struct); + present.set(0, ByteBuffer.allocate(21)); + + writer.write(0, record(schema, present)); + writer.write(0, record(schema, null)); + + Map> metrics = metricsById(writer); + assertThat(metrics.get(2).nullValueCount()).isEqualTo(1); + assertThat(metrics.get(2).valueCount()).isEqualTo(2); + // the size of the one non-null value is still reported + assertThat(metrics.get(2).avgValueSizeInBytes()).isEqualTo(21); + } + + private static Record record(Schema schema, Record struct) { + Record record = GenericRecord.create(schema); + record.set(0, struct); + return record; + } + + /** Returns a writer for the given schema, with a mocked column store. */ + private static ParquetValueWriter writerFor(Schema schema) { + MessageType parquetSchema = ParquetSchemaUtil.convert(schema, "table"); + ParquetValueWriter writer = InternalWriter.createWriter(schema, parquetSchema); + + ColumnWriteStore columnStore = mock(ColumnWriteStore.class); + when(columnStore.getColumnWriter(any())).thenReturn(mock(ColumnWriter.class)); + writer.setColumnStore(columnStore); + + return writer; + } + + private static Map> metricsById(ParquetValueWriter writer) { + return writer.metrics().collect(Collectors.toMap(FieldMetrics::id, Function.identity())); + } }