Guosmilesmile commented on code in PR #17553:
URL: https://github.com/apache/iceberg/pull/17553#discussion_r3828031109


##########
parquet/src/test/java/org/apache/iceberg/parquet/TestVariantWriters.java:
##########
@@ -261,42 +271,157 @@ public void testMixedShredding(Variant variant) throws 
IOException {
     }
   }
 
+  @Test
+  public void testPartialShreddingWithShreddedObject() throws IOException {
+    VariantMetadata metadata = Variants.metadata("id", "name", "city");
+
+    List<Record> records = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      ShreddedObject obj = Variants.object(metadata);
+      obj.put("id", Variants.of(1000L + i));
+      obj.put("name", Variants.of("user_" + i));
+      obj.put("city", Variants.of("city_" + i));
+
+      Variant variant = Variant.of(metadata, obj);
+      Record record = RECORD.copy("id", i, "var", variant);
+      records.add(record);
+    }
+
+    VariantShreddingFunction partialShredding = (id, name) -> shredOnly("id");
+
+    List<Record> actual = writeAndRead(SCHEMA, partialShredding, records, 
"var");
+
+    assertThat(actual).hasSameSizeAs(records);
+    for (int i = 0; i < records.size(); i++) {
+      Record expected = records.get(i);
+      Record read = actual.get(i);
+
+      InternalTestHelpers.assertEquals(SCHEMA.asStruct(), expected, read);
+
+      Variant readVariant = (Variant) read.getField("var");
+      VariantObject readObj = readVariant.value().asObject();
+      assertThat(readObj.numFields()).isEqualTo(3);
+      assertThat(readObj.get("id").asPrimitive().get()).isEqualTo(1000L + i);
+      assertThat(readObj.get("name").asPrimitive().get()).isEqualTo("user_" + 
i);
+      assertThat(readObj.get("city").asPrimitive().get()).isEqualTo("city_" + 
i);
+    }
+  }
+
+  @Test
+  public void testPartialShreddingMultipleColumns() throws IOException {
+    VariantMetadata metadata1 = Variants.metadata("id", "name", "city");
+    VariantMetadata metadata2 = Variants.metadata("key", "val", "extra");
+
+    List<Record> records = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      ShreddedObject object1 = Variants.object(metadata1);
+      object1.put("id", Variants.of(1000L + i));
+      object1.put("name", Variants.of("user_" + i));
+      object1.put("city", Variants.of("city_" + i));
+
+      ShreddedObject object2 = Variants.object(metadata2);
+      object2.put("key", Variants.of(2000L + i));
+      object2.put("val", Variants.of("val_" + i));
+      object2.put("extra", Variants.of("extra_" + i));
+
+      records.add(
+          RECORD_TWO_VARIANTS.copy(
+              "id", i,
+              "var1", Variant.of(metadata1, object1),
+              "var2", Variant.of(metadata2, object2)));
+    }
+
+    VariantShreddingFunction partialShredding =
+        (id, name) -> {
+          if (name.equals("var1")) {
+            return shredOnly("id");
+          } else if (name.equals("var2")) {
+            return shredOnly("key");
+          }
+          return null;
+        };
+
+    List<Record> actual =
+        writeAndRead(SCHEMA_TWO_VARIANTS, partialShredding, records, "var1", 
"var2");
+
+    assertThat(actual).hasSameSizeAs(records);
+    for (int i = 0; i < records.size(); i++) {
+      InternalTestHelpers.assertEquals(
+          SCHEMA_TWO_VARIANTS.asStruct(), records.get(i), actual.get(i));
+
+      VariantObject readObject1 = ((Variant) 
actual.get(i).getField("var1")).value().asObject();
+      assertThat(readObject1.numFields()).isEqualTo(3);
+      assertThat(readObject1.get("id").asPrimitive().get()).isEqualTo(1000L + 
i);
+      
assertThat(readObject1.get("name").asPrimitive().get()).isEqualTo("user_" + i);
+      
assertThat(readObject1.get("city").asPrimitive().get()).isEqualTo("city_" + i);
+
+      VariantObject readObject2 = ((Variant) 
actual.get(i).getField("var2")).value().asObject();
+      assertThat(readObject2.numFields()).isEqualTo(3);
+      assertThat(readObject2.get("key").asPrimitive().get()).isEqualTo(2000L + 
i);
+      assertThat(readObject2.get("val").asPrimitive().get()).isEqualTo("val_" 
+ i);
+      
assertThat(readObject2.get("extra").asPrimitive().get()).isEqualTo("extra_" + 
i);
+    }
+  }
+
   private static Record writeAndRead(VariantShreddingFunction shreddingFunc, 
Record record)
       throws IOException {
     return Iterables.getOnlyElement(writeAndRead(shreddingFunc, 
List.of(record)));
   }
 
   private static List<Record> writeAndRead(
       VariantShreddingFunction shreddingFunc, List<Record> records) throws 
IOException {
+    return writeAndRead(SCHEMA, shreddingFunc, records);
+  }
+
+  private static List<Record> writeAndRead(
+      Schema schema, VariantShreddingFunction shreddingFunc, List<Record> 
records)
+      throws IOException {
+    return writeAndRead(schema, shreddingFunc, records, new String[0]);
+  }

Review Comment:
   Can we remove this method only use 
   ```
   private static List<Record> writeAndRead(
         Schema schema,
         VariantShreddingFunction shreddingFunc,
         List<Record> records,
         String... shreddedColumns)
   ``` 
   



##########
spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestSparkVariantRead.java:
##########
@@ -478,37 +434,96 @@ public void testReadShreddedWithMetricsDisabled(String 
metricsMode)
             + "'write.metadata.metrics.default'='%s')",
         noStatsTable, metricsMode);
 
-    spark.conf().set("spark.sql.iceberg.shred-variants", "true");
-    try {
-      sql(
-          "INSERT INTO %s VALUES "
-              + "(1, parse_json('{\"name\":\"alice\",\"age\":30}')), "
-              + "(2, parse_json('{\"name\":\"bob\",\"age\":25}'))",
-          noStatsTable);
-    } finally {
-      spark.conf().unset("spark.sql.iceberg.shred-variants");
-    }
+    insertShredded(
+        "INSERT INTO %s VALUES "
+            + "(1, parse_json('{\"name\":\"alice\",\"age\":30}')), "
+            + "(2, parse_json('{\"name\":\"bob\",\"age\":25}'))",
+        noStatsTable);
 
     Table table = Spark3Util.loadIcebergTable(spark, noStatsTable);
     assertHasTypedValueSubtree(table);
     setVectorization(noStatsTable, true);
 
     List<Row> rows = spark.table(noStatsTable).select("id", 
"v").orderBy("id").collectAsList();
     assertThat(rows).hasSize(2);
-    Variant v1 =
-        new Variant(
-            ((VariantVal) rows.get(0).get(1)).getValue(),
-            ((VariantVal) rows.get(0).get(1)).getMetadata());
+    Variant v1 = asVariant(rows.get(0), 1);
     assertThat(v1.getFieldByKey("name").getString()).isEqualTo("alice");
-    Variant v2 =
-        new Variant(
-            ((VariantVal) rows.get(1).get(1)).getValue(),
-            ((VariantVal) rows.get(1).get(1)).getMetadata());
+    Variant v2 = asVariant(rows.get(1), 1);
     assertThat(v2.getFieldByKey("name").getString()).isEqualTo("bob");
 
     sql("DROP TABLE IF EXISTS %s", noStatsTable);
   }
 
+  @ParameterizedTest
+  @ValueSource(booleans = {false, true})
+  public void testMultipleUnshreddedVariantColumns(boolean vectorized)
+      throws IOException, NoSuchTableException, ParseException {
+    assertMultipleVariantColumns(false, vectorized);
+  }
+
+  @Test
+  public void testMultipleShreddedVariantColumns()
+      throws IOException, NoSuchTableException, ParseException {
+    assertMultipleVariantColumns(true, true);

Review Comment:
   Should we test novectorized?



##########
parquet/src/test/java/org/apache/iceberg/parquet/TestVariantWriters.java:
##########
@@ -261,42 +271,157 @@ public void testMixedShredding(Variant variant) throws 
IOException {
     }
   }
 
+  @Test
+  public void testPartialShreddingWithShreddedObject() throws IOException {
+    VariantMetadata metadata = Variants.metadata("id", "name", "city");
+
+    List<Record> records = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      ShreddedObject obj = Variants.object(metadata);
+      obj.put("id", Variants.of(1000L + i));
+      obj.put("name", Variants.of("user_" + i));
+      obj.put("city", Variants.of("city_" + i));
+
+      Variant variant = Variant.of(metadata, obj);
+      Record record = RECORD.copy("id", i, "var", variant);
+      records.add(record);
+    }
+
+    VariantShreddingFunction partialShredding = (id, name) -> shredOnly("id");
+
+    List<Record> actual = writeAndRead(SCHEMA, partialShredding, records, 
"var");
+
+    assertThat(actual).hasSameSizeAs(records);
+    for (int i = 0; i < records.size(); i++) {
+      Record expected = records.get(i);
+      Record read = actual.get(i);
+
+      InternalTestHelpers.assertEquals(SCHEMA.asStruct(), expected, read);
+
+      Variant readVariant = (Variant) read.getField("var");
+      VariantObject readObj = readVariant.value().asObject();
+      assertThat(readObj.numFields()).isEqualTo(3);
+      assertThat(readObj.get("id").asPrimitive().get()).isEqualTo(1000L + i);
+      assertThat(readObj.get("name").asPrimitive().get()).isEqualTo("user_" + 
i);
+      assertThat(readObj.get("city").asPrimitive().get()).isEqualTo("city_" + 
i);
+    }
+  }
+
+  @Test
+  public void testPartialShreddingMultipleColumns() throws IOException {
+    VariantMetadata metadata1 = Variants.metadata("id", "name", "city");
+    VariantMetadata metadata2 = Variants.metadata("key", "val", "extra");
+
+    List<Record> records = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      ShreddedObject object1 = Variants.object(metadata1);
+      object1.put("id", Variants.of(1000L + i));
+      object1.put("name", Variants.of("user_" + i));
+      object1.put("city", Variants.of("city_" + i));
+
+      ShreddedObject object2 = Variants.object(metadata2);
+      object2.put("key", Variants.of(2000L + i));
+      object2.put("val", Variants.of("val_" + i));
+      object2.put("extra", Variants.of("extra_" + i));
+
+      records.add(
+          RECORD_TWO_VARIANTS.copy(
+              "id", i,
+              "var1", Variant.of(metadata1, object1),
+              "var2", Variant.of(metadata2, object2)));
+    }
+
+    VariantShreddingFunction partialShredding =
+        (id, name) -> {
+          if (name.equals("var1")) {
+            return shredOnly("id");
+          } else if (name.equals("var2")) {
+            return shredOnly("key");
+          }
+          return null;
+        };
+
+    List<Record> actual =
+        writeAndRead(SCHEMA_TWO_VARIANTS, partialShredding, records, "var1", 
"var2");
+
+    assertThat(actual).hasSameSizeAs(records);
+    for (int i = 0; i < records.size(); i++) {
+      InternalTestHelpers.assertEquals(
+          SCHEMA_TWO_VARIANTS.asStruct(), records.get(i), actual.get(i));
+
+      VariantObject readObject1 = ((Variant) 
actual.get(i).getField("var1")).value().asObject();
+      assertThat(readObject1.numFields()).isEqualTo(3);
+      assertThat(readObject1.get("id").asPrimitive().get()).isEqualTo(1000L + 
i);
+      
assertThat(readObject1.get("name").asPrimitive().get()).isEqualTo("user_" + i);
+      
assertThat(readObject1.get("city").asPrimitive().get()).isEqualTo("city_" + i);
+
+      VariantObject readObject2 = ((Variant) 
actual.get(i).getField("var2")).value().asObject();
+      assertThat(readObject2.numFields()).isEqualTo(3);
+      assertThat(readObject2.get("key").asPrimitive().get()).isEqualTo(2000L + 
i);
+      assertThat(readObject2.get("val").asPrimitive().get()).isEqualTo("val_" 
+ i);
+      
assertThat(readObject2.get("extra").asPrimitive().get()).isEqualTo("extra_" + 
i);
+    }
+  }
+
   private static Record writeAndRead(VariantShreddingFunction shreddingFunc, 
Record record)
       throws IOException {
     return Iterables.getOnlyElement(writeAndRead(shreddingFunc, 
List.of(record)));
   }
 
   private static List<Record> writeAndRead(
       VariantShreddingFunction shreddingFunc, List<Record> records) throws 
IOException {
+    return writeAndRead(SCHEMA, shreddingFunc, records);
+  }
+
+  private static List<Record> writeAndRead(
+      Schema schema, VariantShreddingFunction shreddingFunc, List<Record> 
records)
+      throws IOException {
+    return writeAndRead(schema, shreddingFunc, records, new String[0]);
+  }
+
+  private static List<Record> writeAndRead(
+      Schema schema,
+      VariantShreddingFunction shreddingFunc,
+      List<Record> records,
+      String... shreddedColumns)
+      throws IOException {
     OutputFile outputFile = new InMemoryOutputFile();
 
     try (FileAppender<Record> writer =
         Parquet.write(outputFile)
-            .schema(SCHEMA)
+            .schema(schema)
             .variantShreddingFunc(shreddingFunc)
-            .createWriterFunc(fileSchema -> 
InternalWriter.create(SCHEMA.asStruct(), fileSchema))
+            .createWriterFunc(fileSchema -> 
InternalWriter.create(schema.asStruct(), fileSchema))
             .build()) {
       for (Record record : records) {
         writer.add(record);
       }
     }
 
+    List<String> expectedShredded = ImmutableList.copyOf(shreddedColumns);
     try (ParquetFileReader reader =
         ParquetFileReader.open(ParquetIO.file(outputFile.toInputFile()))) {
-      MessageType schema = reader.getFileMetaData().getSchema();
-      for (Types.NestedField column : SCHEMA.columns()) {
+      MessageType parquetSchema = reader.getFileMetaData().getSchema();
+      for (Types.NestedField column : schema.columns()) {
         if (column.type() == Types.VariantType.get()) {

Review Comment:
   May be if (column.type().isVariantType()) ?



##########
spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestSparkVariantRead.java:
##########
@@ -522,15 +537,36 @@ private void setVectorization(String table, boolean on) {
   }
 
   private static void assertHasTypedValueSubtree(Table table) throws 
IOException {
+    forEachDataFileSchema(
+        table,
+        schema ->
+            assertThat(containsTypedValue(schema))
+                .as("Expected variant column to be shredded with a typed_value 
subtree")
+                .isTrue());
+  }
+
+  private static void assertHasTypedValueSubtree(Table table, String... 
columns)
+      throws IOException {
+    forEachDataFileSchema(
+        table,
+        schema -> {
+          for (String column : columns) {
+            assertThat(containsTypedValue(schema.getType(column)))
+                .as("Expected column %s to be shredded with a typed_value 
subtree", column)
+                .isTrue();
+          }
+        });
+  }
+
+  private static void forEachDataFileSchema(
+      Table table, Consumer<org.apache.parquet.schema.MessageType> 
schemaCheck) throws IOException {

Review Comment:
   nit: import 



##########
parquet/src/test/java/org/apache/iceberg/parquet/TestVariantWriters.java:
##########
@@ -261,42 +271,157 @@ public void testMixedShredding(Variant variant) throws 
IOException {
     }
   }
 
+  @Test
+  public void testPartialShreddingWithShreddedObject() throws IOException {
+    VariantMetadata metadata = Variants.metadata("id", "name", "city");
+
+    List<Record> records = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      ShreddedObject obj = Variants.object(metadata);
+      obj.put("id", Variants.of(1000L + i));
+      obj.put("name", Variants.of("user_" + i));
+      obj.put("city", Variants.of("city_" + i));
+
+      Variant variant = Variant.of(metadata, obj);
+      Record record = RECORD.copy("id", i, "var", variant);
+      records.add(record);
+    }
+
+    VariantShreddingFunction partialShredding = (id, name) -> shredOnly("id");
+
+    List<Record> actual = writeAndRead(SCHEMA, partialShredding, records, 
"var");
+
+    assertThat(actual).hasSameSizeAs(records);
+    for (int i = 0; i < records.size(); i++) {
+      Record expected = records.get(i);
+      Record read = actual.get(i);
+
+      InternalTestHelpers.assertEquals(SCHEMA.asStruct(), expected, read);
+
+      Variant readVariant = (Variant) read.getField("var");
+      VariantObject readObj = readVariant.value().asObject();
+      assertThat(readObj.numFields()).isEqualTo(3);
+      assertThat(readObj.get("id").asPrimitive().get()).isEqualTo(1000L + i);
+      assertThat(readObj.get("name").asPrimitive().get()).isEqualTo("user_" + 
i);
+      assertThat(readObj.get("city").asPrimitive().get()).isEqualTo("city_" + 
i);
+    }

Review Comment:
   This test is named “partial shredding,” but I’m not sure it actually 
distinguishes between partial and full shredding. containsTypedValue() only 
checks for the existence of a typed_value, and the round-trip assertions would 
still pass even if the file were fully shredded. Please correct me if I’m 
misunderstanding this.



##########
parquet/src/test/java/org/apache/iceberg/parquet/TestVariantWriters.java:
##########
@@ -261,42 +271,157 @@ public void testMixedShredding(Variant variant) throws 
IOException {
     }
   }
 
+  @Test
+  public void testPartialShreddingWithShreddedObject() throws IOException {
+    VariantMetadata metadata = Variants.metadata("id", "name", "city");
+
+    List<Record> records = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      ShreddedObject obj = Variants.object(metadata);
+      obj.put("id", Variants.of(1000L + i));
+      obj.put("name", Variants.of("user_" + i));
+      obj.put("city", Variants.of("city_" + i));
+
+      Variant variant = Variant.of(metadata, obj);
+      Record record = RECORD.copy("id", i, "var", variant);
+      records.add(record);
+    }
+
+    VariantShreddingFunction partialShredding = (id, name) -> shredOnly("id");
+
+    List<Record> actual = writeAndRead(SCHEMA, partialShredding, records, 
"var");
+
+    assertThat(actual).hasSameSizeAs(records);
+    for (int i = 0; i < records.size(); i++) {
+      Record expected = records.get(i);
+      Record read = actual.get(i);
+
+      InternalTestHelpers.assertEquals(SCHEMA.asStruct(), expected, read);
+
+      Variant readVariant = (Variant) read.getField("var");
+      VariantObject readObj = readVariant.value().asObject();
+      assertThat(readObj.numFields()).isEqualTo(3);
+      assertThat(readObj.get("id").asPrimitive().get()).isEqualTo(1000L + i);
+      assertThat(readObj.get("name").asPrimitive().get()).isEqualTo("user_" + 
i);
+      assertThat(readObj.get("city").asPrimitive().get()).isEqualTo("city_" + 
i);
+    }
+  }
+
+  @Test
+  public void testPartialShreddingMultipleColumns() throws IOException {
+    VariantMetadata metadata1 = Variants.metadata("id", "name", "city");
+    VariantMetadata metadata2 = Variants.metadata("key", "val", "extra");
+
+    List<Record> records = Lists.newArrayList();
+    for (int i = 0; i < 3; i++) {
+      ShreddedObject object1 = Variants.object(metadata1);
+      object1.put("id", Variants.of(1000L + i));
+      object1.put("name", Variants.of("user_" + i));
+      object1.put("city", Variants.of("city_" + i));
+
+      ShreddedObject object2 = Variants.object(metadata2);
+      object2.put("key", Variants.of(2000L + i));
+      object2.put("val", Variants.of("val_" + i));
+      object2.put("extra", Variants.of("extra_" + i));
+
+      records.add(
+          RECORD_TWO_VARIANTS.copy(
+              "id", i,
+              "var1", Variant.of(metadata1, object1),
+              "var2", Variant.of(metadata2, object2)));
+    }
+
+    VariantShreddingFunction partialShredding =
+        (id, name) -> {
+          if (name.equals("var1")) {
+            return shredOnly("id");
+          } else if (name.equals("var2")) {
+            return shredOnly("key");
+          }
+          return null;
+        };
+
+    List<Record> actual =
+        writeAndRead(SCHEMA_TWO_VARIANTS, partialShredding, records, "var1", 
"var2");
+
+    assertThat(actual).hasSameSizeAs(records);
+    for (int i = 0; i < records.size(); i++) {
+      InternalTestHelpers.assertEquals(
+          SCHEMA_TWO_VARIANTS.asStruct(), records.get(i), actual.get(i));
+
+      VariantObject readObject1 = ((Variant) 
actual.get(i).getField("var1")).value().asObject();
+      assertThat(readObject1.numFields()).isEqualTo(3);
+      assertThat(readObject1.get("id").asPrimitive().get()).isEqualTo(1000L + 
i);
+      
assertThat(readObject1.get("name").asPrimitive().get()).isEqualTo("user_" + i);
+      
assertThat(readObject1.get("city").asPrimitive().get()).isEqualTo("city_" + i);
+
+      VariantObject readObject2 = ((Variant) 
actual.get(i).getField("var2")).value().asObject();
+      assertThat(readObject2.numFields()).isEqualTo(3);
+      assertThat(readObject2.get("key").asPrimitive().get()).isEqualTo(2000L + 
i);
+      assertThat(readObject2.get("val").asPrimitive().get()).isEqualTo("val_" 
+ i);
+      
assertThat(readObject2.get("extra").asPrimitive().get()).isEqualTo("extra_" + 
i);
+    }
+  }
+
   private static Record writeAndRead(VariantShreddingFunction shreddingFunc, 
Record record)
       throws IOException {
     return Iterables.getOnlyElement(writeAndRead(shreddingFunc, 
List.of(record)));
   }
 
   private static List<Record> writeAndRead(
       VariantShreddingFunction shreddingFunc, List<Record> records) throws 
IOException {
+    return writeAndRead(SCHEMA, shreddingFunc, records);
+  }
+
+  private static List<Record> writeAndRead(
+      Schema schema, VariantShreddingFunction shreddingFunc, List<Record> 
records)
+      throws IOException {
+    return writeAndRead(schema, shreddingFunc, records, new String[0]);
+  }
+
+  private static List<Record> writeAndRead(
+      Schema schema,
+      VariantShreddingFunction shreddingFunc,
+      List<Record> records,
+      String... shreddedColumns)
+      throws IOException {
     OutputFile outputFile = new InMemoryOutputFile();
 
     try (FileAppender<Record> writer =
         Parquet.write(outputFile)
-            .schema(SCHEMA)
+            .schema(schema)
             .variantShreddingFunc(shreddingFunc)
-            .createWriterFunc(fileSchema -> 
InternalWriter.create(SCHEMA.asStruct(), fileSchema))
+            .createWriterFunc(fileSchema -> 
InternalWriter.create(schema.asStruct(), fileSchema))
             .build()) {
       for (Record record : records) {
         writer.add(record);
       }
     }
 
+    List<String> expectedShredded = ImmutableList.copyOf(shreddedColumns);
     try (ParquetFileReader reader =
         ParquetFileReader.open(ParquetIO.file(outputFile.toInputFile()))) {
-      MessageType schema = reader.getFileMetaData().getSchema();
-      for (Types.NestedField column : SCHEMA.columns()) {
+      MessageType parquetSchema = reader.getFileMetaData().getSchema();
+      for (Types.NestedField column : schema.columns()) {
         if (column.type() == Types.VariantType.get()) {
-          int fieldIndex = schema.getFieldIndex(column.name());
-          
assertThat(schema.getFields().get(fieldIndex).getLogicalTypeAnnotation())
+          int fieldIndex = parquetSchema.getFieldIndex(column.name());
+          
assertThat(parquetSchema.getFields().get(fieldIndex).getLogicalTypeAnnotation())

Review Comment:
   Can we use `parquetSchema.getType(column.name()).getLogicalTypeAnnotation()` 
instead ? 



##########
spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestSparkVariantRead.java:
##########
@@ -522,15 +537,36 @@ private void setVectorization(String table, boolean on) {
   }
 
   private static void assertHasTypedValueSubtree(Table table) throws 
IOException {
+    forEachDataFileSchema(
+        table,
+        schema ->
+            assertThat(containsTypedValue(schema))
+                .as("Expected variant column to be shredded with a typed_value 
subtree")
+                .isTrue());
+  }
+
+  private static void assertHasTypedValueSubtree(Table table, String... 
columns)
+      throws IOException {
+    forEachDataFileSchema(
+        table,
+        schema -> {
+          for (String column : columns) {
+            assertThat(containsTypedValue(schema.getType(column)))
+                .as("Expected column %s to be shredded with a typed_value 
subtree", column)
+                .isTrue();
+          }
+        });
+  }
+
+  private static void forEachDataFileSchema(
+      Table table, Consumer<org.apache.parquet.schema.MessageType> 
schemaCheck) throws IOException {

Review Comment:
   nit: import 



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to