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); + } }