diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvro.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvro.java index 9783880a5985..e2afff44f3f1 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvro.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvro.java @@ -38,6 +38,7 @@ import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.types.TypeUtil; +import org.apache.iceberg.variants.Variant; class ParquetAvro { @@ -220,6 +221,9 @@ public Conversion getConversionByClass( } } else if ("uuid".equals(logicalType.getName())) { return (Conversion) uuidConversion; + } else if ("variant".equals(logicalType.getName()) + && Variant.class.isAssignableFrom(datumClass)) { + return (Conversion) variantConversion; } return super.getConversionByClass(datumClass, logicalType); } diff --git a/parquet/src/test/java/org/apache/iceberg/parquet/TestParquet.java b/parquet/src/test/java/org/apache/iceberg/parquet/TestParquet.java index 54f664c02c36..12861a2cb526 100644 --- a/parquet/src/test/java/org/apache/iceberg/parquet/TestParquet.java +++ b/parquet/src/test/java/org/apache/iceberg/parquet/TestParquet.java @@ -58,6 +58,7 @@ import org.apache.iceberg.avro.AvroSchemaUtil; import org.apache.iceberg.data.Record; import org.apache.iceberg.data.parquet.GenericParquetWriter; +import org.apache.iceberg.data.parquet.InternalReader; import org.apache.iceberg.expressions.Expression; import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.io.FileAppender; @@ -73,6 +74,8 @@ import org.apache.iceberg.util.Pair; import org.apache.iceberg.util.RandomUtil; import org.apache.iceberg.variants.Variant; +import org.apache.iceberg.variants.VariantTestUtil; +import org.apache.iceberg.variants.Variants; import org.apache.parquet.avro.AvroParquetWriter; import org.apache.parquet.column.Encoding; import org.apache.parquet.column.statistics.Statistics; @@ -602,6 +605,31 @@ public void testAvroWriterRejectsVariantType() { .hasMessage("Avro writer does not support variant types"); } + @Test + void writesVariantWithDefaultAvroWriter() throws IOException { + Schema schema = new Schema(required(1, "v", Types.VariantType.get())); + GenericData.Record record = + new GenericData.Record(AvroSchemaUtil.convert(schema.asStruct(), "table")); + Variant expected = Variant.of(Variants.emptyMetadata(), Variants.of(34)); + record.put("v", expected); + + File file = createTempFile(temp); + try (FileAppender writer = + Parquet.write(Files.localOutput(file)).schema(schema).build()) { + writer.add(record); + } + + try (CloseableIterable rows = + Parquet.read(Files.localInput(file)) + .project(schema) + .createReaderFunc(fileSchema -> InternalReader.create(schema, fileSchema)) + .build()) { + Variant actual = (Variant) getOnlyElement(rows).get(0); + VariantTestUtil.assertEqual(expected.metadata(), actual.metadata()); + VariantTestUtil.assertEqual(expected.value(), actual.value()); + } + } + @Test public void adaptiveBloomFilterSizingShrinksFile() throws IOException { // when PARQUET_BLOOM_FILTER_ADAPTIVE_ENABLED is not set (the default), the writer