From 4bdbfafdf3e26929803bb88428e4a87c604e7845 Mon Sep 17 00:00:00 2001 From: Xin Huang Date: Fri, 31 Jul 2026 07:47:59 -0700 Subject: [PATCH] Core, Parquet: Carry average value sizes for v4 Propagate average non-null value sizes from Parquet metrics through Metrics and ContentFile so v4 content stats adapters can preserve them. Retain compatibility with older serialized Metrics instances. Generated-by: Codex --- .../java/org/apache/iceberg/ContentFile.java | 16 +++-- .../main/java/org/apache/iceberg/Metrics.java | 58 ++++++++++++++++++- .../iceberg/TestMetricsSerialization.java | 28 ++++++++- .../java/org/apache/iceberg/BaseFile.java | 11 ++++ .../apache/iceberg/ContentStatsBackedMap.java | 8 +++ .../java/org/apache/iceberg/DataFiles.java | 7 ++- .../java/org/apache/iceberg/Delegates.java | 5 ++ .../java/org/apache/iceberg/FieldStats.java | 2 +- .../java/org/apache/iceberg/FileMetadata.java | 7 ++- .../org/apache/iceberg/GenericDataFile.java | 1 + .../org/apache/iceberg/GenericDeleteFile.java | 1 + .../java/org/apache/iceberg/MetricsUtil.java | 6 +- .../apache/iceberg/TrackedFileAdapters.java | 5 ++ .../apache/iceberg/util/ContentFileUtil.java | 8 ++- .../org/apache/iceberg/StatsTestUtil.java | 15 +++++ .../iceberg/TestContentStatsBackedMap.java | 13 +++++ .../java/org/apache/iceberg/TestMetrics.java | 8 ++- .../iceberg/TestTrackedFileAdapters.java | 36 ++++++++---- .../iceberg/parquet/ParquetMetrics.java | 8 ++- .../parquet/TestParquetDataWriter.java | 16 ++++- 20 files changed, 229 insertions(+), 30 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/ContentFile.java b/api/src/main/java/org/apache/iceberg/ContentFile.java index c4cd3c06adfd..cdf2a291d115 100644 --- a/api/src/main/java/org/apache/iceberg/ContentFile.java +++ b/api/src/main/java/org/apache/iceberg/ContentFile.java @@ -93,6 +93,14 @@ default String location() { /** Returns if collected, map from column ID to its NaN value count, null otherwise. */ Map nanValueCounts(); + /** + * Returns if collected, map from column ID to its average non-null value size in bytes, null + * otherwise. + */ + default Map avgValueSizes() { + return null; + } + /** Returns if collected, map from column ID to value lower bounds, null otherwise. */ Map lowerBounds(); @@ -187,7 +195,7 @@ default Long firstRowId() { * to copy data without stats when collecting files. * * @return a copy of this data file, without lower bounds, upper bounds, value counts, null value - * counts, or nan value counts + * counts, nan value counts, or average value sizes */ F copyWithoutStats(); @@ -198,7 +206,7 @@ default Long firstRowId() { * * @param requestedColumnIds column IDs for which to keep stats. * @return a copy of data file, with lower bounds, upper bounds, value counts, null value counts, - * and nan value counts for only specific columns. + * nan value counts, and average value sizes for only specific columns. */ default F copyWithStats(Set requestedColumnIds) { throw new UnsupportedOperationException( @@ -211,8 +219,8 @@ default F copyWithStats(Set requestedColumnIds) { * * @param withStats Will copy this file without file stats if set to false. * @return a copy of this data file. If withStats is set to false the - * file will not contain lower bounds, upper bounds, value counts, null value counts, or nan - * value counts + * file will not contain lower bounds, upper bounds, value counts, null value counts, nan + * value counts, or average value sizes */ default F copy(boolean withStats) { return withStats ? copy() : copyWithoutStats(); diff --git a/api/src/main/java/org/apache/iceberg/Metrics.java b/api/src/main/java/org/apache/iceberg/Metrics.java index eb13caf5e80a..633db32c173a 100644 --- a/api/src/main/java/org/apache/iceberg/Metrics.java +++ b/api/src/main/java/org/apache/iceberg/Metrics.java @@ -21,6 +21,7 @@ import java.io.IOException; import java.io.ObjectInputStream; import java.io.ObjectOutputStream; +import java.io.OptionalDataException; import java.io.Serializable; import java.nio.ByteBuffer; import java.util.Map; @@ -31,11 +32,14 @@ /** Iceberg file format metrics. */ public class Metrics implements Serializable { + private static final long serialVersionUID = 5337054944254511702L; + private Long rowCount = null; private Map columnSizes = null; private Map valueCounts = null; private Map nullValueCounts = null; private Map nanValueCounts = null; + private Map avgValueSizes = null; private Map lowerBounds = null; private Map upperBounds = null; // this is not serialized with all the other fields @@ -53,7 +57,16 @@ public Metrics( Map valueCounts, Map nullValueCounts, Map nanValueCounts) { - this(rowCount, columnSizes, valueCounts, nullValueCounts, nanValueCounts, null, null, null); + this( + rowCount, + columnSizes, + valueCounts, + nullValueCounts, + nanValueCounts, + null, + null, + null, + null); } public Metrics( @@ -72,6 +85,7 @@ public Metrics( nanValueCounts, lowerBounds, upperBounds, + null, null); } @@ -84,11 +98,34 @@ public Metrics( Map lowerBounds, Map upperBounds, Map originalTypes) { + this( + rowCount, + columnSizes, + valueCounts, + nullValueCounts, + nanValueCounts, + lowerBounds, + upperBounds, + originalTypes, + null); + } + + public Metrics( + Long rowCount, + Map columnSizes, + Map valueCounts, + Map nullValueCounts, + Map nanValueCounts, + Map lowerBounds, + Map upperBounds, + Map originalTypes, + Map avgValueSizes) { this.rowCount = rowCount; this.columnSizes = columnSizes; this.valueCounts = valueCounts; this.nullValueCounts = nullValueCounts; this.nanValueCounts = nanValueCounts; + this.avgValueSizes = avgValueSizes; this.lowerBounds = lowerBounds; this.upperBounds = upperBounds; this.originalTypes = originalTypes; @@ -139,6 +176,15 @@ public Map nanValueCounts() { return nanValueCounts; } + /** + * Get the average non-null value size in bytes for all fields where it was collected. + * + * @return a Map of fieldId to average value size in bytes + */ + public Map avgValueSizes() { + return avgValueSizes; + } + /** * Get the non-null lower bound values for all fields in a file. * @@ -186,6 +232,7 @@ private void writeObject(ObjectOutputStream out) throws IOException { writeByteBufferMap(out, lowerBounds); writeByteBufferMap(out, upperBounds); + out.writeObject(avgValueSizes); } private static void writeByteBufferMap( @@ -222,6 +269,15 @@ private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundE lowerBounds = readByteBufferMap(in); upperBounds = readByteBufferMap(in); + try { + avgValueSizes = (Map) in.readObject(); + } catch (OptionalDataException e) { + if (!e.eof) { + throw e; + } + + avgValueSizes = null; + } } @SuppressWarnings("DangerousJavaDeserialization") diff --git a/api/src/test/java/org/apache/iceberg/TestMetricsSerialization.java b/api/src/test/java/org/apache/iceberg/TestMetricsSerialization.java index a97aa5cfe8c8..5cc2a90c976e 100644 --- a/api/src/test/java/org/apache/iceberg/TestMetricsSerialization.java +++ b/api/src/test/java/org/apache/iceberg/TestMetricsSerialization.java @@ -26,6 +26,7 @@ import java.io.ObjectInputStream; import java.io.ObjectOutputStream; import java.nio.ByteBuffer; +import java.util.Base64; import java.util.Map; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; import org.apache.iceberg.relocated.com.google.common.collect.Maps; @@ -35,6 +36,14 @@ public class TestMetricsSerialization { + private static final String OLD_SERIALIZED_METRICS = + "rO0ABXNyABpvcmcuYXBhY2hlLmljZWJlcmcuTWV0cmljc0oRBzHjCEpWAwAITAALY29sdW1uU2l6" + + "ZXN0AA9MamF2YS91dGlsL01hcDtMAAtsb3dlckJvdW5kc3EAfgABTAAObmFuVmFsdWVDb3VudHNx" + + "AH4AAUwAD251bGxWYWx1ZUNvdW50c3EAfgABTAANb3JpZ2luYWxUeXBlc3EAfgABTAAIcm93Q291" + + "bnR0ABBMamF2YS9sYW5nL0xvbmc7TAALdXBwZXJCb3VuZHNxAH4AAUwAC3ZhbHVlQ291bnRzcQB+" + + "AAF4cHNyAA5qYXZhLmxhbmcuTG9uZzuL5JDMjyPfAgABSgAFdmFsdWV4cgAQamF2YS5sYW5nLk51" + + "bWJlcoaslR0LlOCLAgAAeHAAAAAAAAAAB3BwcHB3CP//////////eA=="; + @Test public void testSerialization() throws IOException, ClassNotFoundException { Metrics original = generateMetrics(); @@ -55,6 +64,14 @@ public void testSerializationWithNulls() throws IOException, ClassNotFoundExcept assertEquals(original, result); } + @Test + void oldSerializationCompatibility() throws IOException, ClassNotFoundException { + Metrics metrics = deserialize(Base64.getDecoder().decode(OLD_SERIALIZED_METRICS)); + + assertThat(metrics.recordCount()).isEqualTo(7L); + assertThat(metrics.avgValueSizes()).isNull(); + } + private static byte[] serialize(Metrics metrics) throws IOException { try (ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream()) { ObjectOutputStream objectOutputStream = new ObjectOutputStream(byteArrayOutputStream); @@ -93,7 +110,9 @@ private static Metrics generateMetrics() { Map originalTypes = ImmutableMap.of(1, Types.IntegerType.get(), 2, Types.IntegerType.get()); - return new Metrics(0L, longMap1, longMap2, longMap3, null, byteMap1, byteMap2, originalTypes); + Map avgValueSizes = ImmutableMap.of(9, 10); + return new Metrics( + 0L, longMap1, longMap2, longMap3, null, byteMap1, byteMap2, originalTypes, avgValueSizes); } private static Metrics generateMetricsWithNulls() { @@ -106,7 +125,11 @@ private static Metrics generateMetricsWithNulls() { byteMap.put(4, null); Map originalTypes = ImmutableMap.of(4, Types.IntegerType.get()); - return new Metrics(null, null, longMap, longMap, null, null, byteMap, originalTypes); + Map avgValueSizes = Maps.newHashMap(); + avgValueSizes.put(null, 1); + avgValueSizes.put(2, null); + return new Metrics( + null, null, longMap, longMap, null, null, byteMap, originalTypes, avgValueSizes); } private static void assertEquals(Metrics expected, Metrics actual) { @@ -114,6 +137,7 @@ private static void assertEquals(Metrics expected, Metrics actual) { assertThat(actual.columnSizes()).isEqualTo(expected.columnSizes()); assertThat(actual.valueCounts()).isEqualTo(expected.valueCounts()); assertThat(actual.nullValueCounts()).isEqualTo(expected.nullValueCounts()); + assertThat(actual.avgValueSizes()).isEqualTo(expected.avgValueSizes()); assertEquals(expected.lowerBounds(), actual.lowerBounds()); assertEquals(expected.upperBounds(), actual.upperBounds()); diff --git a/core/src/main/java/org/apache/iceberg/BaseFile.java b/core/src/main/java/org/apache/iceberg/BaseFile.java index 9c6a4d5ef607..67554d36b5ec 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFile.java +++ b/core/src/main/java/org/apache/iceberg/BaseFile.java @@ -68,6 +68,7 @@ abstract class BaseFile extends SupportsIndexProjection private Map valueCounts = null; private Map nullValueCounts = null; private Map nanValueCounts = null; + private Map avgValueSizes = null; private Map lowerBounds = null; private Map upperBounds = null; private long[] splitOffsets = null; @@ -146,6 +147,7 @@ abstract class BaseFile extends SupportsIndexProjection Map valueCounts, Map nullValueCounts, Map nanValueCounts, + Map avgValueSizes, Map lowerBounds, Map upperBounds, List splitOffsets, @@ -178,6 +180,7 @@ abstract class BaseFile extends SupportsIndexProjection this.valueCounts = valueCounts; this.nullValueCounts = nullValueCounts; this.nanValueCounts = nanValueCounts; + this.avgValueSizes = avgValueSizes; this.lowerBounds = SerializableByteBufferMap.wrap(lowerBounds); this.upperBounds = SerializableByteBufferMap.wrap(upperBounds); this.splitOffsets = ArrayUtil.toLongArray(splitOffsets); @@ -215,6 +218,7 @@ abstract class BaseFile extends SupportsIndexProjection this.valueCounts = copyMap(toCopy.valueCounts, requestedColumnIds); this.nullValueCounts = copyMap(toCopy.nullValueCounts, requestedColumnIds); this.nanValueCounts = copyMap(toCopy.nanValueCounts, requestedColumnIds); + this.avgValueSizes = copyMap(toCopy.avgValueSizes, requestedColumnIds); this.lowerBounds = copyByteBufferMap(toCopy.lowerBounds, requestedColumnIds); this.upperBounds = copyByteBufferMap(toCopy.upperBounds, requestedColumnIds); } else { @@ -222,6 +226,7 @@ abstract class BaseFile extends SupportsIndexProjection this.valueCounts = null; this.nullValueCounts = null; this.nanValueCounts = null; + this.avgValueSizes = null; this.lowerBounds = null; this.upperBounds = null; } @@ -511,6 +516,11 @@ public Map nanValueCounts() { return toReadableMap(nanValueCounts); } + @Override + public Map avgValueSizes() { + return toReadableMap(avgValueSizes); + } + @Override public Map lowerBounds() { return toReadableByteBufferMap(lowerBounds); @@ -648,6 +658,7 @@ public String toString() { .add("value_counts", valueCounts) .add("null_value_counts", nullValueCounts) .add("nan_value_counts", nanValueCounts) + .add("avg_value_sizes", avgValueSizes) .add("lower_bounds", lowerBounds) .add("upper_bounds", upperBounds) .add("key_metadata", keyMetadata == null ? "null" : "(redacted)") diff --git a/core/src/main/java/org/apache/iceberg/ContentStatsBackedMap.java b/core/src/main/java/org/apache/iceberg/ContentStatsBackedMap.java index 3f2599fc2dac..fa33880010fb 100644 --- a/core/src/main/java/org/apache/iceberg/ContentStatsBackedMap.java +++ b/core/src/main/java/org/apache/iceberg/ContentStatsBackedMap.java @@ -35,6 +35,7 @@ private enum Kind { VALUE_COUNT, NULL_VALUE_COUNT, NAN_VALUE_COUNT, + AVG_VALUE_SIZE, LOWER_BOUND, UPPER_BOUND } @@ -54,6 +55,11 @@ static Map nanValueCounts(ContentStats stats) { return viewOrNull(stats, Kind.NAN_VALUE_COUNT); } + /** Per-column average non-null value sizes, or null if no column tracks the average size. */ + static Map avgValueSizes(ContentStats stats) { + return viewOrNull(stats, Kind.AVG_VALUE_SIZE); + } + /** Per-column lower bounds, or null if no column tracks a lower bound. */ static Map lowerBounds(ContentStats stats) { return viewOrNull(stats, Kind.LOWER_BOUND); @@ -142,6 +148,7 @@ private static boolean isKnown(FieldStats fieldStats, Kind kind) { case VALUE_COUNT -> fieldStats.hasValueCount(); case NULL_VALUE_COUNT -> fieldStats.hasNullValueCount(); case NAN_VALUE_COUNT -> fieldStats.hasNanValueCount(); + case AVG_VALUE_SIZE -> fieldStats.avgValueSizeInBytes() != null; case LOWER_BOUND -> fieldStats.lowerBound() != null; case UPPER_BOUND -> fieldStats.upperBound() != null; }; @@ -156,6 +163,7 @@ private static V statValue(FieldStats fieldStats, Kind kind) { fieldStats.hasNullValueCount() ? (V) Long.valueOf(fieldStats.nullValueCount()) : null; case NAN_VALUE_COUNT -> fieldStats.hasNanValueCount() ? (V) Long.valueOf(fieldStats.nanValueCount()) : null; + case AVG_VALUE_SIZE -> (V) fieldStats.avgValueSizeInBytes(); case LOWER_BOUND -> (V) bound(fieldStats, fieldStats.lowerBound(), StatsUtil.LOWER_BOUND_NAME); case UPPER_BOUND -> diff --git a/core/src/main/java/org/apache/iceberg/DataFiles.java b/core/src/main/java/org/apache/iceberg/DataFiles.java index 4a991f623051..d62a82d1f99c 100644 --- a/core/src/main/java/org/apache/iceberg/DataFiles.java +++ b/core/src/main/java/org/apache/iceberg/DataFiles.java @@ -150,6 +150,7 @@ public static class Builder { private Map valueCounts = null; private Map nullValueCounts = null; private Map nanValueCounts = null; + private Map avgValueSizes = null; private Map lowerBounds = null; private Map upperBounds = null; private Map originalTypes = null; @@ -177,6 +178,7 @@ public void clear() { this.valueCounts = null; this.nullValueCounts = null; this.nanValueCounts = null; + this.avgValueSizes = null; this.lowerBounds = null; this.upperBounds = null; this.splitOffsets = null; @@ -198,6 +200,7 @@ public Builder copy(DataFile toCopy) { this.valueCounts = toCopy.valueCounts(); this.nullValueCounts = toCopy.nullValueCounts(); this.nanValueCounts = toCopy.nanValueCounts(); + this.avgValueSizes = toCopy.avgValueSizes(); this.lowerBounds = toCopy.lowerBounds(); this.upperBounds = toCopy.upperBounds(); this.keyMetadata = @@ -290,6 +293,7 @@ public Builder withMetrics(Metrics metrics) { this.valueCounts = metrics.valueCounts(); this.nullValueCounts = metrics.nullValueCounts(); this.nanValueCounts = metrics.nanValueCounts(); + this.avgValueSizes = metrics.avgValueSizes(); this.lowerBounds = metrics.lowerBounds(); this.upperBounds = metrics.upperBounds(); this.originalTypes = metrics.originalTypes(); @@ -354,7 +358,8 @@ public DataFile build() { nanValueCounts, lowerBounds, upperBounds, - originalTypes), + originalTypes, + avgValueSizes), keyMetadata, splitOffsets, sortOrderId, diff --git a/core/src/main/java/org/apache/iceberg/Delegates.java b/core/src/main/java/org/apache/iceberg/Delegates.java index 324fc1cccdc9..d89e9ded7f38 100644 --- a/core/src/main/java/org/apache/iceberg/Delegates.java +++ b/core/src/main/java/org/apache/iceberg/Delegates.java @@ -232,6 +232,11 @@ public Map nanValueCounts() { return wrapped.nanValueCounts(); } + @Override + public Map avgValueSizes() { + return wrapped.avgValueSizes(); + } + @Override public Map lowerBounds() { return wrapped.lowerBounds(); diff --git a/core/src/main/java/org/apache/iceberg/FieldStats.java b/core/src/main/java/org/apache/iceberg/FieldStats.java index 495dddba913d..e679b6bd47e4 100644 --- a/core/src/main/java/org/apache/iceberg/FieldStats.java +++ b/core/src/main/java/org/apache/iceberg/FieldStats.java @@ -63,7 +63,7 @@ interface FieldStats { /** * The avg value size in memory (uncompressed) in bytes for variable-length types (string, binary, - * variant) + * variant, geometry, geography) */ Integer avgValueSizeInBytes(); diff --git a/core/src/main/java/org/apache/iceberg/FileMetadata.java b/core/src/main/java/org/apache/iceberg/FileMetadata.java index a5266101c252..026b413adb18 100644 --- a/core/src/main/java/org/apache/iceberg/FileMetadata.java +++ b/core/src/main/java/org/apache/iceberg/FileMetadata.java @@ -55,6 +55,7 @@ public static class Builder { private Map valueCounts = null; private Map nullValueCounts = null; private Map nanValueCounts = null; + private Map avgValueSizes = null; private Map lowerBounds = null; private Map upperBounds = null; private Map originalTypes = null; @@ -84,6 +85,7 @@ public void clear() { this.valueCounts = null; this.nullValueCounts = null; this.nanValueCounts = null; + this.avgValueSizes = null; this.lowerBounds = null; this.upperBounds = null; this.sortOrderId = null; @@ -104,6 +106,7 @@ public Builder copy(DeleteFile toCopy) { this.valueCounts = toCopy.valueCounts(); this.nullValueCounts = toCopy.nullValueCounts(); this.nanValueCounts = toCopy.nanValueCounts(); + this.avgValueSizes = toCopy.avgValueSizes(); this.lowerBounds = toCopy.lowerBounds(); this.upperBounds = toCopy.upperBounds(); this.keyMetadata = @@ -200,6 +203,7 @@ public Builder withMetrics(Metrics metrics) { this.valueCounts = metrics.valueCounts(); this.nullValueCounts = metrics.nullValueCounts(); this.nanValueCounts = metrics.nanValueCounts(); + this.avgValueSizes = metrics.avgValueSizes(); this.lowerBounds = metrics.lowerBounds(); this.upperBounds = metrics.upperBounds(); this.originalTypes = metrics.originalTypes(); @@ -300,7 +304,8 @@ public DeleteFile build() { nanValueCounts, lowerBounds, upperBounds, - originalTypes), + originalTypes, + avgValueSizes), equalityFieldIds, sortOrderId, splitOffsets, diff --git a/core/src/main/java/org/apache/iceberg/GenericDataFile.java b/core/src/main/java/org/apache/iceberg/GenericDataFile.java index 9033d0718bb1..4e96a268951e 100644 --- a/core/src/main/java/org/apache/iceberg/GenericDataFile.java +++ b/core/src/main/java/org/apache/iceberg/GenericDataFile.java @@ -60,6 +60,7 @@ class GenericDataFile extends BaseFile implements DataFile { metrics.valueCounts(), metrics.nullValueCounts(), metrics.nanValueCounts(), + metrics.avgValueSizes(), metrics.lowerBounds(), metrics.upperBounds(), splitOffsets, diff --git a/core/src/main/java/org/apache/iceberg/GenericDeleteFile.java b/core/src/main/java/org/apache/iceberg/GenericDeleteFile.java index 897d77c2fb5a..86cd423c785c 100644 --- a/core/src/main/java/org/apache/iceberg/GenericDeleteFile.java +++ b/core/src/main/java/org/apache/iceberg/GenericDeleteFile.java @@ -64,6 +64,7 @@ class GenericDeleteFile extends BaseFile implements DeleteFile { metrics.valueCounts(), metrics.nullValueCounts(), metrics.nanValueCounts(), + metrics.avgValueSizes(), metrics.lowerBounds(), metrics.upperBounds(), splitOffsets, diff --git a/core/src/main/java/org/apache/iceberg/MetricsUtil.java b/core/src/main/java/org/apache/iceberg/MetricsUtil.java index 96645bf074ae..d6df03222eb3 100644 --- a/core/src/main/java/org/apache/iceberg/MetricsUtil.java +++ b/core/src/main/java/org/apache/iceberg/MetricsUtil.java @@ -56,7 +56,8 @@ public static Metrics copyWithoutFieldCounts(Metrics metrics, Set exclu copyWithoutKeys(metrics.nanValueCounts(), excludedFieldIds), metrics.lowerBounds(), metrics.upperBounds(), - metrics.originalTypes()); + metrics.originalTypes(), + copyWithoutKeys(metrics.avgValueSizes(), excludedFieldIds)); } /** @@ -75,7 +76,8 @@ public static Metrics copyWithoutFieldCountsAndBounds( copyWithoutKeys(metrics.nanValueCounts(), excludedFieldIds), copyWithoutKeys(metrics.lowerBounds(), excludedFieldIds), copyWithoutKeys(metrics.upperBounds(), excludedFieldIds), - copyWithoutKeys(metrics.originalTypes(), excludedFieldIds)); + copyWithoutKeys(metrics.originalTypes(), excludedFieldIds), + copyWithoutKeys(metrics.avgValueSizes(), excludedFieldIds)); } private static Map copyWithoutKeys(Map map, Set keys) { diff --git a/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java b/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java index 355e9602af35..1f5d1260ff5e 100644 --- a/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java +++ b/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java @@ -187,6 +187,11 @@ public Map nanValueCounts() { return ContentStatsBackedMap.nanValueCounts(file().contentStats()); } + @Override + public Map avgValueSizes() { + return ContentStatsBackedMap.avgValueSizes(file().contentStats()); + } + @Override public Map lowerBounds() { return ContentStatsBackedMap.lowerBounds(file().contentStats()); diff --git a/core/src/main/java/org/apache/iceberg/util/ContentFileUtil.java b/core/src/main/java/org/apache/iceberg/util/ContentFileUtil.java index ac7de6644c51..6f8e5d6bbf86 100644 --- a/core/src/main/java/org/apache/iceberg/util/ContentFileUtil.java +++ b/core/src/main/java/org/apache/iceberg/util/ContentFileUtil.java @@ -175,7 +175,9 @@ private static Metrics metricsWithoutPathBounds(DeleteFile file) { file.nullValueCounts(), file.nanValueCounts(), lowerBounds == null ? null : Collections.unmodifiableMap(lowerBounds), - upperBounds == null ? null : Collections.unmodifiableMap(upperBounds)); + upperBounds == null ? null : Collections.unmodifiableMap(upperBounds), + null, + file.avgValueSizes()); } private static Metrics metricsWithPathBounds(DeleteFile file, ByteBuffer bound) { @@ -197,6 +199,8 @@ private static Metrics metricsWithPathBounds(DeleteFile file, ByteBuffer bound) file.nullValueCounts(), file.nanValueCounts(), lowerBounds == null ? null : Collections.unmodifiableMap(lowerBounds), - upperBounds == null ? null : Collections.unmodifiableMap(upperBounds)); + upperBounds == null ? null : Collections.unmodifiableMap(upperBounds), + null, + file.avgValueSizes()); } } diff --git a/core/src/test/java/org/apache/iceberg/StatsTestUtil.java b/core/src/test/java/org/apache/iceberg/StatsTestUtil.java index 1a6ee1923292..3c05153d937d 100644 --- a/core/src/test/java/org/apache/iceberg/StatsTestUtil.java +++ b/core/src/test/java/org/apache/iceberg/StatsTestUtil.java @@ -38,6 +38,19 @@ static FieldStats mockFieldStats( Long valueCount, Long nullCount, Long nanCount) { + return mockFieldStats(type, id, lower, upper, valueCount, nullCount, nanCount, null); + } + + @SuppressWarnings("unchecked") + static FieldStats mockFieldStats( + Types.StructType type, + int id, + Object lower, + Object upper, + Long valueCount, + Long nullCount, + Long nanCount, + Integer avgValueSize) { FieldStats stats = Mockito.mock(FieldStats.class); Mockito.when(stats.fieldId()).thenReturn(id); Mockito.when(stats.type()).thenReturn(type); @@ -58,6 +71,8 @@ static FieldStats mockFieldStats( Mockito.when(stats.nanValueCount()).thenReturn(nanCount); } + Mockito.when(stats.avgValueSizeInBytes()).thenReturn(avgValueSize); + return stats; } } diff --git a/core/src/test/java/org/apache/iceberg/TestContentStatsBackedMap.java b/core/src/test/java/org/apache/iceberg/TestContentStatsBackedMap.java index 865ce0582792..31e6eb55d105 100644 --- a/core/src/test/java/org/apache/iceberg/TestContentStatsBackedMap.java +++ b/core/src/test/java/org/apache/iceberg/TestContentStatsBackedMap.java @@ -79,6 +79,18 @@ public void testNanValueCountsOnlyForFloatingColumns() { assertThat(map).containsOnly(Map.entry(3, 4L)); } + @Test + public void avgValueSizes() { + Schema schema = new Schema(optional(4, "str", Types.StringType.get())); + Types.StructType statsType = StatsUtil.statsReadSchema(schema, List.of(4)); + Types.StructType fieldStatsType = statsType.field("str").type().asStructType(); + ContentStatsStruct stats = new ContentStatsStruct(statsType); + stats.setStats(4, StatsTestUtil.mockFieldStats(fieldStatsType, 4, "a", "z", 10L, 1L, null, 12)); + + Map map = ContentStatsBackedMap.avgValueSizes(stats); + assertThat(map).containsOnly(Map.entry(4, 12)); + } + @Test public void testLowerBounds() { Map lower = ContentStatsBackedMap.lowerBounds(POPULATED_STATS); @@ -118,6 +130,7 @@ public void testFactoryReturnsNullWhenNoColumnTracksMetric() { // only a required long column: it tracks neither null_value_count nor nan_value_count assertThat(ContentStatsBackedMap.nullValueCounts(ONLY_REQUIRED_STATS)).isNull(); assertThat(ContentStatsBackedMap.nanValueCounts(ONLY_REQUIRED_STATS)).isNull(); + assertThat(ContentStatsBackedMap.avgValueSizes(ONLY_REQUIRED_STATS)).isNull(); Map valueCounts = ContentStatsBackedMap.valueCounts(ONLY_REQUIRED_STATS); assertThat(valueCounts).isNotNull().containsOnly(Map.entry(1, 10L)); diff --git a/core/src/test/java/org/apache/iceberg/TestMetrics.java b/core/src/test/java/org/apache/iceberg/TestMetrics.java index 095feb29382f..3c68b0e39d44 100644 --- a/core/src/test/java/org/apache/iceberg/TestMetrics.java +++ b/core/src/test/java/org/apache/iceberg/TestMetrics.java @@ -271,9 +271,11 @@ public void testMetricsForGeospatialTypes() throws IOException { optional(3, "geog", Types.GeographyType.crs84())); Record first = GenericRecord.create(schema); + ByteBuffer geom = wkbPoint(30, 10); + ByteBuffer geog = wkbPoint(-5, 40); first.setField("id", 1L); - first.setField("geom", wkbPoint(30, 10)); - first.setField("geog", wkbPoint(-5, 40)); + first.setField("geom", geom); + first.setField("geog", geog); Record second = GenericRecord.create(schema); second.setField("id", 2L); // both geo columns are left null @@ -287,6 +289,8 @@ public void testMetricsForGeospatialTypes() throws IOException { assertBounds(2, Types.GeometryType.crs84(), null, null, metrics); assertCounts(3, 2L, 1L, metrics); assertBounds(3, Types.GeographyType.crs84(), null, null, metrics); + assertThat(metrics.avgValueSizes()) + .containsOnly(Map.entry(2, geom.remaining()), Map.entry(3, geog.remaining())); } @TestTemplate diff --git a/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java b/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java index 199239050690..bb466540372d 100644 --- a/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java +++ b/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java @@ -64,21 +64,27 @@ class TestTrackedFileAdapters { private static final Schema TABLE_SCHEMA = new Schema( - optional(1, "id", Types.IntegerType.get()), optional(2, "score", Types.FloatType.get())); + optional(1, "id", Types.IntegerType.get()), + optional(2, "score", Types.FloatType.get()), + optional(3, "name", Types.StringType.get())); private static final Types.StructType CONTENT_STATS_TYPE = - StatsUtil.statsReadSchema(TABLE_SCHEMA, ImmutableList.of(1, 2)); + StatsUtil.statsReadSchema(TABLE_SCHEMA, ImmutableList.of(1, 2, 3)); private static final FieldStats ID_STATS = StatsTestUtil.mockFieldStats( CONTENT_STATS_TYPE.fieldType("id").asStructType(), 1, 1, 1000, 100L, 5L, null); private static final FieldStats SCORE_STATS = StatsTestUtil.mockFieldStats( CONTENT_STATS_TYPE.fieldType("score").asStructType(), 2, 1.0f, 100.0f, 100L, 10L, 3L); + private static final FieldStats NAME_STATS = + StatsTestUtil.mockFieldStats( + CONTENT_STATS_TYPE.fieldType("name").asStructType(), 3, "a", "z", 100L, 20L, null, 12); private static final ContentStatsStruct CONTENT_STATS = new ContentStatsStruct(CONTENT_STATS_TYPE); static { CONTENT_STATS.setStats(1, ID_STATS); CONTENT_STATS.setStats(2, SCORE_STATS); + CONTENT_STATS.setStats(3, NAME_STATS); } @Test @@ -134,17 +140,22 @@ void dataFileAdapterDelegation() { assertThat(dataFile.manifestLocation()).isEqualTo(MANIFEST_LOCATION); assertThat(dataFile.equalityFieldIds()).isNull(); assertThat(dataFile.columnSizes()).isNull(); - assertThat(dataFile.valueCounts()).containsOnly(Map.entry(1, 100L), Map.entry(2, 100L)); - assertThat(dataFile.nullValueCounts()).containsOnly(Map.entry(1, 5L), Map.entry(2, 10L)); + assertThat(dataFile.valueCounts()) + .containsOnly(Map.entry(1, 100L), Map.entry(2, 100L), Map.entry(3, 100L)); + assertThat(dataFile.nullValueCounts()) + .containsOnly(Map.entry(1, 5L), Map.entry(2, 10L), Map.entry(3, 20L)); assertThat(dataFile.nanValueCounts()).containsOnly(Map.entry(2, 3L)); + assertThat(dataFile.avgValueSizes()).containsOnly(Map.entry(3, 12)); assertThat(dataFile.lowerBounds()) .containsOnly( Map.entry(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1)), - Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 1.0f))); + Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 1.0f)), + Map.entry(3, Conversions.toByteBuffer(Types.StringType.get(), "a"))); assertThat(dataFile.upperBounds()) .containsOnly( Map.entry(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1000)), - Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 100.0f))); + Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 100.0f)), + Map.entry(3, Conversions.toByteBuffer(Types.StringType.get(), "z"))); } @ParameterizedTest @@ -211,17 +222,22 @@ void equalityDeleteFileAdapterDelegation() { assertThat(deleteFile.manifestLocation()).isEqualTo(MANIFEST_LOCATION); assertThat(deleteFile.equalityFieldIds()).containsExactly(1, 2, 3); assertThat(deleteFile.columnSizes()).isNull(); - assertThat(deleteFile.valueCounts()).containsOnly(Map.entry(1, 100L), Map.entry(2, 100L)); - assertThat(deleteFile.nullValueCounts()).containsOnly(Map.entry(1, 5L), Map.entry(2, 10L)); + assertThat(deleteFile.valueCounts()) + .containsOnly(Map.entry(1, 100L), Map.entry(2, 100L), Map.entry(3, 100L)); + assertThat(deleteFile.nullValueCounts()) + .containsOnly(Map.entry(1, 5L), Map.entry(2, 10L), Map.entry(3, 20L)); assertThat(deleteFile.nanValueCounts()).containsOnly(Map.entry(2, 3L)); + assertThat(deleteFile.avgValueSizes()).containsOnly(Map.entry(3, 12)); assertThat(deleteFile.lowerBounds()) .containsOnly( Map.entry(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1)), - Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 1.0f))); + Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 1.0f)), + Map.entry(3, Conversions.toByteBuffer(Types.StringType.get(), "a"))); assertThat(deleteFile.upperBounds()) .containsOnly( Map.entry(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1000)), - Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 100.0f))); + Map.entry(2, Conversions.toByteBuffer(Types.FloatType.get(), 100.0f)), + Map.entry(3, Conversions.toByteBuffer(Types.StringType.get(), "z"))); } @ParameterizedTest diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetrics.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetrics.java index 0412ebc6953f..d7899703e258 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetrics.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetrics.java @@ -129,6 +129,7 @@ static Metrics metrics( Map valueCounts = Maps.newHashMap(); Map nullValueCounts = Maps.newHashMap(); Map nanValueCounts = Maps.newHashMap(); + Map avgValueSizes = Maps.newHashMap(); Map lowerBounds = Maps.newHashMap(); Map upperBounds = Maps.newHashMap(); Map originalTypes = Maps.newHashMap(); @@ -151,6 +152,10 @@ static Metrics metrics( nanValueCounts.put(id, metrics.nanValueCount()); } + if (metrics.avgValueSizeInBytes() != null) { + avgValueSizes.put(id, metrics.avgValueSizeInBytes()); + } + if (metrics.lowerBound() != null) { ByteBuffer lowerBound = Conversions.toByteBuffer(metrics.originalType(), metrics.lowerBound()); @@ -172,7 +177,8 @@ static Metrics metrics( nanValueCounts, lowerBounds, upperBounds, - originalTypes); + originalTypes, + avgValueSizes); } private static class MetricsVisitor extends TypeWithSchemaVisitor>> { diff --git a/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java b/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java index 00891b507eef..6c4f85d81821 100644 --- a/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java +++ b/parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java @@ -25,8 +25,10 @@ import java.nio.ByteBuffer; import java.nio.file.Path; import java.util.List; +import java.util.Map; import java.util.Optional; import java.util.Random; +import java.util.Set; import java.util.function.Function; import java.util.function.UnaryOperator; import org.apache.iceberg.DataFile; @@ -112,10 +114,11 @@ public void testGeospatialRoundTrip() throws IOException { Types.NestedField.optional(3, "geog", Types.GeographyType.crs84())); GenericRecord record = GenericRecord.create(schema); + ByteBuffer geom = wkbPoint(30, 10); + ByteBuffer geog = wkbPoint(-5, 40); List geoRecords = ImmutableList.of( - record.copy( - ImmutableMap.of("id", 1L, "geom", wkbPoint(30, 10), "geog", wkbPoint(-5, 40))), + record.copy(ImmutableMap.of("id", 1L, "geom", geom, "geog", geog)), // geog is left null record.copy(ImmutableMap.of("id", 2L, "geom", wkbPoint(0, 0))), // both geo columns are left null @@ -135,7 +138,14 @@ public void testGeospatialRoundTrip() throws IOException { } } - assertThat(dataWriter.toDataFile().recordCount()).isEqualTo(geoRecords.size()); + DataFile dataFile = dataWriter.toDataFile(); + assertThat(dataFile.recordCount()).isEqualTo(geoRecords.size()); + assertThat(dataFile.avgValueSizes()) + .containsOnly(Map.entry(2, geom.remaining()), Map.entry(3, geog.remaining())); + assertThat(dataFile.copy().avgValueSizes()).isEqualTo(dataFile.avgValueSizes()); + assertThat(dataFile.copyWithStats(Set.of(2)).avgValueSizes()) + .containsOnly(Map.entry(2, geom.remaining())); + assertThat(dataFile.copyWithoutStats().avgValueSizes()).isNull(); List writtenRecords; try (CloseableIterable reader =