Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -1370,7 +1371,8 @@ private List<IcebergManifestFileMeta> createNewlyAddedManifestFileMetas(
paimonFileMeta.fileSize(),
schemaCache.get(paimonFileMeta.schemaId()),
paimonFileMeta.valueStats(),
paimonFileMeta.valueStatsCols());
paimonFileMeta.valueStatsCols(),
paimonFileMeta.writeCols());
return new IcebergManifestEntry(
IcebergManifestEntry.Status.ADDED,
currentSnapshotId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -154,8 +156,31 @@ public static IcebergDataFileMeta create(
long fileSizeInBytes,
IcebergSchema icebergSchema,
SimpleStats stats,
@Nullable List<String> statsColumns) {
@Nullable List<String> statsColumns,
@Nullable List<String> 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<String> 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<String, Integer> indexMap = new HashMap<>();
if (statsColumns == null) {
for (int i = 0; i < numFields; i++) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,7 @@ void testUnknownNullCountOmitted() {
100,
icebergSchema,
new SimpleStats(values, values, nullCounts),
null,
null);

assertThat(meta.nullValueCounts().size()).isEqualTo(1);
Expand Down Expand Up @@ -128,6 +129,7 @@ void testRequiredFieldWithUnknownStats() {
100,
icebergSchema,
new SimpleStats(values, values, nullCounts),
null,
null);

assertThat(meta.nullValueCounts().size()).isEqualTo(1);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -219,6 +222,7 @@ void testRequiredNestedFieldSkipped() {
100,
icebergSchema,
new SimpleStats(minValues, maxValues, nullCounts),
null,
null);

assertThat(meta.lowerBounds().size()).isEqualTo(1);
Expand Down Expand Up @@ -267,6 +271,7 @@ void testGeospatialBoundsSkipped() {
100,
icebergSchema,
new SimpleStats(values, values, nullCounts),
null,
null);

assertThat(meta.nullValueCounts().size()).isEqualTo(2);
Expand All @@ -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);
}
}
Loading