This is an automated email from the ASF dual-hosted git repository.
huaxingao pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new e0afb38714 Parquet: Honor column metrics truncate length for variant
shredded bounds (#17342)
e0afb38714 is described below
commit e0afb3871498d776b6aaf8fb10901b403680ef62
Author: Neelesh Salian <[email protected]>
AuthorDate: Mon Aug 3 21:24:31 2026 -0700
Parquet: Honor column metrics truncate length for variant shredded bounds
(#17342)
* Parquet: Honor column metrics truncate length for variant shredded bounds
* test cleanup
* Bump testShreddedStringBoundsFull input above default truncation
---
.../org/apache/iceberg/parquet/ParquetMetrics.java | 12 +-
.../apache/iceberg/parquet/ParquetVariantUtil.java | 22 ++-
.../apache/iceberg/parquet/TestVariantMetrics.java | 197 +++++++++++++++++++--
3 files changed, 205 insertions(+), 26 deletions(-)
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 0412ebc695..8695cd2156 100644
--- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetrics.java
+++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetrics.java
@@ -390,7 +390,8 @@ class ParquetMetrics {
List<ParquetVariantUtil.VariantMetrics> results =
Lists.newArrayList(
- ParquetVariantVisitor.visit(variant, new
MetricsVariantVisitor(currentPath())));
+ ParquetVariantVisitor.visit(
+ variant, new MetricsVariantVisitor(currentPath(),
truncateLength(mode))));
if (results.isEmpty()) {
return ImmutableList.of();
@@ -442,9 +443,11 @@ class ParquetMetrics {
extends
ParquetVariantVisitor<Iterable<ParquetVariantUtil.VariantMetrics>> {
private final Deque<String> fieldNames = Lists.newLinkedList();
private final String[] basePath;
+ private final int truncateLength;
- private MetricsVariantVisitor(String[] basePath) {
+ private MetricsVariantVisitor(String[] basePath, int truncateLength) {
this.basePath = basePath;
+ this.truncateLength = truncateLength;
}
@Override
@@ -624,10 +627,11 @@ class ParquetMetrics {
return null;
}
- if (lowerBound != null && upperBound != null) {
+ if (lowerBound != null && upperBound != null && truncateLength > 0) {
VariantValue lower = Variants.of(variantType, lowerBound);
VariantValue upper = Variants.of(variantType, upperBound);
- return new ParquetVariantUtil.VariantMetrics(valueCount, nullCount,
lower, upper);
+ return new ParquetVariantUtil.VariantMetrics(
+ valueCount, nullCount, lower, upper, truncateLength);
} else {
return new ParquetVariantUtil.VariantMetrics(valueCount, nullCount);
}
diff --git
a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetVariantUtil.java
b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetVariantUtil.java
index a9b088aaf2..144407bcee 100644
--- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetVariantUtil.java
+++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetVariantUtil.java
@@ -307,11 +307,15 @@ class ParquetVariantUtil {
}
VariantMetrics(
- long valueCount, long nullCount, VariantValue lowerBound, VariantValue
upperBound) {
+ long valueCount,
+ long nullCount,
+ VariantValue lowerBound,
+ VariantValue upperBound,
+ int truncateLength) {
this.valueCount = valueCount;
this.nullCount = nullCount;
- this.lowerBound = truncateLowerBound(lowerBound);
- this.upperBound = truncateUpperBound(upperBound);
+ this.lowerBound = truncateLowerBound(lowerBound, truncateLength);
+ this.upperBound = truncateUpperBound(upperBound, truncateLength);
}
VariantMetrics prependFieldName(String name) {
@@ -344,30 +348,30 @@ class ParquetVariantUtil {
return upperBound;
}
- private static VariantValue truncateLowerBound(VariantValue value) {
+ private static VariantValue truncateLowerBound(VariantValue value, int
length) {
switch (value.type()) {
case STRING:
return Variants.of(
PhysicalType.STRING,
- UnicodeUtil.truncateStringMin((String)
value.asPrimitive().get(), 16));
+ UnicodeUtil.truncateStringMin((String)
value.asPrimitive().get(), length));
case BINARY:
return Variants.of(
PhysicalType.BINARY,
- BinaryUtil.truncateBinaryMin((ByteBuffer)
value.asPrimitive().get(), 16));
+ BinaryUtil.truncateBinaryMin((ByteBuffer)
value.asPrimitive().get(), length));
default:
return value;
}
}
- private static VariantValue truncateUpperBound(VariantValue value) {
+ private static VariantValue truncateUpperBound(VariantValue value, int
length) {
switch (value.type()) {
case STRING:
String truncatedString =
- UnicodeUtil.truncateStringMax((String)
value.asPrimitive().get(), 16);
+ UnicodeUtil.truncateStringMax((String)
value.asPrimitive().get(), length);
return truncatedString != null ? Variants.of(PhysicalType.STRING,
truncatedString) : null;
case BINARY:
ByteBuffer truncatedBuffer =
- BinaryUtil.truncateBinaryMax((ByteBuffer)
value.asPrimitive().get(), 16);
+ BinaryUtil.truncateBinaryMax((ByteBuffer)
value.asPrimitive().get(), length);
return truncatedBuffer != null ? Variants.of(PhysicalType.BINARY,
truncatedBuffer) : null;
default:
return value;
diff --git
a/parquet/src/test/java/org/apache/iceberg/parquet/TestVariantMetrics.java
b/parquet/src/test/java/org/apache/iceberg/parquet/TestVariantMetrics.java
index 8381052f58..4392452a14 100644
--- a/parquet/src/test/java/org/apache/iceberg/parquet/TestVariantMetrics.java
+++ b/parquet/src/test/java/org/apache/iceberg/parquet/TestVariantMetrics.java
@@ -30,6 +30,7 @@ import org.apache.hadoop.conf.Configuration;
import org.apache.iceberg.Metrics;
import org.apache.iceberg.MetricsConfig;
import org.apache.iceberg.Schema;
+import org.apache.iceberg.TableProperties;
import org.apache.iceberg.data.GenericRecord;
import org.apache.iceberg.data.Record;
import org.apache.iceberg.data.parquet.InternalWriter;
@@ -72,6 +73,16 @@ public class TestVariantMetrics {
private static final String ROOT_FIELD = "$";
+ private static final byte[] BINARY_20_BYTES = new byte[20];
+ private static final byte[] BINARY_20_BYTES_ALL_FF = new byte[20];
+
+ static {
+ for (int i = 0; i < 20; i += 1) {
+ BINARY_20_BYTES[i] = (byte) (i + 1);
+ BINARY_20_BYTES_ALL_FF[i] = (byte) 0xFF;
+ }
+ }
+
private static final VariantValue[] PRIMITIVES =
new VariantValue[] {
Variants.of(true),
@@ -226,11 +237,7 @@ public class TestVariantMetrics {
@Test
public void testShreddedBinaryBoundsTruncation() throws IOException {
// binary longer than the 16-byte truncation length so the bounds are
truncated
- byte[] bytes = new byte[20];
- for (int i = 0; i < bytes.length; i += 1) {
- bytes[i] = (byte) (i + 1);
- }
- VariantValue value = Variants.of(ByteBuffer.wrap(bytes));
+ VariantValue value = Variants.of(ByteBuffer.wrap(BINARY_20_BYTES));
Metrics metrics =
writeParquet(
@@ -241,21 +248,17 @@ public class TestVariantMetrics {
assertThat(metrics.lowerBounds().get(2))
.extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
-
.isEqualTo(Variants.of(BinaryUtil.truncateBinaryMin(ByteBuffer.wrap(bytes),
16)));
+
.isEqualTo(Variants.of(BinaryUtil.truncateBinaryMin(ByteBuffer.wrap(BINARY_20_BYTES),
16)));
assertThat(metrics.upperBounds().get(2))
.extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
-
.isEqualTo(Variants.of(BinaryUtil.truncateBinaryMax(ByteBuffer.wrap(bytes),
16)));
+
.isEqualTo(Variants.of(BinaryUtil.truncateBinaryMax(ByteBuffer.wrap(BINARY_20_BYTES),
16)));
}
@Test
public void testShreddedBinaryUpperBoundOverflow() throws IOException {
// an all-0xFF binary cannot be truncated up so the upper bound is omitted
- byte[] bytes = new byte[20];
- for (int i = 0; i < bytes.length; i += 1) {
- bytes[i] = (byte) 0xFF;
- }
- VariantValue value = Variants.of(ByteBuffer.wrap(bytes));
+ VariantValue value = Variants.of(ByteBuffer.wrap(BINARY_20_BYTES_ALL_FF));
Metrics metrics =
writeParquet(
@@ -266,13 +269,174 @@ public class TestVariantMetrics {
assertThat(metrics.lowerBounds().get(2))
.extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
-
.isEqualTo(Variants.of(BinaryUtil.truncateBinaryMin(ByteBuffer.wrap(bytes),
16)));
+ .isEqualTo(
+
Variants.of(BinaryUtil.truncateBinaryMin(ByteBuffer.wrap(BINARY_20_BYTES_ALL_FF),
16)));
assertThat(metrics.upperBounds().get(2))
.extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
.isNull();
}
+ @Test
+ public void testShreddedBinaryBoundsTruncateLength() throws IOException {
+ // a per-column truncate(8) overrides the default 16-byte truncation
+ VariantValue value = Variants.of(ByteBuffer.wrap(BINARY_20_BYTES));
+
+ MetricsConfig metricsConfig =
+ MetricsConfig.from(
+ ImmutableMap.of(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX +
"var", "truncate(8)"),
+ SCHEMA,
+ null);
+
+ Metrics metrics =
+ writeParquetWithMetricsConfig(
+ (id, name) -> ParquetVariantUtil.toParquetSchema(value),
+ metricsConfig,
+ Variant.of(EMPTY, value),
+ Variant.of(EMPTY, Variants.ofNull()),
+ null);
+
+ assertThat(metrics.lowerBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+
.isEqualTo(Variants.of(BinaryUtil.truncateBinaryMin(ByteBuffer.wrap(BINARY_20_BYTES),
8)));
+
+ assertThat(metrics.upperBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+
.isEqualTo(Variants.of(BinaryUtil.truncateBinaryMax(ByteBuffer.wrap(BINARY_20_BYTES),
8)));
+ }
+
+ @Test
+ public void testShreddedBinaryBoundsFull() throws IOException {
+ // full mode leaves the bounds untruncated
+ VariantValue value = Variants.of(ByteBuffer.wrap(BINARY_20_BYTES));
+
+ MetricsConfig metricsConfig =
+ MetricsConfig.from(
+ ImmutableMap.of(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX +
"var", "full"),
+ SCHEMA,
+ null);
+
+ Metrics metrics =
+ writeParquetWithMetricsConfig(
+ (id, name) -> ParquetVariantUtil.toParquetSchema(value),
+ metricsConfig,
+ Variant.of(EMPTY, value),
+ Variant.of(EMPTY, Variants.ofNull()),
+ null);
+
+ assertThat(metrics.lowerBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+ .isEqualTo(Variants.of(ByteBuffer.wrap(BINARY_20_BYTES)));
+
+ assertThat(metrics.upperBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+ .isEqualTo(Variants.of(ByteBuffer.wrap(BINARY_20_BYTES)));
+ }
+
+ @Test
+ public void testShreddedBinaryBoundsCounts() throws IOException {
+ // counts mode drops shredded bounds
+ VariantValue value = Variants.of(ByteBuffer.wrap(BINARY_20_BYTES));
+
+ MetricsConfig metricsConfig =
+ MetricsConfig.from(
+ ImmutableMap.of(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX +
"var", "counts"),
+ SCHEMA,
+ null);
+
+ Metrics metrics =
+ writeParquetWithMetricsConfig(
+ (id, name) -> ParquetVariantUtil.toParquetSchema(value),
+ metricsConfig,
+ Variant.of(EMPTY, value),
+ Variant.of(EMPTY, Variants.ofNull()),
+ null);
+
+ assertThat(metrics.valueCounts()).containsKey(2);
+ assertThat(metrics.lowerBounds()).doesNotContainKey(2);
+ assertThat(metrics.upperBounds()).doesNotContainKey(2);
+ }
+
+ @Test
+ public void testShreddedStringBoundsTruncateLength() throws IOException {
+ // a per-column truncate(8) overrides the default 16-char truncation
+ VariantValue value = Variants.of("iceberg_variant");
+
+ MetricsConfig metricsConfig =
+ MetricsConfig.from(
+ ImmutableMap.of(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX +
"var", "truncate(8)"),
+ SCHEMA,
+ null);
+
+ Metrics metrics =
+ writeParquetWithMetricsConfig(
+ (id, name) -> ParquetVariantUtil.toParquetSchema(value),
+ metricsConfig,
+ Variant.of(EMPTY, value),
+ Variant.of(EMPTY, Variants.ofNull()),
+ null);
+
+ assertThat(metrics.lowerBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+
.isEqualTo(Variants.of(UnicodeUtil.truncateStringMin("iceberg_variant", 8)));
+
+ assertThat(metrics.upperBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+
.isEqualTo(Variants.of(UnicodeUtil.truncateStringMax("iceberg_variant", 8)));
+ }
+
+ @Test
+ public void testShreddedStringBoundsFull() throws IOException {
+ // full mode leaves the string bound untruncated
+ VariantValue value = Variants.of("iceberg_variant_full");
+
+ MetricsConfig metricsConfig =
+ MetricsConfig.from(
+ ImmutableMap.of(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX +
"var", "full"),
+ SCHEMA,
+ null);
+
+ Metrics metrics =
+ writeParquetWithMetricsConfig(
+ (id, name) -> ParquetVariantUtil.toParquetSchema(value),
+ metricsConfig,
+ Variant.of(EMPTY, value),
+ Variant.of(EMPTY, Variants.ofNull()),
+ null);
+
+ assertThat(metrics.lowerBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+ .isEqualTo(Variants.of("iceberg_variant_full"));
+
+ assertThat(metrics.upperBounds().get(2))
+ .extracting(b -> Variant.from(b).value().asObject().get(ROOT_FIELD))
+ .isEqualTo(Variants.of("iceberg_variant_full"));
+ }
+
+ @Test
+ public void testShreddedStringBoundsCounts() throws IOException {
+ // counts mode must not truncate the shredded string bound: truncate
length 0 would throw
+ VariantValue value = Variants.of("iceberg_variant");
+
+ MetricsConfig metricsConfig =
+ MetricsConfig.from(
+ ImmutableMap.of(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX +
"var", "counts"),
+ SCHEMA,
+ null);
+
+ Metrics metrics =
+ writeParquetWithMetricsConfig(
+ (id, name) -> ParquetVariantUtil.toParquetSchema(value),
+ metricsConfig,
+ Variant.of(EMPTY, value),
+ Variant.of(EMPTY, Variants.ofNull()),
+ null);
+
+ assertThat(metrics.valueCounts()).containsKey(2);
+ assertThat(metrics.lowerBounds()).doesNotContainKey(2);
+ assertThat(metrics.upperBounds()).doesNotContainKey(2);
+ }
+
@Test
public void testVariantFloatNaN() throws IOException {
// NaN values are not counted because there is no ID for FieldMetrics
@@ -578,6 +742,12 @@ public class TestVariantMetrics {
private Metrics writeParquet(VariantShreddingFunction shredding, Variant...
variants)
throws IOException {
+ return writeParquetWithMetricsConfig(shredding,
MetricsConfig.getDefault(), variants);
+ }
+
+ private Metrics writeParquetWithMetricsConfig(
+ VariantShreddingFunction shredding, MetricsConfig metricsConfig,
Variant... variants)
+ throws IOException {
OutputFile out = new InMemoryOutputFile();
GenericRecord record = GenericRecord.create(SCHEMA);
@@ -585,6 +755,7 @@ public class TestVariantMetrics {
Parquet.write(out)
.schema(SCHEMA)
.variantShreddingFunc(shredding)
+ .metricsConfig(metricsConfig)
.createWriterFunc(fileSchema ->
InternalWriter.create(SCHEMA.asStruct(), fileSchema))
.build();