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();
 

Reply via email to