From 77f436f5108c454af60dc4387967c38a026283e4 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 16:17:42 +0800 Subject: [PATCH] [core] Align null stats columns of partial writes by write schema MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit IcebergDataFileMeta.create interpreted null stats columns as covering the whole table schema, but a null value of valueStatsCols means the stats cover the whole WRITE schema. For a partial write the stats row follows the write-column order — which can differ from the table schema order, for example for a MERGE INTO whose SET clause lists columns differently — and nested writes record leaf paths, so bounds and null counts drifted onto the wrong fields in the Iceberg manifest. Derive the explicit stats column list from the write columns in their recorded order, mapping leaf paths to their top-level field. Assisted-by: GLM-5.3 --- .../paimon/iceberg/IcebergCommitCallback.java | 6 +- .../iceberg/manifest/IcebergDataFileMeta.java | 27 ++++- .../manifest/IcebergDataFileMetaTest.java | 103 +++++++++++++++++- 3 files changed, 132 insertions(+), 4 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java index ab731753969e..c423457287c6 100644 --- a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java +++ b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java @@ -674,7 +674,8 @@ private void dataSplitToManifestEntries( rawFile.fileSize(), schemaCache.get(paimonFileMeta.schemaId()), paimonFileMeta.valueStats(), - paimonFileMeta.valueStatsCols()); + paimonFileMeta.valueStatsCols(), + paimonFileMeta.writeCols()); dataFileEntries.add( new IcebergManifestEntry( IcebergManifestEntry.Status.ADDED, @@ -1370,7 +1371,8 @@ private List createNewlyAddedManifestFileMetas( paimonFileMeta.fileSize(), schemaCache.get(paimonFileMeta.schemaId()), paimonFileMeta.valueStats(), - paimonFileMeta.valueStatsCols()); + paimonFileMeta.valueStatsCols(), + paimonFileMeta.writeCols()); return new IcebergManifestEntry( IcebergManifestEntry.Status.ADDED, currentSnapshotId, diff --git a/paimon-core/src/main/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMeta.java index f9ce30bad4e3..b5a5cfd9f55e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMeta.java @@ -34,9 +34,11 @@ import java.util.ArrayList; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Set; /** * Iceberg data file meta. @@ -154,8 +156,31 @@ public static IcebergDataFileMeta create( long fileSizeInBytes, IcebergSchema icebergSchema, SimpleStats stats, - @Nullable List statsColumns) { + @Nullable List statsColumns, + @Nullable List writeCols) { int numFields = icebergSchema.fields().size(); + if (statsColumns == null && writeCols != null) { + // null stats columns mean the stats cover the whole WRITE schema: the stats row + // follows the write-column order (which may differ from the table schema order, + // for example for a MERGE INTO whose SET clause lists columns differently), and + // nested writes record leaf paths that map to their top-level field + Set icebergNames = new HashSet<>(); + for (IcebergDataField field : icebergSchema.fields()) { + icebergNames.add(field.name()); + } + statsColumns = new ArrayList<>(); + for (String writeCol : writeCols) { + // try the exact name first so a column whose own name contains a dot is + // not split + String topLevel = + icebergNames.contains(writeCol) + ? writeCol + : writeCol.substring(0, Math.max(writeCol.indexOf('.'), 0)); + if (icebergNames.contains(topLevel) && !statsColumns.contains(topLevel)) { + statsColumns.add(topLevel); + } + } + } Map indexMap = new HashMap<>(); if (statsColumns == null) { for (int i = 0; i < numFields; i++) { diff --git a/paimon-core/src/test/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMetaTest.java b/paimon-core/src/test/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMetaTest.java index b0b7154b18d2..af9799af8281 100644 --- a/paimon-core/src/test/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMetaTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/iceberg/manifest/IcebergDataFileMetaTest.java @@ -89,6 +89,7 @@ void testUnknownNullCountOmitted() { 100, icebergSchema, new SimpleStats(values, values, nullCounts), + null, null); assertThat(meta.nullValueCounts().size()).isEqualTo(1); @@ -128,6 +129,7 @@ void testRequiredFieldWithUnknownStats() { 100, icebergSchema, new SimpleStats(values, values, nullCounts), + null, null); assertThat(meta.nullValueCounts().size()).isEqualTo(1); @@ -167,7 +169,8 @@ void testStatsColumnsSubsetAlignment() { 100, icebergSchema, new SimpleStats(values, values, nullCounts), - Arrays.asList("b")); + Arrays.asList("b"), + null); assertThat(meta.nullValueCounts().size()).isEqualTo(1); assertThat(((GenericMap) meta.nullValueCounts()).get(2)).isEqualTo(2L); @@ -219,6 +222,7 @@ void testRequiredNestedFieldSkipped() { 100, icebergSchema, new SimpleStats(minValues, maxValues, nullCounts), + null, null); assertThat(meta.lowerBounds().size()).isEqualTo(1); @@ -267,6 +271,7 @@ void testGeospatialBoundsSkipped() { 100, icebergSchema, new SimpleStats(values, values, nullCounts), + null, null); assertThat(meta.nullValueCounts().size()).isEqualTo(2); @@ -275,4 +280,100 @@ void testGeospatialBoundsSkipped() { assertThat(meta.lowerBounds().size()).isZero(); assertThat(meta.upperBounds().size()).isZero(); } + + @Test + @DisplayName("Null stats columns of a partial write align by the write columns") + void testNullStatsColumnsWithWriteColsAlignByWriteSchema() { + IcebergSchema icebergSchema = + new IcebergSchema( + 0, + Arrays.asList( + new IcebergDataField(1, "k", false, "int", null), + new IcebergDataField(2, "a", false, "int", null), + new IcebergDataField(3, "b", false, "int", null))); + + // partial write of (b, k) in SET-clause order: the stats row follows the + // write-column order, so slot 0 is b and slot 1 is k, and the bounds must not + // drift onto a + BinaryRow values = new BinaryRow(2); + BinaryRowWriter rowWriter = new BinaryRowWriter(values); + rowWriter.writeInt(0, 100); + rowWriter.writeInt(1, 1); + rowWriter.complete(); + + BinaryArray nullCounts = new BinaryArray(); + BinaryArrayWriter arrayWriter = new BinaryArrayWriter(nullCounts, 2, 8); + arrayWriter.writeLong(0, 0L); + arrayWriter.writeLong(1, 0L); + arrayWriter.complete(); + + IcebergDataFileMeta meta = + IcebergDataFileMeta.create( + IcebergDataFileMeta.Content.DATA, + "path", + "parquet", + BinaryRow.EMPTY_ROW, + 10, + 100, + icebergSchema, + new SimpleStats(values, values, nullCounts), + null, + Arrays.asList("b", "k")); + + assertThat(meta.lowerBounds().size()).isEqualTo(2); + byte[] int1 = {1, 0, 0, 0}; + byte[] int100 = {100, 0, 0, 0}; + assertThat((byte[]) ((GenericMap) meta.lowerBounds()).get(3)).isEqualTo(int100); + assertThat((byte[]) ((GenericMap) meta.lowerBounds()).get(1)).isEqualTo(int1); + assertThat((byte[]) ((GenericMap) meta.upperBounds()).get(1)).isEqualTo(int1); + assertThat(((GenericMap) meta.lowerBounds()).get(2)).isNull(); + assertThat(((GenericMap) meta.nullValueCounts()).get(2)).isNull(); + } + + @Test + @DisplayName("Null stats columns of a nested partial write map leaf paths to the field") + void testNullStatsColumnsWithNestedWriteColsMapToTopLevel() { + IcebergSchema icebergSchema = + new IcebergSchema( + 0, + Arrays.asList( + new IcebergDataField(1, "k", false, "int", null), + // primitive here: the leaf-path mapping under test does + // not depend on the field's type + new IcebergDataField(2, "nest", false, "int", null))); + + // partial write of (nest.x, k) in SET-clause order: the write schema holds the + // top-level nest field and k, and the stats follow that order + BinaryRow values = new BinaryRow(2); + BinaryRowWriter rowWriter = new BinaryRowWriter(values); + rowWriter.setNullAt(0); + rowWriter.writeInt(1, 7); + rowWriter.complete(); + + BinaryArray nullCounts = new BinaryArray(); + BinaryArrayWriter arrayWriter = new BinaryArrayWriter(nullCounts, 2, 8); + arrayWriter.writeLong(0, 3L); + arrayWriter.writeLong(1, 0L); + arrayWriter.complete(); + + IcebergDataFileMeta meta = + IcebergDataFileMeta.create( + IcebergDataFileMeta.Content.DATA, + "path", + "parquet", + BinaryRow.EMPTY_ROW, + 10, + 100, + icebergSchema, + new SimpleStats(values, values, nullCounts), + null, + Arrays.asList("nest.x", "k")); + + // the leaf path maps to the top-level field: nest's stats land on field 2, not k + assertThat(meta.lowerBounds().size()).isEqualTo(1); + assertThat((byte[]) ((GenericMap) meta.lowerBounds()).get(1)) + .isEqualTo(new byte[] {7, 0, 0, 0}); + assertThat(((GenericMap) meta.nullValueCounts()).get(2)).isEqualTo(3L); + assertThat(((GenericMap) meta.nullValueCounts()).get(1)).isEqualTo(0L); + } }