From b8f721f6cfb97285f55a6525c320ef0b13729766 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Thu, 12 Sep 2024 20:58:11 -0700 Subject: [PATCH 01/11] Core: add variant type support --- .../apache/iceberg/transforms/Identity.java | 5 +- .../apache/iceberg/types/AssignFreshIds.java | 5 + .../apache/iceberg/types/GetProjectedIds.java | 4 +- .../org/apache/iceberg/types/ReassignIds.java | 5 + .../java/org/apache/iceberg/types/Type.java | 4 + .../org/apache/iceberg/types/TypeUtil.java | 7 + .../java/org/apache/iceberg/types/Types.java | 5 + .../apache/iceberg/types/TestTypeUtil.java | 162 +++++++++--------- .../java/org/apache/iceberg/SchemaParser.java | 7 + .../iceberg/avro/BuildAvroProjection.java | 15 ++ .../org/apache/iceberg/avro/TypeToSchema.java | 14 ++ .../iceberg/TestMetadataUpdateParser.java | 25 ++- .../avro/TestAvroSchemaProjection.java | 14 ++ .../iceberg/avro/TestBuildAvroProjection.java | 29 ++++ 14 files changed, 216 insertions(+), 85 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/transforms/Identity.java b/api/src/main/java/org/apache/iceberg/transforms/Identity.java index 099a99cc3cf4..f0ca938a402c 100644 --- a/api/src/main/java/org/apache/iceberg/transforms/Identity.java +++ b/api/src/main/java/org/apache/iceberg/transforms/Identity.java @@ -19,14 +19,17 @@ package org.apache.iceberg.transforms; import java.io.ObjectStreamException; +import java.util.Set; import org.apache.iceberg.expressions.BoundPredicate; import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.expressions.UnboundPredicate; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; import org.apache.iceberg.util.SerializableFunction; class Identity implements Transform { + private static final Set UNSUPPORTED_TYPES = Set.of(Types.VariantType.get()); private static final Identity INSTANCE = new Identity<>(); private final Type type; @@ -39,7 +42,7 @@ class Identity implements Transform { @Deprecated public static Identity get(Type type) { Preconditions.checkArgument( - type.typeId() != Type.TypeID.VARIANT, "Unsupported type for identity: %s", type); + !UNSUPPORTED_TYPES.contains(type), "Unsupported type for identity: %s", type); return new Identity<>(type); } diff --git a/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java b/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java index e58f76a8de56..f6422671bacb 100644 --- a/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java +++ b/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java @@ -124,6 +124,11 @@ public Type map(Types.MapType map, Supplier keyFuture, Supplier valu } } + @Override + public Type variant() { + return Types.VariantType.get(); + } + @Override public Type primitive(Type.PrimitiveType primitive) { return primitive; diff --git a/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java b/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java index a8a7de065ece..4692a8c400f8 100644 --- a/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java +++ b/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java @@ -47,7 +47,9 @@ public Set struct(Types.StructType struct, List> fieldResu @Override public Set field(Types.NestedField field, Set fieldResult) { - if ((includeStructIds && field.type().isStructType()) || field.type().isPrimitiveType()) { + if ((includeStructIds && field.type().isStructType()) + || field.type().isPrimitiveType() + || field.type().isVariantType()) { fieldIds.add(field.fieldId()); } return fieldIds; diff --git a/api/src/main/java/org/apache/iceberg/types/ReassignIds.java b/api/src/main/java/org/apache/iceberg/types/ReassignIds.java index 565ceee2a901..dd737f5308d3 100644 --- a/api/src/main/java/org/apache/iceberg/types/ReassignIds.java +++ b/api/src/main/java/org/apache/iceberg/types/ReassignIds.java @@ -157,6 +157,11 @@ public Type map(Types.MapType map, Supplier keyTypeFuture, Supplier } } + @Override + public Type variant() { + return Types.VariantType.get(); + } + @Override public Type primitive(Type.PrimitiveType primitive) { return primitive; // nothing to reassign diff --git a/api/src/main/java/org/apache/iceberg/types/Type.java b/api/src/main/java/org/apache/iceberg/types/Type.java index f4c6f22134a5..53018ffac65b 100644 --- a/api/src/main/java/org/apache/iceberg/types/Type.java +++ b/api/src/main/java/org/apache/iceberg/types/Type.java @@ -98,6 +98,10 @@ default boolean isMapType() { return false; } + default boolean isVariantType() { + return false; + } + default NestedType asNestedType() { throw new IllegalArgumentException("Not a nested type: " + this); } diff --git a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java index 39f2898757a6..19ef2d155f26 100644 --- a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java +++ b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java @@ -712,6 +712,10 @@ public T map(Types.MapType map, Supplier keyResult, Supplier valueResult) return null; } + public T variant() { + return null; + } + public T primitive(Type.PrimitiveType primitive) { return null; } @@ -788,6 +792,9 @@ public static T visit(Type type, CustomOrderSchemaVisitor visitor) { new VisitFuture<>(map.keyType(), visitor), new VisitFuture<>(map.valueType(), visitor)); + case VARIANT: + return visitor.variant(); + default: return visitor.primitive(type.asPrimitiveType()); } diff --git a/api/src/main/java/org/apache/iceberg/types/Types.java b/api/src/main/java/org/apache/iceberg/types/Types.java index 6882f718508b..2c7c3ed81846 100644 --- a/api/src/main/java/org/apache/iceberg/types/Types.java +++ b/api/src/main/java/org/apache/iceberg/types/Types.java @@ -430,6 +430,11 @@ public String toString() { return "variant"; } + @Override + public boolean isVariantType() { + return true; + } + @Override public boolean equals(Object o) { if (this == o) { diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java index 36384d232af3..e6740d8efc8b 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java @@ -24,38 +24,43 @@ import static org.assertj.core.api.Assertions.assertThatThrownBy; import java.util.Set; +import java.util.stream.Stream; import org.apache.iceberg.Schema; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Sets; -import org.apache.iceberg.types.Types.IntegerType; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; public class TestTypeUtil { - @Test - public void testReassignIdsDuplicateColumns() { + private static Stream testTypes() { + return Stream.of(Arguments.of(Types.IntegerType.get()), Arguments.of(Types.VariantType.get())); + } + + @ParameterizedTest + @MethodSource("testTypes") + public void testReassignIdsDuplicateColumns(Type testType) { Schema schema = - new Schema( - required(0, "a", Types.IntegerType.get()), required(1, "A", Types.IntegerType.get())); + new Schema(required(0, "a", testType), required(1, "A", Types.IntegerType.get())); Schema sourceSchema = - new Schema( - required(1, "a", Types.IntegerType.get()), required(2, "A", Types.IntegerType.get())); + new Schema(required(1, "a", testType), required(2, "A", Types.IntegerType.get())); final Schema actualSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(actualSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); } - @Test - public void testReassignIdsWithIdentifier() { + @ParameterizedTest + @MethodSource("testTypes") + public void testReassignIdsWithIdentifier(Type testType) { Schema schema = new Schema( Lists.newArrayList( - required(0, "a", Types.IntegerType.get()), - required(1, "A", Types.IntegerType.get())), + required(0, "a", Types.IntegerType.get()), required(1, "A", testType)), Sets.newHashSet(0)); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), - required(2, "A", Types.IntegerType.get())), + required(1, "a", Types.IntegerType.get()), required(2, "A", testType)), Sets.newHashSet(1)); final Schema actualSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(actualSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); @@ -64,19 +69,18 @@ public void testReassignIdsWithIdentifier() { .isEqualTo(sourceSchema.identifierFieldIds()); } - @Test - public void testAssignIncreasingFreshIdWithIdentifier() { + @ParameterizedTest + @MethodSource("testTypes") + public void testAssignIncreasingFreshIdWithIdentifier(Type testType) { Schema schema = new Schema( Lists.newArrayList( - required(10, "a", Types.IntegerType.get()), - required(11, "A", Types.IntegerType.get())), + required(10, "a", Types.IntegerType.get()), required(11, "A", testType)), Sets.newHashSet(10)); Schema expectedSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), - required(2, "A", Types.IntegerType.get())), + required(1, "a", Types.IntegerType.get()), required(2, "A", testType)), Sets.newHashSet(1)); final Schema actualSchema = TypeUtil.assignIncreasingFreshIds(schema); assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); @@ -85,19 +89,18 @@ public void testAssignIncreasingFreshIdWithIdentifier() { .isEqualTo(expectedSchema.identifierFieldIds()); } - @Test - public void testAssignIncreasingFreshIdNewIdentifier() { + @ParameterizedTest + @MethodSource("testTypes") + public void testAssignIncreasingFreshIdNewIdentifier(Type testType) { Schema schema = new Schema( Lists.newArrayList( - required(10, "a", Types.IntegerType.get()), - required(11, "A", Types.IntegerType.get())), + required(10, "a", Types.IntegerType.get()), required(11, "A", testType)), Sets.newHashSet(10)); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), - required(2, "A", Types.IntegerType.get()))); + required(1, "a", Types.IntegerType.get()), required(2, "A", testType))); final Schema actualSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(actualSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); assertThat(actualSchema.identifierFieldIds()) @@ -105,13 +108,14 @@ public void testAssignIncreasingFreshIdNewIdentifier() { .isEqualTo(Sets.newHashSet(sourceSchema.findField("a").fieldId())); } - @Test - public void testProject() { + @ParameterizedTest + @MethodSource("testTypes") + public void testProject(Type testType) { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", Types.IntegerType.get()), + required(11, "A", testType), required( 12, "someStruct", @@ -125,7 +129,7 @@ public void testProject() { required(16, "c", Types.IntegerType.get()), required(17, "C", Types.IntegerType.get()))))))); - Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", Types.IntegerType.get()))); + Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", testType))); Schema actualTop = TypeUtil.project(schema, Sets.newHashSet(11)); assertThat(actualTop.asStruct()).isEqualTo(expectedTop.asStruct()); @@ -145,7 +149,7 @@ public void testProject() { Schema expectedDepthTwo = new Schema( Lists.newArrayList( - required(11, "A", Types.IntegerType.get()), + required(11, "A", testType), required( 12, "someStruct", @@ -210,13 +214,14 @@ public void testProjectNaturallyEmpty() { assertThat(actualDepthThreeChildren.asStruct()).isEqualTo(expectedDepthThree.asStruct()); } - @Test - public void testProjectEmpty() { + @ParameterizedTest + @MethodSource("testTypes") + public void testProjectEmpty(Type testType) { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", Types.IntegerType.get()), + required(11, "A", testType), required( 12, "someStruct", @@ -248,13 +253,14 @@ public void testProjectEmpty() { assertThat(actualDepthTwo.asStruct()).isEqualTo(expectedDepthTwo.asStruct()); } - @Test - public void testSelect() { + @ParameterizedTest + @MethodSource("testTypes") + public void testSelect(Type testType) { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", Types.IntegerType.get()), + required(11, "A", testType), required( 12, "someStruct", @@ -268,7 +274,7 @@ public void testSelect() { required(16, "c", Types.IntegerType.get()), required(17, "C", Types.IntegerType.get()))))))); - Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", Types.IntegerType.get()))); + Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", testType))); Schema actualTop = TypeUtil.select(schema, Sets.newHashSet(11)); assertThat(actualTop.asStruct()).isEqualTo(expectedTop.asStruct()); @@ -296,7 +302,7 @@ public void testSelect() { Schema expectedDepthTwo = new Schema( Lists.newArrayList( - required(11, "A", Types.IntegerType.get()), + required(11, "A", testType), required( 12, "someStruct", @@ -310,13 +316,14 @@ public void testSelect() { assertThat(actualDepthTwo.asStruct()).isEqualTo(expectedDepthTwo.asStruct()); } - @Test - public void testProjectMap() { + @ParameterizedTest + @MethodSource("testTypes") + public void testProjectMap(Type testType) { // We can't partially project keys because it changes key equality Schema schema = new Schema( Lists.newArrayList( - required(10, "a", Types.IntegerType.get()), + required(10, "a", testType), required(11, "A", Types.IntegerType.get()), required( 12, @@ -348,15 +355,14 @@ public void testProjectMap() { .isInstanceOf(IllegalArgumentException.class) .hasMessageContaining("Cannot explicitly project List or Map types"); - Schema expectedTopLevel = - new Schema(Lists.newArrayList(required(10, "a", Types.IntegerType.get()))); + Schema expectedTopLevel = new Schema(Lists.newArrayList(required(10, "a", testType))); Schema actualTopLevel = TypeUtil.project(schema, Sets.newHashSet(10)); assertThat(actualTopLevel.asStruct()).isEqualTo(expectedTopLevel.asStruct()); Schema expectedDepthOne = new Schema( Lists.newArrayList( - required(10, "a", Types.IntegerType.get()), + required(10, "a", testType), required( 12, "map", @@ -375,7 +381,7 @@ public void testProjectMap() { Schema expectedDepthTwo = new Schema( Lists.newArrayList( - required(10, "a", Types.IntegerType.get()), + required(10, "a", testType), required( 12, "map", @@ -397,13 +403,14 @@ public void testProjectMap() { assertThat(actualDepthTwo.asStruct()).isEqualTo(expectedDepthTwo.asStruct()); } - @Test - public void testGetProjectedIds() { + @ParameterizedTest + @MethodSource("testTypes") + public void testGetProjectedIds(Type testType) { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", Types.IntegerType.get()), + required(11, "A", testType), required(35, "emptyStruct", Types.StructType.of()), required( 12, @@ -424,8 +431,9 @@ public void testGetProjectedIds() { assertThat(actualIds).isEqualTo(expectedIds); } - @Test - public void testProjectListNested() { + @ParameterizedTest + @MethodSource("testTypes") + public void testProjectListNested(Type testType) { Schema schema = new Schema( Lists.newArrayList( @@ -439,7 +447,7 @@ public void testProjectListNested() { Types.MapType.ofRequired( 15, 16, - IntegerType.get(), + testType, Types.StructType.of( required(17, "x", Types.IntegerType.get()), required(18, "y", Types.IntegerType.get())))))))); @@ -466,15 +474,15 @@ public void testProjectListNested() { 13, Types.ListType.ofRequired( 14, - Types.MapType.ofRequired( - 15, 16, IntegerType.get(), Types.StructType.of())))))); + Types.MapType.ofRequired(15, 16, testType, Types.StructType.of())))))); Schema actual = TypeUtil.project(schema, Sets.newHashSet(16)); assertThat(actual.asStruct()).isEqualTo(expected.asStruct()); } - @Test - public void testProjectMapNested() { + @ParameterizedTest + @MethodSource("testTypes") + public void testProjectMapNested(Type testType) { Schema schema = new Schema( Lists.newArrayList( @@ -488,7 +496,7 @@ public void testProjectMapNested() { Types.MapType.ofRequired( 15, 16, - Types.IntegerType.get(), + testType, Types.ListType.ofRequired( 17, Types.StructType.of( @@ -520,33 +528,34 @@ public void testProjectMapNested() { Types.MapType.ofRequired( 15, 16, - Types.IntegerType.get(), + testType, Types.ListType.ofRequired(17, Types.StructType.of())))))); Schema actual = TypeUtil.project(schema, Sets.newHashSet(17)); assertThat(actual.asStruct()).isEqualTo(expected.asStruct()); } - @Test - public void testReassignIdsIllegalArgumentException() { + @ParameterizedTest + @MethodSource("testTypes") + public void testReassignIdsIllegalArgumentException(Type testType) { Schema schema = - new Schema( - required(1, "a", Types.IntegerType.get()), required(2, "b", Types.IntegerType.get())); + new Schema(required(1, "a", Types.IntegerType.get()), required(2, "b", testType)); Schema sourceSchema = new Schema(required(1, "a", Types.IntegerType.get())); assertThatThrownBy(() -> TypeUtil.reassignIds(schema, sourceSchema)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Field b not found in source schema"); } - @Test - public void testValidateSchemaViaIndexByName() { + @ParameterizedTest + @MethodSource("testTypes") + public void testValidateSchemaViaIndexByName(Type testType) { Types.NestedField nestedType = Types.NestedField.required( 1, "a", Types.StructType.of( required(2, "b", Types.StructType.of(required(3, "c", Types.BooleanType.get()))), - required(4, "b.c", Types.BooleanType.get()))); + required(4, "b.c", testType))); assertThatThrownBy(() -> TypeUtil.indexByName(Types.StructType.of(nestedType))) .isInstanceOf(RuntimeException.class) @@ -576,8 +585,8 @@ public void testSelectNot() { required(3, "lat", Types.DoubleType.get()), required(4, "long", Types.DoubleType.get()))))); - Schema actualNoPrimitve = TypeUtil.selectNot(schema, Sets.newHashSet(1)); - assertThat(actualNoPrimitve.asStruct()).isEqualTo(expectedNoPrimitive.asStruct()); + Schema actualNoPrimitive = TypeUtil.selectNot(schema, Sets.newHashSet(1)); + assertThat(actualNoPrimitive.asStruct()).isEqualTo(expectedNoPrimitive.asStruct()); // Expected legacy behavior is to completely remove structs if their elements are removed Schema expectedNoStructElements = new Schema(required(1, "id", Types.LongType.get())); @@ -589,15 +598,16 @@ public void testSelectNot() { assertThat(actualNoStruct.asStruct()).isEqualTo(schema.asStruct()); } - @Test - public void testReassignOrRefreshIds() { + @ParameterizedTest + @MethodSource("testTypes") + public void testReassignOrRefreshIds(Type testType) { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), Types.NestedField.required("c") .withId(11) - .ofType(Types.IntegerType.get()) + .ofType(testType) .withInitialDefault(23) .withWriteDefault(34) .build(), @@ -625,24 +635,22 @@ public void testReassignOrRefreshIds() { assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } - @Test - public void testReassignOrRefreshIdsCaseInsensitive() { + @ParameterizedTest + @MethodSource("testTypes") + public void testReassignOrRefreshIdsCaseInsensitive(Type testType) { Schema schema = new Schema( Lists.newArrayList( - required(1, "FIELD1", Types.IntegerType.get()), - required(2, "FIELD2", Types.IntegerType.get()))); + required(1, "FIELD1", Types.IntegerType.get()), required(2, "FIELD2", testType))); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "field1", Types.IntegerType.get()), - required(2, "field2", Types.IntegerType.get()))); + required(1, "field1", Types.IntegerType.get()), required(2, "field2", testType))); final Schema actualSchema = TypeUtil.reassignOrRefreshIds(schema, sourceSchema, false); final Schema expectedSchema = new Schema( Lists.newArrayList( - required(1, "FIELD1", Types.IntegerType.get()), - required(2, "FIELD2", Types.IntegerType.get()))); + required(1, "FIELD1", Types.IntegerType.get()), required(2, "FIELD2", testType))); assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } } diff --git a/core/src/main/java/org/apache/iceberg/SchemaParser.java b/core/src/main/java/org/apache/iceberg/SchemaParser.java index 27e6ed048712..346aabc120ae 100644 --- a/core/src/main/java/org/apache/iceberg/SchemaParser.java +++ b/core/src/main/java/org/apache/iceberg/SchemaParser.java @@ -42,6 +42,7 @@ private SchemaParser() {} private static final String STRUCT = "struct"; private static final String LIST = "list"; private static final String MAP = "map"; + private static final String VARIANT = "variant"; private static final String FIELDS = "fields"; private static final String ELEMENT = "element"; private static final String KEY = "key"; @@ -145,6 +146,8 @@ static void toJson(Type.PrimitiveType primitive, JsonGenerator generator) throws static void toJson(Type type, JsonGenerator generator) throws IOException { if (type.isPrimitiveType()) { toJson(type.asPrimitiveType(), generator); + } else if (type.isVariantType()) { + generator.writeString(type.toString()); } else { Type.NestedType nested = type.asNestedType(); switch (type.typeId()) { @@ -179,6 +182,10 @@ public static String toJson(Schema schema, boolean pretty) { private static Type typeFromJson(JsonNode json) { if (json.isTextual()) { + if (VARIANT.equalsIgnoreCase(json.asText())) { + return Types.VariantType.get(); + } + return Types.fromPrimitiveString(json.asText()); } else if (json.isObject()) { JsonNode typeObj = json.get(TYPE); diff --git a/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java b/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java index c5c78dd1472a..aa31adae2692 100644 --- a/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java +++ b/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java @@ -56,6 +56,10 @@ class BuildAvroProjection extends AvroCustomOrderSchemaVisitor names, Iterable schemaIterable) { + if (current.isVariantType()) { + return variant(record); + } + Preconditions.checkArgument( current.isNestedType() && current.asNestedType().isStructType(), "Cannot project non-struct: %s", @@ -130,6 +134,17 @@ public Schema record(Schema record, List names, Iterable s return record; } + private Schema variant(Schema record) { + Preconditions.checkArgument( + current.isVariantType() + && record.getField("value") != null + && record.getField("metadata") != null, + "Expect variant type with value and metadata fields: %s", + current); + + return record; + } + @Override public Schema.Field field(Schema.Field field, Supplier fieldResult) { Types.StructType struct = current.asNestedType().asStructType(); diff --git a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java index 05ce4e618662..05f8afaba10d 100644 --- a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java +++ b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java @@ -49,6 +49,15 @@ abstract class TypeToSchema extends TypeUtil.SchemaVisitor { private static final Schema UUID_SCHEMA = LogicalTypes.uuid().addToSchema(Schema.createFixed("uuid_fixed", null, null, 16)); private static final Schema BINARY_SCHEMA = Schema.create(Schema.Type.BYTES); + private static final Schema VARIANT_SCHEMA = + Schema.createRecord( + "variant", + null, + null, + false, + List.of( + new Schema.Field("metadata", BINARY_SCHEMA), + new Schema.Field("value", BINARY_SCHEMA))); static { TIMESTAMP_SCHEMA.addProp(AvroSchemaUtil.ADJUST_TO_UTC_PROP, false); @@ -187,6 +196,11 @@ public Schema map(Types.MapType map, Schema keySchema, Schema valueSchema) { return mapSchema; } + @Override + public Schema variant() { + return VARIANT_SCHEMA; + } + @Override public Schema primitive(Type.PrimitiveType primitive) { Schema primitiveSchema; diff --git a/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java b/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java index 741184d612f1..cc6648533ffd 100644 --- a/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java +++ b/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java @@ -30,6 +30,7 @@ import java.util.Map; import java.util.Set; import java.util.stream.IntStream; +import java.util.stream.Stream; import org.apache.iceberg.catalog.Namespace; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; @@ -42,6 +43,9 @@ import org.apache.iceberg.view.ViewVersion; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; public class TestMetadataUpdateParser { @@ -52,6 +56,15 @@ public class TestMetadataUpdateParser { Types.NestedField.required(1, "id", Types.IntegerType.get()), Types.NestedField.optional(2, "data", Types.StringType.get())); + private static final Schema ID_VARIANTDATA_SCHEMA = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "data", Types.VariantType.get())); + + private static Stream testSchemas() { + return Stream.of(Arguments.of(ID_DATA_SCHEMA), Arguments.of(ID_VARIANTDATA_SCHEMA)); + } + @Test public void testMetadataUpdateWithoutActionCannotDeserialize() { List invalidJson = @@ -108,19 +121,19 @@ public void testUpgradeFormatVersionFromJson() { } /** AddSchema * */ - @Test - public void testAddSchemaFromJson() { + @ParameterizedTest + @MethodSource("testSchemas") + public void testAddSchemaFromJson(Schema schema) { String action = MetadataUpdateParser.ADD_SCHEMA; - Schema schema = ID_DATA_SCHEMA; String json = String.format("{\"action\":\"add-schema\",\"schema\":%s}", SchemaParser.toJson(schema)); MetadataUpdate actualUpdate = new MetadataUpdate.AddSchema(schema); assertEquals(action, actualUpdate, MetadataUpdateParser.fromJson(json)); } - @Test - public void testAddSchemaToJson() { - Schema schema = ID_DATA_SCHEMA; + @ParameterizedTest + @MethodSource("testSchemas") + public void testAddSchemaToJson(Schema schema) { int lastColumnId = schema.highestFieldId(); String expected = String.format( diff --git a/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java b/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java index b71280daf8db..cff99443db1b 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java @@ -18,11 +18,13 @@ */ package org.apache.iceberg.avro; +import static org.apache.iceberg.types.Types.NestedField.required; import static org.assertj.core.api.Assertions.assertThat; import java.util.Collections; import org.apache.avro.SchemaBuilder; import org.apache.iceberg.Schema; +import org.apache.iceberg.types.Types; import org.junit.jupiter.api.Test; public class TestAvroSchemaProjection { @@ -150,4 +152,16 @@ public void projectWithMapSchemaChanged() { .as("Result of buildAvroProjection is missing some IDs") .isFalse(); } + + @Test + public void testVariantConversion() { + Schema schema = new Schema(required(1, "variantCol", Types.VariantType.get())); + org.apache.avro.Schema avroSchema = AvroSchemaUtil.convert(schema.asStruct()); + + org.apache.avro.Schema variantSchema = avroSchema.getField("variantCol").schema(); + assertThat(variantSchema.getType()).isEqualTo(org.apache.avro.Schema.Type.RECORD); + assertThat(variantSchema.getFields().size()).isEqualTo(2); + assertThat(variantSchema.getField("metadata")).isNotNull(); + assertThat(variantSchema.getField("value")).isNotNull(); + } } diff --git a/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java b/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java index eaea4394dbfd..e1b795d93e71 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java @@ -22,6 +22,7 @@ import static org.assertj.core.api.Assertions.assertThat; import java.util.Collections; +import java.util.List; import java.util.function.Supplier; import org.apache.avro.SchemaBuilder; import org.apache.iceberg.types.Type; @@ -401,4 +402,32 @@ public void projectMapWithLessFieldInValueSchema() { .as("Unexpected value ID discovered on the projected map schema") .isEqualTo(1); } + + @Test + public void projectVariantSchemaUnchanged() { + final Type icebergType = Types.VariantType.get(); + + final org.apache.avro.Schema expected = + SchemaBuilder.record("variant") + .namespace("unit.test") + .fields() + .name("metadata") + .prop(AvroSchemaUtil.FIELD_ID_PROP, "1") + .type() + .bytesType() + .noDefault() + .name("value") + .prop(AvroSchemaUtil.FIELD_ID_PROP, "2") + .type() + .bytesType() + .noDefault() + .endRecord(); + + final BuildAvroProjection testSubject = + new BuildAvroProjection(icebergType, Collections.emptyMap()); + final org.apache.avro.Schema actual = testSubject.record(expected, List.of(), null); + assertThat(actual) + .as("Variant projection produced undesired variant schema") + .isEqualTo(expected); + } } From 9a07342099c1214cba74b18fa7c43de5de6b3888 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Fri, 20 Dec 2024 20:17:35 -0800 Subject: [PATCH 02/11] Address comments --- .../main/java/org/apache/iceberg/types/Types.java | 8 ++++++++ .../org/apache/iceberg/types/TestTypeUtil.java | 9 ++++----- .../main/java/org/apache/iceberg/SchemaParser.java | 7 +------ .../iceberg/avro/TestAvroSchemaProjection.java | 14 -------------- .../apache/iceberg/avro/TestSchemaConversions.java | 13 +++++++++++++ 5 files changed, 26 insertions(+), 25 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/types/Types.java b/api/src/main/java/org/apache/iceberg/types/Types.java index 2c7c3ed81846..fc091065db52 100644 --- a/api/src/main/java/org/apache/iceberg/types/Types.java +++ b/api/src/main/java/org/apache/iceberg/types/Types.java @@ -62,6 +62,14 @@ private Types() {} private static final Pattern DECIMAL = Pattern.compile("decimal\\(\\s*(\\d+)\\s*,\\s*(\\d+)\\s*\\)"); + public static Type typeFromTypeString(String typeString) { + if (VariantType.get().toString().equalsIgnoreCase(typeString)) { + return Types.VariantType.get(); + } + + return Types.fromPrimitiveString(typeString); + } + public static PrimitiveType fromPrimitiveString(String typeString) { String lowerTypeString = typeString.toLowerCase(Locale.ROOT); if (TYPES.containsKey(lowerTypeString)) { diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java index e6740d8efc8b..7d90e7cf6e24 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java @@ -607,17 +607,16 @@ public void testReassignOrRefreshIds(Type testType) { required(10, "a", Types.IntegerType.get()), Types.NestedField.required("c") .withId(11) - .ofType(testType) + .ofType(Types.IntegerType.get()) .withInitialDefault(23) .withWriteDefault(34) .build(), - required(12, "B", Types.IntegerType.get())), + required(12, "B", testType)), Sets.newHashSet(10)); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), - required(15, "B", Types.IntegerType.get()))); + required(1, "a", Types.IntegerType.get()), required(15, "B", testType))); Schema actualSchema = TypeUtil.reassignOrRefreshIds(schema, sourceSchema); Schema expectedSchema = @@ -630,7 +629,7 @@ public void testReassignOrRefreshIds(Type testType) { .withInitialDefault(23) .withWriteDefault(34) .build(), - required(15, "B", Types.IntegerType.get()))); + required(15, "B", testType))); assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } diff --git a/core/src/main/java/org/apache/iceberg/SchemaParser.java b/core/src/main/java/org/apache/iceberg/SchemaParser.java index 346aabc120ae..22ad0e9af708 100644 --- a/core/src/main/java/org/apache/iceberg/SchemaParser.java +++ b/core/src/main/java/org/apache/iceberg/SchemaParser.java @@ -42,7 +42,6 @@ private SchemaParser() {} private static final String STRUCT = "struct"; private static final String LIST = "list"; private static final String MAP = "map"; - private static final String VARIANT = "variant"; private static final String FIELDS = "fields"; private static final String ELEMENT = "element"; private static final String KEY = "key"; @@ -182,11 +181,7 @@ public static String toJson(Schema schema, boolean pretty) { private static Type typeFromJson(JsonNode json) { if (json.isTextual()) { - if (VARIANT.equalsIgnoreCase(json.asText())) { - return Types.VariantType.get(); - } - - return Types.fromPrimitiveString(json.asText()); + return Types.typeFromTypeString(json.asText()); } else if (json.isObject()) { JsonNode typeObj = json.get(TYPE); if (typeObj != null) { diff --git a/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java b/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java index cff99443db1b..b71280daf8db 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestAvroSchemaProjection.java @@ -18,13 +18,11 @@ */ package org.apache.iceberg.avro; -import static org.apache.iceberg.types.Types.NestedField.required; import static org.assertj.core.api.Assertions.assertThat; import java.util.Collections; import org.apache.avro.SchemaBuilder; import org.apache.iceberg.Schema; -import org.apache.iceberg.types.Types; import org.junit.jupiter.api.Test; public class TestAvroSchemaProjection { @@ -152,16 +150,4 @@ public void projectWithMapSchemaChanged() { .as("Result of buildAvroProjection is missing some IDs") .isFalse(); } - - @Test - public void testVariantConversion() { - Schema schema = new Schema(required(1, "variantCol", Types.VariantType.get())); - org.apache.avro.Schema avroSchema = AvroSchemaUtil.convert(schema.asStruct()); - - org.apache.avro.Schema variantSchema = avroSchema.getField("variantCol").schema(); - assertThat(variantSchema.getType()).isEqualTo(org.apache.avro.Schema.Type.RECORD); - assertThat(variantSchema.getFields().size()).isEqualTo(2); - assertThat(variantSchema.getField("metadata")).isNotNull(); - assertThat(variantSchema.getField("value")).isNotNull(); - } } diff --git a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java index d9dc49d17257..ac86a09e7d04 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java @@ -370,4 +370,17 @@ public void testFieldDocsArePreserved() { Lists.newArrayList(Iterables.transform(origSchema.columns(), Types.NestedField::doc)); assertThat(fieldDocs).isEqualTo(origFieldDocs); } + + @Test + public void testVariantConversion() { + org.apache.iceberg.Schema schema = + new org.apache.iceberg.Schema(required(1, "variantCol", Types.VariantType.get())); + org.apache.avro.Schema avroSchema = AvroSchemaUtil.convert(schema.asStruct()); + + org.apache.avro.Schema variantSchema = avroSchema.getField("variantCol").schema(); + assertThat(variantSchema.getType()).isEqualTo(org.apache.avro.Schema.Type.RECORD); + assertThat(variantSchema.getFields().size()).isEqualTo(2); + assertThat(variantSchema.getField("metadata")).isNotNull(); + assertThat(variantSchema.getField("value")).isNotNull(); + } } From 31b0ee3f262a351964b2b5cb74436162a161097e Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Mon, 6 Jan 2025 10:08:06 -0800 Subject: [PATCH 03/11] Revert unrelated change --- .../main/java/org/apache/iceberg/transforms/Identity.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/transforms/Identity.java b/api/src/main/java/org/apache/iceberg/transforms/Identity.java index f0ca938a402c..099a99cc3cf4 100644 --- a/api/src/main/java/org/apache/iceberg/transforms/Identity.java +++ b/api/src/main/java/org/apache/iceberg/transforms/Identity.java @@ -19,17 +19,14 @@ package org.apache.iceberg.transforms; import java.io.ObjectStreamException; -import java.util.Set; import org.apache.iceberg.expressions.BoundPredicate; import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.expressions.UnboundPredicate; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.types.Type; -import org.apache.iceberg.types.Types; import org.apache.iceberg.util.SerializableFunction; class Identity implements Transform { - private static final Set UNSUPPORTED_TYPES = Set.of(Types.VariantType.get()); private static final Identity INSTANCE = new Identity<>(); private final Type type; @@ -42,7 +39,7 @@ class Identity implements Transform { @Deprecated public static Identity get(Type type) { Preconditions.checkArgument( - !UNSUPPORTED_TYPES.contains(type), "Unsupported type for identity: %s", type); + type.typeId() != Type.TypeID.VARIANT, "Unsupported type for identity: %s", type); return new Identity<>(type); } From 82f654cd4cef18155af89ef8f0e72bfc11692029 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Fri, 24 Jan 2025 22:40:29 -0800 Subject: [PATCH 04/11] Update tests --- .../java/org/apache/iceberg/Accessors.java | 5 + .../apache/iceberg/types/GetProjectedIds.java | 5 + .../org/apache/iceberg/types/IndexById.java | 5 + .../apache/iceberg/types/PruneColumns.java | 5 + .../org/apache/iceberg/types/TypeUtil.java | 2 +- .../java/org/apache/iceberg/types/Types.java | 15 +- .../apache/iceberg/types/TestTypeUtil.java | 202 ++++++++++-------- .../org/apache/iceberg/types/TestTypes.java | 3 + .../java/org/apache/iceberg/SchemaParser.java | 2 +- .../iceberg/avro/BuildAvroProjection.java | 2 + .../org/apache/iceberg/avro/TypeToSchema.java | 18 +- .../iceberg/TestMetadataUpdateParser.java | 25 +-- .../org/apache/iceberg/TestSchemaParser.java | 10 + .../iceberg/avro/TestBuildAvroProjection.java | 3 +- .../iceberg/data/avro/TestGenericData.java | 35 +++ 15 files changed, 206 insertions(+), 131 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/Accessors.java b/api/src/main/java/org/apache/iceberg/Accessors.java index 08233624f244..b372f6a87233 100644 --- a/api/src/main/java/org/apache/iceberg/Accessors.java +++ b/api/src/main/java/org/apache/iceberg/Accessors.java @@ -232,6 +232,11 @@ public Map> struct( return accessors; } + @Override + public Map> variant() { + return null; + } + @Override public Map> field( Types.NestedField field, Map> fieldResult) { diff --git a/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java b/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java index 4692a8c400f8..814bb72f201c 100644 --- a/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java +++ b/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java @@ -74,4 +74,9 @@ public Set map(Types.MapType map, Set keyResult, Set } return fieldIds; } + + @Override + public Set variant() { + return null; + } } diff --git a/api/src/main/java/org/apache/iceberg/types/IndexById.java b/api/src/main/java/org/apache/iceberg/types/IndexById.java index 40280c5ed9dd..7c36a3ec241a 100644 --- a/api/src/main/java/org/apache/iceberg/types/IndexById.java +++ b/api/src/main/java/org/apache/iceberg/types/IndexById.java @@ -64,4 +64,9 @@ public Map map( } return null; } + + @Override + public Map variant() { + return null; + } } diff --git a/api/src/main/java/org/apache/iceberg/types/PruneColumns.java b/api/src/main/java/org/apache/iceberg/types/PruneColumns.java index daf2e6bbc0ca..d9230ce6e02b 100644 --- a/api/src/main/java/org/apache/iceberg/types/PruneColumns.java +++ b/api/src/main/java/org/apache/iceberg/types/PruneColumns.java @@ -159,6 +159,11 @@ public Type map(Types.MapType map, Type ignored, Type valueResult) { return null; } + @Override + public Type variant() { + return null; + } + @Override public Type primitive(Type.PrimitiveType primitive) { return null; diff --git a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java index 19ef2d155f26..c5c97e723834 100644 --- a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java +++ b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java @@ -617,7 +617,7 @@ public T map(Types.MapType map, T keyResult, T valueResult) { } public T variant() { - return null; + throw new UnsupportedOperationException("Unsupported type: variant"); } public T primitive(Type.PrimitiveType primitive) { diff --git a/api/src/main/java/org/apache/iceberg/types/Types.java b/api/src/main/java/org/apache/iceberg/types/Types.java index fc091065db52..0bbecfff5d1d 100644 --- a/api/src/main/java/org/apache/iceberg/types/Types.java +++ b/api/src/main/java/org/apache/iceberg/types/Types.java @@ -39,8 +39,8 @@ public class Types { private Types() {} - private static final ImmutableMap TYPES = - ImmutableMap.builder() + private static final ImmutableMap TYPES = + ImmutableMap.builder() .put(BooleanType.get().toString(), BooleanType.get()) .put(IntegerType.get().toString(), IntegerType.get()) .put(LongType.get().toString(), LongType.get()) @@ -56,21 +56,14 @@ private Types() {} .put(UUIDType.get().toString(), UUIDType.get()) .put(BinaryType.get().toString(), BinaryType.get()) .put(UnknownType.get().toString(), UnknownType.get()) + .put(VariantType.get().toString(), VariantType.get()) .buildOrThrow(); private static final Pattern FIXED = Pattern.compile("fixed\\[\\s*(\\d+)\\s*\\]"); private static final Pattern DECIMAL = Pattern.compile("decimal\\(\\s*(\\d+)\\s*,\\s*(\\d+)\\s*\\)"); - public static Type typeFromTypeString(String typeString) { - if (VariantType.get().toString().equalsIgnoreCase(typeString)) { - return Types.VariantType.get(); - } - - return Types.fromPrimitiveString(typeString); - } - - public static PrimitiveType fromPrimitiveString(String typeString) { + public static Type fromPrimitiveString(String typeString) { String lowerTypeString = typeString.toLowerCase(Locale.ROOT); if (TYPES.containsKey(lowerTypeString)) { return TYPES.get(lowerTypeString); diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java index 7d90e7cf6e24..272f58bda5ba 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java @@ -23,44 +23,41 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import java.util.Map; import java.util.Set; -import java.util.stream.Stream; +import java.util.concurrent.atomic.AtomicInteger; import org.apache.iceberg.Schema; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.Types.IntegerType; import org.junit.jupiter.api.Test; -import org.junit.jupiter.params.ParameterizedTest; -import org.junit.jupiter.params.provider.Arguments; -import org.junit.jupiter.params.provider.MethodSource; public class TestTypeUtil { - private static Stream testTypes() { - return Stream.of(Arguments.of(Types.IntegerType.get()), Arguments.of(Types.VariantType.get())); - } - - @ParameterizedTest - @MethodSource("testTypes") - public void testReassignIdsDuplicateColumns(Type testType) { + @Test + public void testReassignIdsDuplicateColumns() { Schema schema = - new Schema(required(0, "a", testType), required(1, "A", Types.IntegerType.get())); + new Schema( + required(0, "a", Types.IntegerType.get()), required(1, "A", Types.IntegerType.get())); Schema sourceSchema = - new Schema(required(1, "a", testType), required(2, "A", Types.IntegerType.get())); + new Schema( + required(1, "a", Types.IntegerType.get()), required(2, "A", Types.IntegerType.get())); final Schema actualSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(actualSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testReassignIdsWithIdentifier(Type testType) { + @Test + public void testReassignIdsWithIdentifier() { Schema schema = new Schema( Lists.newArrayList( - required(0, "a", Types.IntegerType.get()), required(1, "A", testType)), + required(0, "a", Types.IntegerType.get()), + required(1, "A", Types.IntegerType.get())), Sets.newHashSet(0)); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), required(2, "A", testType)), + required(1, "a", Types.IntegerType.get()), + required(2, "A", Types.IntegerType.get())), Sets.newHashSet(1)); final Schema actualSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(actualSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); @@ -69,18 +66,19 @@ public void testReassignIdsWithIdentifier(Type testType) { .isEqualTo(sourceSchema.identifierFieldIds()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testAssignIncreasingFreshIdWithIdentifier(Type testType) { + @Test + public void testAssignIncreasingFreshIdWithIdentifier() { Schema schema = new Schema( Lists.newArrayList( - required(10, "a", Types.IntegerType.get()), required(11, "A", testType)), + required(10, "a", Types.IntegerType.get()), + required(11, "A", Types.IntegerType.get())), Sets.newHashSet(10)); Schema expectedSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), required(2, "A", testType)), + required(1, "a", Types.IntegerType.get()), + required(2, "A", Types.IntegerType.get())), Sets.newHashSet(1)); final Schema actualSchema = TypeUtil.assignIncreasingFreshIds(schema); assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); @@ -89,18 +87,19 @@ public void testAssignIncreasingFreshIdWithIdentifier(Type testType) { .isEqualTo(expectedSchema.identifierFieldIds()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testAssignIncreasingFreshIdNewIdentifier(Type testType) { + @Test + public void testAssignIncreasingFreshIdNewIdentifier() { Schema schema = new Schema( Lists.newArrayList( - required(10, "a", Types.IntegerType.get()), required(11, "A", testType)), + required(10, "a", Types.IntegerType.get()), + required(11, "A", Types.IntegerType.get())), Sets.newHashSet(10)); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), required(2, "A", testType))); + required(1, "a", Types.IntegerType.get()), + required(2, "A", Types.IntegerType.get()))); final Schema actualSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(actualSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); assertThat(actualSchema.identifierFieldIds()) @@ -108,14 +107,13 @@ public void testAssignIncreasingFreshIdNewIdentifier(Type testType) { .isEqualTo(Sets.newHashSet(sourceSchema.findField("a").fieldId())); } - @ParameterizedTest - @MethodSource("testTypes") - public void testProject(Type testType) { + @Test + public void testProject() { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", testType), + required(11, "A", Types.IntegerType.get()), required( 12, "someStruct", @@ -129,7 +127,7 @@ public void testProject(Type testType) { required(16, "c", Types.IntegerType.get()), required(17, "C", Types.IntegerType.get()))))))); - Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", testType))); + Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", Types.IntegerType.get()))); Schema actualTop = TypeUtil.project(schema, Sets.newHashSet(11)); assertThat(actualTop.asStruct()).isEqualTo(expectedTop.asStruct()); @@ -149,7 +147,7 @@ public void testProject(Type testType) { Schema expectedDepthTwo = new Schema( Lists.newArrayList( - required(11, "A", testType), + required(11, "A", Types.IntegerType.get()), required( 12, "someStruct", @@ -214,14 +212,13 @@ public void testProjectNaturallyEmpty() { assertThat(actualDepthThreeChildren.asStruct()).isEqualTo(expectedDepthThree.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testProjectEmpty(Type testType) { + @Test + public void testProjectEmpty() { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", testType), + required(11, "A", Types.IntegerType.get()), required( 12, "someStruct", @@ -253,14 +250,13 @@ public void testProjectEmpty(Type testType) { assertThat(actualDepthTwo.asStruct()).isEqualTo(expectedDepthTwo.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testSelect(Type testType) { + @Test + public void testSelect() { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", testType), + required(11, "A", Types.IntegerType.get()), required( 12, "someStruct", @@ -274,7 +270,7 @@ public void testSelect(Type testType) { required(16, "c", Types.IntegerType.get()), required(17, "C", Types.IntegerType.get()))))))); - Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", testType))); + Schema expectedTop = new Schema(Lists.newArrayList(required(11, "A", Types.IntegerType.get()))); Schema actualTop = TypeUtil.select(schema, Sets.newHashSet(11)); assertThat(actualTop.asStruct()).isEqualTo(expectedTop.asStruct()); @@ -302,7 +298,7 @@ public void testSelect(Type testType) { Schema expectedDepthTwo = new Schema( Lists.newArrayList( - required(11, "A", testType), + required(11, "A", Types.IntegerType.get()), required( 12, "someStruct", @@ -316,14 +312,13 @@ public void testSelect(Type testType) { assertThat(actualDepthTwo.asStruct()).isEqualTo(expectedDepthTwo.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testProjectMap(Type testType) { + @Test + public void testProjectMap() { // We can't partially project keys because it changes key equality Schema schema = new Schema( Lists.newArrayList( - required(10, "a", testType), + required(10, "a", Types.IntegerType.get()), required(11, "A", Types.IntegerType.get()), required( 12, @@ -355,14 +350,15 @@ public void testProjectMap(Type testType) { .isInstanceOf(IllegalArgumentException.class) .hasMessageContaining("Cannot explicitly project List or Map types"); - Schema expectedTopLevel = new Schema(Lists.newArrayList(required(10, "a", testType))); + Schema expectedTopLevel = + new Schema(Lists.newArrayList(required(10, "a", Types.IntegerType.get()))); Schema actualTopLevel = TypeUtil.project(schema, Sets.newHashSet(10)); assertThat(actualTopLevel.asStruct()).isEqualTo(expectedTopLevel.asStruct()); Schema expectedDepthOne = new Schema( Lists.newArrayList( - required(10, "a", testType), + required(10, "a", Types.IntegerType.get()), required( 12, "map", @@ -381,7 +377,7 @@ public void testProjectMap(Type testType) { Schema expectedDepthTwo = new Schema( Lists.newArrayList( - required(10, "a", testType), + required(10, "a", Types.IntegerType.get()), required( 12, "map", @@ -403,14 +399,13 @@ public void testProjectMap(Type testType) { assertThat(actualDepthTwo.asStruct()).isEqualTo(expectedDepthTwo.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testGetProjectedIds(Type testType) { + @Test + public void testGetProjectedIds() { Schema schema = new Schema( Lists.newArrayList( required(10, "a", Types.IntegerType.get()), - required(11, "A", testType), + required(11, "A", Types.IntegerType.get()), required(35, "emptyStruct", Types.StructType.of()), required( 12, @@ -431,9 +426,8 @@ public void testGetProjectedIds(Type testType) { assertThat(actualIds).isEqualTo(expectedIds); } - @ParameterizedTest - @MethodSource("testTypes") - public void testProjectListNested(Type testType) { + @Test + public void testProjectListNested() { Schema schema = new Schema( Lists.newArrayList( @@ -447,7 +441,7 @@ public void testProjectListNested(Type testType) { Types.MapType.ofRequired( 15, 16, - testType, + IntegerType.get(), Types.StructType.of( required(17, "x", Types.IntegerType.get()), required(18, "y", Types.IntegerType.get())))))))); @@ -474,15 +468,15 @@ public void testProjectListNested(Type testType) { 13, Types.ListType.ofRequired( 14, - Types.MapType.ofRequired(15, 16, testType, Types.StructType.of())))))); + Types.MapType.ofRequired( + 15, 16, IntegerType.get(), Types.StructType.of())))))); Schema actual = TypeUtil.project(schema, Sets.newHashSet(16)); assertThat(actual.asStruct()).isEqualTo(expected.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testProjectMapNested(Type testType) { + @Test + public void testProjectMapNested() { Schema schema = new Schema( Lists.newArrayList( @@ -496,7 +490,7 @@ public void testProjectMapNested(Type testType) { Types.MapType.ofRequired( 15, 16, - testType, + Types.IntegerType.get(), Types.ListType.ofRequired( 17, Types.StructType.of( @@ -528,34 +522,33 @@ public void testProjectMapNested(Type testType) { Types.MapType.ofRequired( 15, 16, - testType, + Types.IntegerType.get(), Types.ListType.ofRequired(17, Types.StructType.of())))))); Schema actual = TypeUtil.project(schema, Sets.newHashSet(17)); assertThat(actual.asStruct()).isEqualTo(expected.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testReassignIdsIllegalArgumentException(Type testType) { + @Test + public void testReassignIdsIllegalArgumentException() { Schema schema = - new Schema(required(1, "a", Types.IntegerType.get()), required(2, "b", testType)); + new Schema( + required(1, "a", Types.IntegerType.get()), required(2, "b", Types.IntegerType.get())); Schema sourceSchema = new Schema(required(1, "a", Types.IntegerType.get())); assertThatThrownBy(() -> TypeUtil.reassignIds(schema, sourceSchema)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Field b not found in source schema"); } - @ParameterizedTest - @MethodSource("testTypes") - public void testValidateSchemaViaIndexByName(Type testType) { + @Test + public void testValidateSchemaViaIndexByName() { Types.NestedField nestedType = Types.NestedField.required( 1, "a", Types.StructType.of( required(2, "b", Types.StructType.of(required(3, "c", Types.BooleanType.get()))), - required(4, "b.c", testType))); + required(4, "b.c", Types.BooleanType.get()))); assertThatThrownBy(() -> TypeUtil.indexByName(Types.StructType.of(nestedType))) .isInstanceOf(RuntimeException.class) @@ -585,8 +578,8 @@ public void testSelectNot() { required(3, "lat", Types.DoubleType.get()), required(4, "long", Types.DoubleType.get()))))); - Schema actualNoPrimitive = TypeUtil.selectNot(schema, Sets.newHashSet(1)); - assertThat(actualNoPrimitive.asStruct()).isEqualTo(expectedNoPrimitive.asStruct()); + Schema actualNoPrimitve = TypeUtil.selectNot(schema, Sets.newHashSet(1)); + assertThat(actualNoPrimitve.asStruct()).isEqualTo(expectedNoPrimitive.asStruct()); // Expected legacy behavior is to completely remove structs if their elements are removed Schema expectedNoStructElements = new Schema(required(1, "id", Types.LongType.get())); @@ -598,9 +591,8 @@ public void testSelectNot() { assertThat(actualNoStruct.asStruct()).isEqualTo(schema.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testReassignOrRefreshIds(Type testType) { + @Test + public void testReassignOrRefreshIds() { Schema schema = new Schema( Lists.newArrayList( @@ -611,12 +603,13 @@ public void testReassignOrRefreshIds(Type testType) { .withInitialDefault(23) .withWriteDefault(34) .build(), - required(12, "B", testType)), + required(12, "B", Types.IntegerType.get())), Sets.newHashSet(10)); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "a", Types.IntegerType.get()), required(15, "B", testType))); + required(1, "a", Types.IntegerType.get()), + required(15, "B", Types.IntegerType.get()))); Schema actualSchema = TypeUtil.reassignOrRefreshIds(schema, sourceSchema); Schema expectedSchema = @@ -629,27 +622,62 @@ public void testReassignOrRefreshIds(Type testType) { .withInitialDefault(23) .withWriteDefault(34) .build(), - required(15, "B", testType))); + required(15, "B", Types.IntegerType.get()))); assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } - @ParameterizedTest - @MethodSource("testTypes") - public void testReassignOrRefreshIdsCaseInsensitive(Type testType) { + @Test + public void testReassignOrRefreshIdsCaseInsensitive() { Schema schema = new Schema( Lists.newArrayList( - required(1, "FIELD1", Types.IntegerType.get()), required(2, "FIELD2", testType))); + required(1, "FIELD1", Types.IntegerType.get()), + required(2, "FIELD2", Types.IntegerType.get()))); Schema sourceSchema = new Schema( Lists.newArrayList( - required(1, "field1", Types.IntegerType.get()), required(2, "field2", testType))); + required(1, "field1", Types.IntegerType.get()), + required(2, "field2", Types.IntegerType.get()))); final Schema actualSchema = TypeUtil.reassignOrRefreshIds(schema, sourceSchema, false); final Schema expectedSchema = new Schema( Lists.newArrayList( - required(1, "FIELD1", Types.IntegerType.get()), required(2, "FIELD2", testType))); + required(1, "FIELD1", Types.IntegerType.get()), + required(2, "FIELD2", Types.IntegerType.get()))); assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } + + @Test + public void testVariantType() { + Schema schema = + new Schema( + required(0, "v", Types.VariantType.get()), required(1, "A", Types.IntegerType.get())); + Schema sourceSchema = + new Schema( + required(1, "v", Types.VariantType.get()), required(2, "A", Types.IntegerType.get())); + + Schema assignedSchema = + TypeUtil.assignFreshIds(sourceSchema, new AtomicInteger(10)::incrementAndGet); + Schema expectedSchema = + new Schema( + required(11, "v", Types.VariantType.get()), required(12, "A", Types.IntegerType.get())); + assertThat(assignedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); + + final Schema reassignedSchema = TypeUtil.reassignIds(schema, sourceSchema); + assertThat(reassignedSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); + + expectedSchema = new Schema(Lists.newArrayList(required(1, "v", Types.VariantType.get()))); + Schema projectedSchema = TypeUtil.project(sourceSchema, Sets.newHashSet(1)); + assertThat(projectedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); + + Set projectedIds = TypeUtil.getProjectedIds(sourceSchema); + assertThat(Set.of(1, 2)).isEqualTo(projectedIds); + + Map indexByIds = TypeUtil.indexById(sourceSchema.asStruct()); + assertThat(indexByIds.get(1).type().isVariantType()).isTrue(); + + Map indexNameByIds = TypeUtil.indexNameById(sourceSchema.asStruct()); + assertThat(indexNameByIds.get(1)).isEqualTo("v"); + } } diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypes.java b/api/src/test/java/org/apache/iceberg/types/TestTypes.java index cbc37291375f..8063fd688a53 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypes.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypes.java @@ -43,6 +43,9 @@ public void fromPrimitiveString() { assertThat(Types.fromPrimitiveString("Decimal(2,3)")).isEqualTo(Types.DecimalType.of(2, 3)); + assertThat(Types.fromPrimitiveString("variant")).isSameAs(Types.VariantType.get()); + assertThat(Types.fromPrimitiveString("Variant")).isSameAs(Types.VariantType.get()); + assertThatExceptionOfType(IllegalArgumentException.class) .isThrownBy(() -> Types.fromPrimitiveString("abcdefghij")) .withMessage("Cannot parse type string to primitive: abcdefghij"); diff --git a/core/src/main/java/org/apache/iceberg/SchemaParser.java b/core/src/main/java/org/apache/iceberg/SchemaParser.java index 22ad0e9af708..610140ba908c 100644 --- a/core/src/main/java/org/apache/iceberg/SchemaParser.java +++ b/core/src/main/java/org/apache/iceberg/SchemaParser.java @@ -181,7 +181,7 @@ public static String toJson(Schema schema, boolean pretty) { private static Type typeFromJson(JsonNode json) { if (json.isTextual()) { - return Types.typeFromTypeString(json.asText()); + return Types.fromPrimitiveString(json.asText()); } else if (json.isObject()) { JsonNode typeObj = json.get(TYPE); if (typeObj != null) { diff --git a/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java b/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java index aa31adae2692..e7d099bd2683 100644 --- a/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java +++ b/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java @@ -56,6 +56,8 @@ class BuildAvroProjection extends AvroCustomOrderSchemaVisitor names, Iterable schemaIterable) { + // TODO: When the Variant logical type is introduced in Avro, handle the following in a separate + // visitor method if (current.isVariantType()) { return variant(record); } diff --git a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java index 05f8afaba10d..c009eed76e7c 100644 --- a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java +++ b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java @@ -49,15 +49,6 @@ abstract class TypeToSchema extends TypeUtil.SchemaVisitor { private static final Schema UUID_SCHEMA = LogicalTypes.uuid().addToSchema(Schema.createFixed("uuid_fixed", null, null, 16)); private static final Schema BINARY_SCHEMA = Schema.create(Schema.Type.BYTES); - private static final Schema VARIANT_SCHEMA = - Schema.createRecord( - "variant", - null, - null, - false, - List.of( - new Schema.Field("metadata", BINARY_SCHEMA), - new Schema.Field("value", BINARY_SCHEMA))); static { TIMESTAMP_SCHEMA.addProp(AvroSchemaUtil.ADJUST_TO_UTC_PROP, false); @@ -198,7 +189,14 @@ public Schema map(Types.MapType map, Schema keySchema, Schema valueSchema) { @Override public Schema variant() { - return VARIANT_SCHEMA; + String recordName = "r" + fieldIds.peek(); + return Schema.createRecord( + recordName, + null, + null, + false, + List.of( + new Schema.Field("metadata", BINARY_SCHEMA), new Schema.Field("value", BINARY_SCHEMA))); } @Override diff --git a/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java b/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java index cc6648533ffd..741184d612f1 100644 --- a/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java +++ b/core/src/test/java/org/apache/iceberg/TestMetadataUpdateParser.java @@ -30,7 +30,6 @@ import java.util.Map; import java.util.Set; import java.util.stream.IntStream; -import java.util.stream.Stream; import org.apache.iceberg.catalog.Namespace; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; @@ -43,9 +42,6 @@ import org.apache.iceberg.view.ViewVersion; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; -import org.junit.jupiter.params.ParameterizedTest; -import org.junit.jupiter.params.provider.Arguments; -import org.junit.jupiter.params.provider.MethodSource; public class TestMetadataUpdateParser { @@ -56,15 +52,6 @@ public class TestMetadataUpdateParser { Types.NestedField.required(1, "id", Types.IntegerType.get()), Types.NestedField.optional(2, "data", Types.StringType.get())); - private static final Schema ID_VARIANTDATA_SCHEMA = - new Schema( - Types.NestedField.required(1, "id", Types.IntegerType.get()), - Types.NestedField.optional(2, "data", Types.VariantType.get())); - - private static Stream testSchemas() { - return Stream.of(Arguments.of(ID_DATA_SCHEMA), Arguments.of(ID_VARIANTDATA_SCHEMA)); - } - @Test public void testMetadataUpdateWithoutActionCannotDeserialize() { List invalidJson = @@ -121,19 +108,19 @@ public void testUpgradeFormatVersionFromJson() { } /** AddSchema * */ - @ParameterizedTest - @MethodSource("testSchemas") - public void testAddSchemaFromJson(Schema schema) { + @Test + public void testAddSchemaFromJson() { String action = MetadataUpdateParser.ADD_SCHEMA; + Schema schema = ID_DATA_SCHEMA; String json = String.format("{\"action\":\"add-schema\",\"schema\":%s}", SchemaParser.toJson(schema)); MetadataUpdate actualUpdate = new MetadataUpdate.AddSchema(schema); assertEquals(action, actualUpdate, MetadataUpdateParser.fromJson(json)); } - @ParameterizedTest - @MethodSource("testSchemas") - public void testAddSchemaToJson(Schema schema) { + @Test + public void testAddSchemaToJson() { + Schema schema = ID_DATA_SCHEMA; int lastColumnId = schema.highestFieldId(); String expected = String.format( diff --git a/core/src/test/java/org/apache/iceberg/TestSchemaParser.java b/core/src/test/java/org/apache/iceberg/TestSchemaParser.java index ebd197a68af0..40db5cfee2cb 100644 --- a/core/src/test/java/org/apache/iceberg/TestSchemaParser.java +++ b/core/src/test/java/org/apache/iceberg/TestSchemaParser.java @@ -123,4 +123,14 @@ public void testPrimitiveTypeDefaultValues(Type.PrimitiveType type, Object defau assertThat(serialized.findField("col_with_default").initialDefault()).isEqualTo(defaultValue); assertThat(serialized.findField("col_with_default").writeDefault()).isEqualTo(defaultValue); } + + @Test + public void testVariantType() throws IOException { + Schema schema = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "data", Types.VariantType.get())); + + writeAndValidate(schema); + } } diff --git a/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java b/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java index e1b795d93e71..79a2afccdb9a 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java @@ -409,15 +409,14 @@ public void projectVariantSchemaUnchanged() { final org.apache.avro.Schema expected = SchemaBuilder.record("variant") + .prop(AvroSchemaUtil.FIELD_ID_PROP, "1") .namespace("unit.test") .fields() .name("metadata") - .prop(AvroSchemaUtil.FIELD_ID_PROP, "1") .type() .bytesType() .noDefault() .name("value") - .prop(AvroSchemaUtil.FIELD_ID_PROP, "2") .type() .bytesType() .noDefault() diff --git a/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java b/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java index bf5160fd18dc..7141e7c7b7b1 100644 --- a/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java +++ b/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java @@ -18,11 +18,18 @@ */ package org.apache.iceberg.data.avro; +import static org.apache.iceberg.avro.AvroSchemaUtil.convert; +import static org.apache.iceberg.types.Types.NestedField.optional; +import static org.apache.iceberg.types.Types.NestedField.required; import static org.assertj.core.api.Assertions.assertThat; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.JsonNode; import java.io.File; import java.io.IOException; import java.util.List; +import org.apache.avro.generic.GenericData; +import org.apache.avro.generic.GenericRecordBuilder; import org.apache.iceberg.Files; import org.apache.iceberg.Schema; import org.apache.iceberg.avro.Avro; @@ -33,6 +40,9 @@ import org.apache.iceberg.data.Record; import org.apache.iceberg.io.FileAppender; import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.types.Types; +import org.apache.iceberg.util.JsonUtil; +import org.junit.jupiter.api.Test; public class TestGenericData extends DataTest { @Override @@ -76,4 +86,29 @@ protected void writeAndValidate(Schema writeSchema, Schema expectedSchema) throw protected boolean supportsDefaultValues() { return true; } + + @Test + public void testSchemaWithTwoVariants() throws JsonProcessingException { + final Schema FILE_SCHEMA = + new Schema( + required(1, "id", Types.IntegerType.get()), + optional(2, "v1", Types.VariantType.get()), + optional(3, "v2", Types.VariantType.get())); + + GenericRecordBuilder builder = new GenericRecordBuilder(convert(FILE_SCHEMA, "table")); + builder.set("id", 1); + + GenericData.Record record = builder.build(); + String expectedSchema = + "{\"type\":\"record\",\"name\":\"table\"," + + "\"fields\":[{\"name\":\"id\",\"type\":\"int\",\"field-id\":1}," + + "{\"name\":\"v1\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"r2\"," + + "\"fields\":[{\"name\":\"metadata\",\"type\":\"bytes\"}," + + "{\"name\":\"value\",\"type\":\"bytes\"}]}],\"default\":null,\"field-id\":2}," + + "{\"name\":\"v2\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"r3\"," + + "\"fields\":[{\"name\":\"metadata\",\"type\":\"bytes\"}," + + "{\"name\":\"value\",\"type\":\"bytes\"}]}],\"default\":null,\"field-id\":3}]}"; + assertThat(JsonUtil.mapper().readValue(record.getSchema().toString(), JsonNode.class)) + .isEqualTo(JsonUtil.mapper().readValue(expectedSchema, JsonNode.class)); + } } From fcb3b67672035389f4dcfab8ccd736c735a604a1 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Fri, 31 Jan 2025 10:29:36 -0800 Subject: [PATCH 05/11] Update fromPrimitiveString to fromTypeString --- .palantir/revapi.yml | 4 +++ .../apache/iceberg/types/PrimitiveHolder.java | 2 +- .../java/org/apache/iceberg/types/Types.java | 2 +- .../org/apache/iceberg/types/TestTypes.java | 27 +++++++++---------- .../java/org/apache/iceberg/SchemaParser.java | 2 +- .../iceberg/TestSchemaUnionByFieldName.java | 2 +- .../iceberg/avro/TestSchemaConversions.java | 4 +-- .../iceberg/data/avro/TestGenericData.java | 4 +-- 8 files changed, 25 insertions(+), 22 deletions(-) diff --git a/.palantir/revapi.yml b/.palantir/revapi.yml index 18c63fbe7bb1..ab5ea2f3a3ba 100644 --- a/.palantir/revapi.yml +++ b/.palantir/revapi.yml @@ -1146,6 +1146,10 @@ acceptedBreaks: \ org.apache.iceberg.TableMetadata)" justification: "Removing deprecated code" "1.7.0": + org.apache.iceberg:iceberg-api: + - code: "java.method.removed" + old: "method org.apache.iceberg.types.Type.PrimitiveType org.apache.iceberg.types.Types::fromPrimitiveString(java.lang.String)" + justification: "Replace with fromTypeString to support variant type" org.apache.iceberg:iceberg-core: - code: "java.method.removed" old: "method org.apache.iceberg.deletes.PositionDeleteIndex\ diff --git a/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java b/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java index 42f0da38167d..a9330cdd730f 100644 --- a/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java +++ b/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java @@ -33,6 +33,6 @@ class PrimitiveHolder implements Serializable { } Object readResolve() throws ObjectStreamException { - return Types.fromPrimitiveString(typeAsString); + return Types.fromTypeString(typeAsString); } } diff --git a/api/src/main/java/org/apache/iceberg/types/Types.java b/api/src/main/java/org/apache/iceberg/types/Types.java index 0bbecfff5d1d..fb12be08f00a 100644 --- a/api/src/main/java/org/apache/iceberg/types/Types.java +++ b/api/src/main/java/org/apache/iceberg/types/Types.java @@ -63,7 +63,7 @@ private Types() {} private static final Pattern DECIMAL = Pattern.compile("decimal\\(\\s*(\\d+)\\s*,\\s*(\\d+)\\s*\\)"); - public static Type fromPrimitiveString(String typeString) { + public static Type fromTypeString(String typeString) { String lowerTypeString = typeString.toLowerCase(Locale.ROOT); if (TYPES.containsKey(lowerTypeString)) { return TYPES.get(lowerTypeString); diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypes.java b/api/src/test/java/org/apache/iceberg/types/TestTypes.java index 8063fd688a53..7a88a6bdf2f6 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypes.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypes.java @@ -26,28 +26,27 @@ public class TestTypes { @Test - public void fromPrimitiveString() { - assertThat(Types.fromPrimitiveString("boolean")).isSameAs(Types.BooleanType.get()); - assertThat(Types.fromPrimitiveString("BooLean")).isSameAs(Types.BooleanType.get()); + public void fromTypeString() { + assertThat(Types.fromTypeString("boolean")).isSameAs(Types.BooleanType.get()); + assertThat(Types.fromTypeString("BooLean")).isSameAs(Types.BooleanType.get()); - assertThat(Types.fromPrimitiveString("timestamp")).isSameAs(Types.TimestampType.withoutZone()); - assertThat(Types.fromPrimitiveString("timestamptz")).isSameAs(Types.TimestampType.withZone()); - assertThat(Types.fromPrimitiveString("timestamp_ns")) + assertThat(Types.fromTypeString("timestamp")).isSameAs(Types.TimestampType.withoutZone()); + assertThat(Types.fromTypeString("timestamptz")).isSameAs(Types.TimestampType.withZone()); + assertThat(Types.fromTypeString("timestamp_ns")) .isSameAs(Types.TimestampNanoType.withoutZone()); - assertThat(Types.fromPrimitiveString("timestamptz_ns")) - .isSameAs(Types.TimestampNanoType.withZone()); + assertThat(Types.fromTypeString("timestamptz_ns")).isSameAs(Types.TimestampNanoType.withZone()); - assertThat(Types.fromPrimitiveString("Fixed[ 3 ]")).isEqualTo(Types.FixedType.ofLength(3)); + assertThat(Types.fromTypeString("Fixed[ 3 ]")).isEqualTo(Types.FixedType.ofLength(3)); - assertThat(Types.fromPrimitiveString("Decimal( 2 , 3 )")).isEqualTo(Types.DecimalType.of(2, 3)); + assertThat(Types.fromTypeString("Decimal( 2 , 3 )")).isEqualTo(Types.DecimalType.of(2, 3)); - assertThat(Types.fromPrimitiveString("Decimal(2,3)")).isEqualTo(Types.DecimalType.of(2, 3)); + assertThat(Types.fromTypeString("Decimal(2,3)")).isEqualTo(Types.DecimalType.of(2, 3)); - assertThat(Types.fromPrimitiveString("variant")).isSameAs(Types.VariantType.get()); - assertThat(Types.fromPrimitiveString("Variant")).isSameAs(Types.VariantType.get()); + assertThat(Types.fromTypeString("variant")).isSameAs(Types.VariantType.get()); + assertThat(Types.fromTypeString("Variant")).isSameAs(Types.VariantType.get()); assertThatExceptionOfType(IllegalArgumentException.class) - .isThrownBy(() -> Types.fromPrimitiveString("abcdefghij")) + .isThrownBy(() -> Types.fromTypeString("abcdefghij")) .withMessage("Cannot parse type string to primitive: abcdefghij"); } } diff --git a/core/src/main/java/org/apache/iceberg/SchemaParser.java b/core/src/main/java/org/apache/iceberg/SchemaParser.java index 610140ba908c..40ef13b004b4 100644 --- a/core/src/main/java/org/apache/iceberg/SchemaParser.java +++ b/core/src/main/java/org/apache/iceberg/SchemaParser.java @@ -181,7 +181,7 @@ public static String toJson(Schema schema, boolean pretty) { private static Type typeFromJson(JsonNode json) { if (json.isTextual()) { - return Types.fromPrimitiveString(json.asText()); + return Types.fromTypeString(json.asText()); } else if (json.isObject()) { JsonNode typeObj = json.get(TYPE); if (typeObj != null) { diff --git a/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java b/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java index 656e72a0c19c..3a0423e8e2b7 100644 --- a/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java +++ b/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java @@ -75,7 +75,7 @@ private static NestedField[] primitiveFields( optional( atomicInteger.incrementAndGet(), type.toString(), - Types.fromPrimitiveString(type.toString()))) + Types.fromTypeString(type.toString()))) .toArray(NestedField[]::new); } diff --git a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java index ac86a09e7d04..9f4cac7e81f8 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java @@ -380,7 +380,7 @@ public void testVariantConversion() { org.apache.avro.Schema variantSchema = avroSchema.getField("variantCol").schema(); assertThat(variantSchema.getType()).isEqualTo(org.apache.avro.Schema.Type.RECORD); assertThat(variantSchema.getFields().size()).isEqualTo(2); - assertThat(variantSchema.getField("metadata")).isNotNull(); - assertThat(variantSchema.getField("value")).isNotNull(); + assertThat(variantSchema.getField("metadata").schema().getType()).isEqualTo(Schema.Type.BYTES); + assertThat(variantSchema.getField("value").schema().getType()).isEqualTo(Schema.Type.BYTES); } } diff --git a/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java b/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java index 7141e7c7b7b1..22d5e9c92179 100644 --- a/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java +++ b/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java @@ -89,13 +89,13 @@ protected boolean supportsDefaultValues() { @Test public void testSchemaWithTwoVariants() throws JsonProcessingException { - final Schema FILE_SCHEMA = + final Schema schema = new Schema( required(1, "id", Types.IntegerType.get()), optional(2, "v1", Types.VariantType.get()), optional(3, "v2", Types.VariantType.get())); - GenericRecordBuilder builder = new GenericRecordBuilder(convert(FILE_SCHEMA, "table")); + GenericRecordBuilder builder = new GenericRecordBuilder(convert(schema, "table")); builder.set("id", 1); GenericData.Record record = builder.build(); From 50d1482b2ed96745a211629148a4b17b50bdeda9 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Mon, 3 Feb 2025 11:57:09 -0800 Subject: [PATCH 06/11] Deprecate fromPrimitiveString --- .palantir/revapi.yml | 4 --- .../java/org/apache/iceberg/types/Types.java | 10 +++++++ .../org/apache/iceberg/types/TestTypes.java | 30 +++++++++++++++++++ 3 files changed, 40 insertions(+), 4 deletions(-) diff --git a/.palantir/revapi.yml b/.palantir/revapi.yml index ab5ea2f3a3ba..18c63fbe7bb1 100644 --- a/.palantir/revapi.yml +++ b/.palantir/revapi.yml @@ -1146,10 +1146,6 @@ acceptedBreaks: \ org.apache.iceberg.TableMetadata)" justification: "Removing deprecated code" "1.7.0": - org.apache.iceberg:iceberg-api: - - code: "java.method.removed" - old: "method org.apache.iceberg.types.Type.PrimitiveType org.apache.iceberg.types.Types::fromPrimitiveString(java.lang.String)" - justification: "Replace with fromTypeString to support variant type" org.apache.iceberg:iceberg-core: - code: "java.method.removed" old: "method org.apache.iceberg.deletes.PositionDeleteIndex\ diff --git a/api/src/main/java/org/apache/iceberg/types/Types.java b/api/src/main/java/org/apache/iceberg/types/Types.java index fb12be08f00a..50e9c48b63ef 100644 --- a/api/src/main/java/org/apache/iceberg/types/Types.java +++ b/api/src/main/java/org/apache/iceberg/types/Types.java @@ -82,6 +82,16 @@ public static Type fromTypeString(String typeString) { throw new IllegalArgumentException("Cannot parse type string to primitive: " + typeString); } + @Deprecated + public static PrimitiveType fromPrimitiveString(String typeString) { + Type type = fromTypeString(typeString); + if (type.isPrimitiveType()) { + return (PrimitiveType) type; + } + + throw new IllegalArgumentException("Cannot parse type string to primitive: " + typeString); + } + public static class BooleanType extends PrimitiveType { private static final BooleanType INSTANCE = new BooleanType(); diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypes.java b/api/src/test/java/org/apache/iceberg/types/TestTypes.java index 7a88a6bdf2f6..21511e55ed3f 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypes.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypes.java @@ -49,4 +49,34 @@ public void fromTypeString() { .isThrownBy(() -> Types.fromTypeString("abcdefghij")) .withMessage("Cannot parse type string to primitive: abcdefghij"); } + + @Test + public void fromPrimitiveString() { + assertThat(Types.fromPrimitiveString("boolean")).isSameAs(Types.BooleanType.get()); + assertThat(Types.fromPrimitiveString("BooLean")).isSameAs(Types.BooleanType.get()); + + assertThat(Types.fromPrimitiveString("timestamp")).isSameAs(Types.TimestampType.withoutZone()); + assertThat(Types.fromPrimitiveString("timestamptz")).isSameAs(Types.TimestampType.withZone()); + assertThat(Types.fromPrimitiveString("timestamp_ns")) + .isSameAs(Types.TimestampNanoType.withoutZone()); + assertThat(Types.fromPrimitiveString("timestamptz_ns")) + .isSameAs(Types.TimestampNanoType.withZone()); + + assertThat(Types.fromPrimitiveString("Fixed[ 3 ]")).isEqualTo(Types.FixedType.ofLength(3)); + + assertThat(Types.fromPrimitiveString("Decimal( 2 , 3 )")).isEqualTo(Types.DecimalType.of(2, 3)); + + assertThat(Types.fromPrimitiveString("Decimal(2,3)")).isEqualTo(Types.DecimalType.of(2, 3)); + + assertThatExceptionOfType(IllegalArgumentException.class) + .isThrownBy(() -> Types.fromPrimitiveString("variant")) + .withMessage("Cannot parse type string to primitive: variant"); + assertThatExceptionOfType(IllegalArgumentException.class) + .isThrownBy(() -> Types.fromPrimitiveString("Variant")) + .withMessage("Cannot parse type string to primitive: Variant"); + + assertThatExceptionOfType(IllegalArgumentException.class) + .isThrownBy(() -> Types.fromPrimitiveString("abcdefghij")) + .withMessage("Cannot parse type string to primitive: abcdefghij"); + } } From f53769b3804badd827daef203aa61f7289e33014 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Fri, 7 Feb 2025 11:43:33 -0800 Subject: [PATCH 07/11] Update tests --- .../apache/iceberg/types/PrimitiveHolder.java | 2 +- .../java/org/apache/iceberg/types/Types.java | 5 +- .../apache/iceberg/types/TestTypeUtil.java | 74 +++++++++++++++---- .../org/apache/iceberg/types/TestTypes.java | 27 ++++--- .../java/org/apache/iceberg/SchemaParser.java | 6 +- .../iceberg/avro/BuildAvroProjection.java | 17 ----- .../schema/SchemaWithPartnerVisitor.java | 7 +- .../iceberg/schema/UnionByNameVisitor.java | 5 ++ .../iceberg/TestSchemaUnionByFieldName.java | 25 ++++--- .../iceberg/avro/TestBuildAvroProjection.java | 28 ------- .../iceberg/avro/TestSchemaConversions.java | 17 +++-- .../iceberg/data/avro/TestGenericData.java | 35 --------- 12 files changed, 115 insertions(+), 133 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java b/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java index a9330cdd730f..928a65878d3a 100644 --- a/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java +++ b/api/src/main/java/org/apache/iceberg/types/PrimitiveHolder.java @@ -33,6 +33,6 @@ class PrimitiveHolder implements Serializable { } Object readResolve() throws ObjectStreamException { - return Types.fromTypeString(typeAsString); + return Types.fromTypeName(typeAsString); } } diff --git a/api/src/main/java/org/apache/iceberg/types/Types.java b/api/src/main/java/org/apache/iceberg/types/Types.java index 50e9c48b63ef..2bd46df41dee 100644 --- a/api/src/main/java/org/apache/iceberg/types/Types.java +++ b/api/src/main/java/org/apache/iceberg/types/Types.java @@ -63,7 +63,7 @@ private Types() {} private static final Pattern DECIMAL = Pattern.compile("decimal\\(\\s*(\\d+)\\s*,\\s*(\\d+)\\s*\\)"); - public static Type fromTypeString(String typeString) { + public static Type fromTypeName(String typeString) { String lowerTypeString = typeString.toLowerCase(Locale.ROOT); if (TYPES.containsKey(lowerTypeString)) { return TYPES.get(lowerTypeString); @@ -82,9 +82,8 @@ public static Type fromTypeString(String typeString) { throw new IllegalArgumentException("Cannot parse type string to primitive: " + typeString); } - @Deprecated public static PrimitiveType fromPrimitiveString(String typeString) { - Type type = fromTypeString(typeString); + Type type = fromTypeName(typeString); if (type.isPrimitiveType()) { return (PrimitiveType) type; } diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java index 272f58bda5ba..ce69f5de3277 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java @@ -26,11 +26,15 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Stream; import org.apache.iceberg.Schema; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Sets; import org.apache.iceberg.types.Types.IntegerType; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; public class TestTypeUtil { @Test @@ -648,36 +652,76 @@ public void testReassignOrRefreshIdsCaseInsensitive() { assertThat(actualSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } - @Test - public void testVariantType() { + private static Stream testTypes() { + return Stream.of( + Arguments.of(Types.UnknownType.get()), + Arguments.of(Types.VariantType.get()), + Arguments.of(Types.TimestampNanoType.withoutZone()), + Arguments.of(Types.TimestampNanoType.withZone())); + } + + @ParameterizedTest + @MethodSource("testTypes") + public void testAssignFreshIdsWithType(Type testType) { Schema schema = - new Schema( - required(0, "v", Types.VariantType.get()), required(1, "A", Types.IntegerType.get())); + new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); Schema sourceSchema = - new Schema( - required(1, "v", Types.VariantType.get()), required(2, "A", Types.IntegerType.get())); + new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); Schema assignedSchema = TypeUtil.assignFreshIds(sourceSchema, new AtomicInteger(10)::incrementAndGet); Schema expectedSchema = - new Schema( - required(11, "v", Types.VariantType.get()), required(12, "A", Types.IntegerType.get())); + new Schema(required(11, "v", testType), required(12, "A", Types.IntegerType.get())); assertThat(assignedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); + } + + @ParameterizedTest + @MethodSource("testTypes") + public void testReassignIdsWithType(Type testType) { + Schema schema = + new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema sourceSchema = + new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); final Schema reassignedSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(reassignedSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); + } - expectedSchema = new Schema(Lists.newArrayList(required(1, "v", Types.VariantType.get()))); - Schema projectedSchema = TypeUtil.project(sourceSchema, Sets.newHashSet(1)); - assertThat(projectedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); - - Set projectedIds = TypeUtil.getProjectedIds(sourceSchema); - assertThat(Set.of(1, 2)).isEqualTo(projectedIds); + @ParameterizedTest + @MethodSource("testTypes") + public void testIndexByIdWithType(Type testType) { + Schema schema = + new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema sourceSchema = + new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); Map indexByIds = TypeUtil.indexById(sourceSchema.asStruct()); - assertThat(indexByIds.get(1).type().isVariantType()).isTrue(); + assertThat(indexByIds.get(1).type()).isEqualTo(testType); + } + + @ParameterizedTest + @MethodSource("testTypes") + public void testIndexNameByIdWithType(Type testType) { + Schema schema = + new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema sourceSchema = + new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); Map indexNameByIds = TypeUtil.indexNameById(sourceSchema.asStruct()); assertThat(indexNameByIds.get(1)).isEqualTo("v"); } + + @ParameterizedTest + @MethodSource("testTypes") + public void testProjectWithType(Type testType) { + Schema sourceSchema = + new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); + + Schema expectedSchema = new Schema(Lists.newArrayList(required(1, "v", testType))); + Schema projectedSchema = TypeUtil.project(sourceSchema, Sets.newHashSet(1)); + assertThat(projectedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); + + Set projectedIds = TypeUtil.getProjectedIds(sourceSchema); + assertThat(Set.of(1, 2)).isEqualTo(projectedIds); + } } diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypes.java b/api/src/test/java/org/apache/iceberg/types/TestTypes.java index 21511e55ed3f..73b071e821cf 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypes.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypes.java @@ -26,27 +26,26 @@ public class TestTypes { @Test - public void fromTypeString() { - assertThat(Types.fromTypeString("boolean")).isSameAs(Types.BooleanType.get()); - assertThat(Types.fromTypeString("BooLean")).isSameAs(Types.BooleanType.get()); + public void fromTypeName() { + assertThat(Types.fromTypeName("boolean")).isSameAs(Types.BooleanType.get()); + assertThat(Types.fromTypeName("BooLean")).isSameAs(Types.BooleanType.get()); - assertThat(Types.fromTypeString("timestamp")).isSameAs(Types.TimestampType.withoutZone()); - assertThat(Types.fromTypeString("timestamptz")).isSameAs(Types.TimestampType.withZone()); - assertThat(Types.fromTypeString("timestamp_ns")) - .isSameAs(Types.TimestampNanoType.withoutZone()); - assertThat(Types.fromTypeString("timestamptz_ns")).isSameAs(Types.TimestampNanoType.withZone()); + assertThat(Types.fromTypeName("timestamp")).isSameAs(Types.TimestampType.withoutZone()); + assertThat(Types.fromTypeName("timestamptz")).isSameAs(Types.TimestampType.withZone()); + assertThat(Types.fromTypeName("timestamp_ns")).isSameAs(Types.TimestampNanoType.withoutZone()); + assertThat(Types.fromTypeName("timestamptz_ns")).isSameAs(Types.TimestampNanoType.withZone()); - assertThat(Types.fromTypeString("Fixed[ 3 ]")).isEqualTo(Types.FixedType.ofLength(3)); + assertThat(Types.fromTypeName("Fixed[ 3 ]")).isEqualTo(Types.FixedType.ofLength(3)); - assertThat(Types.fromTypeString("Decimal( 2 , 3 )")).isEqualTo(Types.DecimalType.of(2, 3)); + assertThat(Types.fromTypeName("Decimal( 2 , 3 )")).isEqualTo(Types.DecimalType.of(2, 3)); - assertThat(Types.fromTypeString("Decimal(2,3)")).isEqualTo(Types.DecimalType.of(2, 3)); + assertThat(Types.fromTypeName("Decimal(2,3)")).isEqualTo(Types.DecimalType.of(2, 3)); - assertThat(Types.fromTypeString("variant")).isSameAs(Types.VariantType.get()); - assertThat(Types.fromTypeString("Variant")).isSameAs(Types.VariantType.get()); + assertThat(Types.fromTypeName("variant")).isSameAs(Types.VariantType.get()); + assertThat(Types.fromTypeName("Variant")).isSameAs(Types.VariantType.get()); assertThatExceptionOfType(IllegalArgumentException.class) - .isThrownBy(() -> Types.fromTypeString("abcdefghij")) + .isThrownBy(() -> Types.fromTypeName("abcdefghij")) .withMessage("Cannot parse type string to primitive: abcdefghij"); } diff --git a/core/src/main/java/org/apache/iceberg/SchemaParser.java b/core/src/main/java/org/apache/iceberg/SchemaParser.java index 40ef13b004b4..04655ce3f7d7 100644 --- a/core/src/main/java/org/apache/iceberg/SchemaParser.java +++ b/core/src/main/java/org/apache/iceberg/SchemaParser.java @@ -143,9 +143,7 @@ static void toJson(Type.PrimitiveType primitive, JsonGenerator generator) throws } static void toJson(Type type, JsonGenerator generator) throws IOException { - if (type.isPrimitiveType()) { - toJson(type.asPrimitiveType(), generator); - } else if (type.isVariantType()) { + if (type.isPrimitiveType() || type.isVariantType()) { generator.writeString(type.toString()); } else { Type.NestedType nested = type.asNestedType(); @@ -181,7 +179,7 @@ public static String toJson(Schema schema, boolean pretty) { private static Type typeFromJson(JsonNode json) { if (json.isTextual()) { - return Types.fromTypeString(json.asText()); + return Types.fromTypeName(json.asText()); } else if (json.isObject()) { JsonNode typeObj = json.get(TYPE); if (typeObj != null) { diff --git a/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java b/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java index e7d099bd2683..c5c78dd1472a 100644 --- a/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java +++ b/core/src/main/java/org/apache/iceberg/avro/BuildAvroProjection.java @@ -56,12 +56,6 @@ class BuildAvroProjection extends AvroCustomOrderSchemaVisitor names, Iterable schemaIterable) { - // TODO: When the Variant logical type is introduced in Avro, handle the following in a separate - // visitor method - if (current.isVariantType()) { - return variant(record); - } - Preconditions.checkArgument( current.isNestedType() && current.asNestedType().isStructType(), "Cannot project non-struct: %s", @@ -136,17 +130,6 @@ public Schema record(Schema record, List names, Iterable s return record; } - private Schema variant(Schema record) { - Preconditions.checkArgument( - current.isVariantType() - && record.getField("value") != null - && record.getField("metadata") != null, - "Expect variant type with value and metadata fields: %s", - current); - - return record; - } - @Override public Schema.Field field(Schema.Field field, Supplier fieldResult) { Types.StructType struct = current.asNestedType().asStructType(); diff --git a/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java b/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java index 9b2226f5714d..33ddba95f5b6 100644 --- a/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java +++ b/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java @@ -106,7 +106,8 @@ public static T visit( } return visitor.map(map, partner, keyResult, valueResult); - + case VARIANT: + return visitor.variant(partner); default: return visitor.primitive(type.asPrimitiveType(), partner); } @@ -160,6 +161,10 @@ public R map(Types.MapType map, P partner, R keyResult, R valueResult) { return null; } + public R variant(P partner) { + throw new UnsupportedOperationException("Unsupported type: variant"); + } + public R primitive(Type.PrimitiveType primitive, P partner) { return null; } diff --git a/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java b/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java index 68172b7062a6..f97e7ccfcff6 100644 --- a/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java +++ b/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java @@ -142,6 +142,11 @@ public Boolean map( return false; } + @Override + public Boolean variant(Integer partnerId) { + return partnerId == null; + } + @Override public Boolean primitive(Type.PrimitiveType primitive, Integer partnerId) { return partnerId == null; diff --git a/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java b/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java index 3a0423e8e2b7..8649ada99bef 100644 --- a/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java +++ b/core/src/test/java/org/apache/iceberg/TestSchemaUnionByFieldName.java @@ -26,7 +26,7 @@ import java.util.List; import java.util.concurrent.atomic.AtomicInteger; import org.apache.iceberg.relocated.com.google.common.collect.Lists; -import org.apache.iceberg.types.Type.PrimitiveType; +import org.apache.iceberg.types.Type; import org.apache.iceberg.types.Types; import org.apache.iceberg.types.Types.BinaryType; import org.apache.iceberg.types.Types.BooleanType; @@ -42,13 +42,16 @@ import org.apache.iceberg.types.Types.StringType; import org.apache.iceberg.types.Types.StructType; import org.apache.iceberg.types.Types.TimeType; +import org.apache.iceberg.types.Types.TimestampNanoType; import org.apache.iceberg.types.Types.TimestampType; import org.apache.iceberg.types.Types.UUIDType; +import org.apache.iceberg.types.Types.UnknownType; +import org.apache.iceberg.types.Types.VariantType; import org.junit.jupiter.api.Test; public class TestSchemaUnionByFieldName { - private static List primitiveTypes() { + private static List primitiveTypes() { return Lists.newArrayList( StringType.get(), TimeType.get(), @@ -63,11 +66,15 @@ private static List primitiveTypes() { FixedType.ofLength(10), DecimalType.of(10, 2), LongType.get(), - FloatType.get()); + FloatType.get(), + VariantType.get(), + UnknownType.get(), + TimestampNanoType.withoutZone(), + TimestampNanoType.withZone()); } private static NestedField[] primitiveFields( - Integer initialValue, List primitiveTypes) { + Integer initialValue, List primitiveTypes) { AtomicInteger atomicInteger = new AtomicInteger(initialValue); return primitiveTypes.stream() .map( @@ -75,7 +82,7 @@ private static NestedField[] primitiveFields( optional( atomicInteger.incrementAndGet(), type.toString(), - Types.fromTypeString(type.toString()))) + Types.fromTypeName(type.toString()))) .toArray(NestedField[]::new); } @@ -88,7 +95,7 @@ public void testAddTopLevelPrimitives() { @Test public void testAddTopLevelListOfPrimitives() { - for (PrimitiveType primitiveType : primitiveTypes()) { + for (Type primitiveType : primitiveTypes()) { Schema newSchema = new Schema(optional(1, "aList", Types.ListType.ofOptional(2, primitiveType))); Schema applied = new SchemaUpdate(new Schema(), 0).unionByNameWith(newSchema).apply(); @@ -98,7 +105,7 @@ public void testAddTopLevelListOfPrimitives() { @Test public void testAddTopLevelMapOfPrimitives() { - for (PrimitiveType primitiveType : primitiveTypes()) { + for (Type primitiveType : primitiveTypes()) { Schema newSchema = new Schema( optional(1, "aMap", Types.MapType.ofOptional(2, 3, primitiveType, primitiveType))); @@ -109,7 +116,7 @@ public void testAddTopLevelMapOfPrimitives() { @Test public void testAddTopLevelStructOfPrimitives() { - for (PrimitiveType primitiveType : primitiveTypes()) { + for (Type primitiveType : primitiveTypes()) { Schema currentSchema = new Schema( optional(1, "aStruct", Types.StructType.of(optional(2, "primitive", primitiveType)))); @@ -120,7 +127,7 @@ public void testAddTopLevelStructOfPrimitives() { @Test public void testAddNestedPrimitive() { - for (PrimitiveType primitiveType : primitiveTypes()) { + for (Type primitiveType : primitiveTypes()) { Schema currentSchema = new Schema(optional(1, "aStruct", Types.StructType.of())); Schema newSchema = new Schema( diff --git a/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java b/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java index 79a2afccdb9a..eaea4394dbfd 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestBuildAvroProjection.java @@ -22,7 +22,6 @@ import static org.assertj.core.api.Assertions.assertThat; import java.util.Collections; -import java.util.List; import java.util.function.Supplier; import org.apache.avro.SchemaBuilder; import org.apache.iceberg.types.Type; @@ -402,31 +401,4 @@ public void projectMapWithLessFieldInValueSchema() { .as("Unexpected value ID discovered on the projected map schema") .isEqualTo(1); } - - @Test - public void projectVariantSchemaUnchanged() { - final Type icebergType = Types.VariantType.get(); - - final org.apache.avro.Schema expected = - SchemaBuilder.record("variant") - .prop(AvroSchemaUtil.FIELD_ID_PROP, "1") - .namespace("unit.test") - .fields() - .name("metadata") - .type() - .bytesType() - .noDefault() - .name("value") - .type() - .bytesType() - .noDefault() - .endRecord(); - - final BuildAvroProjection testSubject = - new BuildAvroProjection(icebergType, Collections.emptyMap()); - final org.apache.avro.Schema actual = testSubject.record(expected, List.of(), null); - assertThat(actual) - .as("Variant projection produced undesired variant schema") - .isEqualTo(expected); - } } diff --git a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java index 9f4cac7e81f8..bf2c911f2d76 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java @@ -374,13 +374,18 @@ public void testFieldDocsArePreserved() { @Test public void testVariantConversion() { org.apache.iceberg.Schema schema = - new org.apache.iceberg.Schema(required(1, "variantCol", Types.VariantType.get())); + new org.apache.iceberg.Schema( + required(1, "variantCol1", Types.VariantType.get()), + required(2, "variantCol2", Types.VariantType.get())); org.apache.avro.Schema avroSchema = AvroSchemaUtil.convert(schema.asStruct()); - org.apache.avro.Schema variantSchema = avroSchema.getField("variantCol").schema(); - assertThat(variantSchema.getType()).isEqualTo(org.apache.avro.Schema.Type.RECORD); - assertThat(variantSchema.getFields().size()).isEqualTo(2); - assertThat(variantSchema.getField("metadata").schema().getType()).isEqualTo(Schema.Type.BYTES); - assertThat(variantSchema.getField("value").schema().getType()).isEqualTo(Schema.Type.BYTES); + for (int id : Lists.newArrayList(1, 2)) { + org.apache.avro.Schema variantSchema = avroSchema.getField("variantCol" + id).schema(); + assertThat(variantSchema.getType()).isEqualTo(org.apache.avro.Schema.Type.RECORD); + assertThat(variantSchema.getFields().size()).isEqualTo(2); + assertThat(variantSchema.getField("metadata").schema().getType()) + .isEqualTo(Schema.Type.BYTES); + assertThat(variantSchema.getField("value").schema().getType()).isEqualTo(Schema.Type.BYTES); + } } } diff --git a/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java b/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java index 22d5e9c92179..bf5160fd18dc 100644 --- a/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java +++ b/data/src/test/java/org/apache/iceberg/data/avro/TestGenericData.java @@ -18,18 +18,11 @@ */ package org.apache.iceberg.data.avro; -import static org.apache.iceberg.avro.AvroSchemaUtil.convert; -import static org.apache.iceberg.types.Types.NestedField.optional; -import static org.apache.iceberg.types.Types.NestedField.required; import static org.assertj.core.api.Assertions.assertThat; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.JsonNode; import java.io.File; import java.io.IOException; import java.util.List; -import org.apache.avro.generic.GenericData; -import org.apache.avro.generic.GenericRecordBuilder; import org.apache.iceberg.Files; import org.apache.iceberg.Schema; import org.apache.iceberg.avro.Avro; @@ -40,9 +33,6 @@ import org.apache.iceberg.data.Record; import org.apache.iceberg.io.FileAppender; import org.apache.iceberg.relocated.com.google.common.collect.Lists; -import org.apache.iceberg.types.Types; -import org.apache.iceberg.util.JsonUtil; -import org.junit.jupiter.api.Test; public class TestGenericData extends DataTest { @Override @@ -86,29 +76,4 @@ protected void writeAndValidate(Schema writeSchema, Schema expectedSchema) throw protected boolean supportsDefaultValues() { return true; } - - @Test - public void testSchemaWithTwoVariants() throws JsonProcessingException { - final Schema schema = - new Schema( - required(1, "id", Types.IntegerType.get()), - optional(2, "v1", Types.VariantType.get()), - optional(3, "v2", Types.VariantType.get())); - - GenericRecordBuilder builder = new GenericRecordBuilder(convert(schema, "table")); - builder.set("id", 1); - - GenericData.Record record = builder.build(); - String expectedSchema = - "{\"type\":\"record\",\"name\":\"table\"," - + "\"fields\":[{\"name\":\"id\",\"type\":\"int\",\"field-id\":1}," - + "{\"name\":\"v1\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"r2\"," - + "\"fields\":[{\"name\":\"metadata\",\"type\":\"bytes\"}," - + "{\"name\":\"value\",\"type\":\"bytes\"}]}],\"default\":null,\"field-id\":2}," - + "{\"name\":\"v2\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"r3\"," - + "\"fields\":[{\"name\":\"metadata\",\"type\":\"bytes\"}," - + "{\"name\":\"value\",\"type\":\"bytes\"}]}],\"default\":null,\"field-id\":3}]}"; - assertThat(JsonUtil.mapper().readValue(record.getSchema().toString(), JsonNode.class)) - .isEqualTo(JsonUtil.mapper().readValue(expectedSchema, JsonNode.class)); - } } From 93586e3cf83e493c45e93cced144212f37f6524a Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Wed, 12 Feb 2025 18:31:03 -0800 Subject: [PATCH 08/11] Return same variant instance and update test --- .../java/org/apache/iceberg/Accessors.java | 2 +- .../apache/iceberg/types/AssignFreshIds.java | 4 +-- .../apache/iceberg/types/FindTypeVisitor.java | 6 ++-- .../apache/iceberg/types/GetProjectedIds.java | 2 +- .../org/apache/iceberg/types/IndexById.java | 2 +- .../org/apache/iceberg/types/IndexByName.java | 2 +- .../apache/iceberg/types/IndexParents.java | 2 +- .../apache/iceberg/types/PruneColumns.java | 2 +- .../org/apache/iceberg/types/ReassignIds.java | 4 +-- .../java/org/apache/iceberg/types/Type.java | 4 +++ .../org/apache/iceberg/types/TypeUtil.java | 10 +++--- .../java/org/apache/iceberg/types/Types.java | 9 ++++-- .../apache/iceberg/types/TestTypeUtil.java | 31 +++++++------------ .../org/apache/iceberg/types/TestTypes.java | 4 +-- .../org/apache/iceberg/avro/TypeToSchema.java | 2 +- .../schema/SchemaWithPartnerVisitor.java | 4 +-- .../iceberg/schema/UnionByNameVisitor.java | 2 +- 17 files changed, 47 insertions(+), 45 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/Accessors.java b/api/src/main/java/org/apache/iceberg/Accessors.java index b372f6a87233..0b36730fbb4b 100644 --- a/api/src/main/java/org/apache/iceberg/Accessors.java +++ b/api/src/main/java/org/apache/iceberg/Accessors.java @@ -233,7 +233,7 @@ public Map> struct( } @Override - public Map> variant() { + public Map> variant(Types.VariantType variant) { return null; } diff --git a/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java b/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java index f6422671bacb..75055cddc197 100644 --- a/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java +++ b/api/src/main/java/org/apache/iceberg/types/AssignFreshIds.java @@ -125,8 +125,8 @@ public Type map(Types.MapType map, Supplier keyFuture, Supplier valu } @Override - public Type variant() { - return Types.VariantType.get(); + public Type variant(Types.VariantType variant) { + return variant; } @Override diff --git a/api/src/main/java/org/apache/iceberg/types/FindTypeVisitor.java b/api/src/main/java/org/apache/iceberg/types/FindTypeVisitor.java index f0750f337e2e..64faebb48243 100644 --- a/api/src/main/java/org/apache/iceberg/types/FindTypeVisitor.java +++ b/api/src/main/java/org/apache/iceberg/types/FindTypeVisitor.java @@ -77,9 +77,9 @@ public Type map(Types.MapType map, Type keyResult, Type valueResult) { } @Override - public Type variant() { - if (predicate.test(Types.VariantType.get())) { - return Types.VariantType.get(); + public Type variant(Types.VariantType variant) { + if (predicate.test(variant)) { + return variant; } return null; diff --git a/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java b/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java index 814bb72f201c..1ec70b8578bc 100644 --- a/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java +++ b/api/src/main/java/org/apache/iceberg/types/GetProjectedIds.java @@ -76,7 +76,7 @@ public Set map(Types.MapType map, Set keyResult, Set } @Override - public Set variant() { + public Set variant(Types.VariantType variant) { return null; } } diff --git a/api/src/main/java/org/apache/iceberg/types/IndexById.java b/api/src/main/java/org/apache/iceberg/types/IndexById.java index 7c36a3ec241a..a7b96eb381f7 100644 --- a/api/src/main/java/org/apache/iceberg/types/IndexById.java +++ b/api/src/main/java/org/apache/iceberg/types/IndexById.java @@ -66,7 +66,7 @@ public Map map( } @Override - public Map variant() { + public Map variant(Types.VariantType variant) { return null; } } diff --git a/api/src/main/java/org/apache/iceberg/types/IndexByName.java b/api/src/main/java/org/apache/iceberg/types/IndexByName.java index 131434c9a156..60258f5c5c3e 100644 --- a/api/src/main/java/org/apache/iceberg/types/IndexByName.java +++ b/api/src/main/java/org/apache/iceberg/types/IndexByName.java @@ -177,7 +177,7 @@ public Map map( } @Override - public Map variant() { + public Map variant(Types.VariantType variant) { return nameToId; } diff --git a/api/src/main/java/org/apache/iceberg/types/IndexParents.java b/api/src/main/java/org/apache/iceberg/types/IndexParents.java index 952447ed2799..6e611d47e912 100644 --- a/api/src/main/java/org/apache/iceberg/types/IndexParents.java +++ b/api/src/main/java/org/apache/iceberg/types/IndexParents.java @@ -77,7 +77,7 @@ public Map map( } @Override - public Map variant() { + public Map variant(Types.VariantType variant) { return idToParent; } diff --git a/api/src/main/java/org/apache/iceberg/types/PruneColumns.java b/api/src/main/java/org/apache/iceberg/types/PruneColumns.java index d9230ce6e02b..56f01cf34bb5 100644 --- a/api/src/main/java/org/apache/iceberg/types/PruneColumns.java +++ b/api/src/main/java/org/apache/iceberg/types/PruneColumns.java @@ -160,7 +160,7 @@ public Type map(Types.MapType map, Type ignored, Type valueResult) { } @Override - public Type variant() { + public Type variant(Types.VariantType variant) { return null; } diff --git a/api/src/main/java/org/apache/iceberg/types/ReassignIds.java b/api/src/main/java/org/apache/iceberg/types/ReassignIds.java index dd737f5308d3..3d114f093f6b 100644 --- a/api/src/main/java/org/apache/iceberg/types/ReassignIds.java +++ b/api/src/main/java/org/apache/iceberg/types/ReassignIds.java @@ -158,8 +158,8 @@ public Type map(Types.MapType map, Supplier keyTypeFuture, Supplier } @Override - public Type variant() { - return Types.VariantType.get(); + public Type variant(Types.VariantType variant) { + return variant; } @Override diff --git a/api/src/main/java/org/apache/iceberg/types/Type.java b/api/src/main/java/org/apache/iceberg/types/Type.java index 53018ffac65b..67e40df9e939 100644 --- a/api/src/main/java/org/apache/iceberg/types/Type.java +++ b/api/src/main/java/org/apache/iceberg/types/Type.java @@ -82,6 +82,10 @@ default Types.MapType asMapType() { throw new IllegalArgumentException("Not a map type: " + this); } + default Types.VariantType asVariantType() { + throw new IllegalArgumentException("Not a variant type: " + this); + } + default boolean isNestedType() { return false; } diff --git a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java index c5c97e723834..e1cb123d3504 100644 --- a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java +++ b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java @@ -616,7 +616,7 @@ public T map(Types.MapType map, T keyResult, T valueResult) { return null; } - public T variant() { + public T variant(Types.VariantType variant) { throw new UnsupportedOperationException("Unsupported type: variant"); } @@ -684,7 +684,7 @@ public static T visit(Type type, SchemaVisitor visitor) { return visitor.map(map, keyResult, valueResult); case VARIANT: - return visitor.variant(); + return visitor.variant(type.asVariantType()); default: return visitor.primitive(type.asPrimitiveType()); @@ -712,8 +712,8 @@ public T map(Types.MapType map, Supplier keyResult, Supplier valueResult) return null; } - public T variant() { - return null; + public T variant(Types.VariantType variant) { + throw new UnsupportedOperationException("Unsupported type: variant"); } public T primitive(Type.PrimitiveType primitive) { @@ -793,7 +793,7 @@ public static T visit(Type type, CustomOrderSchemaVisitor visitor) { new VisitFuture<>(map.valueType(), visitor)); case VARIANT: - return visitor.variant(); + return visitor.variant(type.asVariantType()); default: return visitor.primitive(type.asPrimitiveType()); diff --git a/api/src/main/java/org/apache/iceberg/types/Types.java b/api/src/main/java/org/apache/iceberg/types/Types.java index 2bd46df41dee..c1935d6980e9 100644 --- a/api/src/main/java/org/apache/iceberg/types/Types.java +++ b/api/src/main/java/org/apache/iceberg/types/Types.java @@ -85,10 +85,10 @@ public static Type fromTypeName(String typeString) { public static PrimitiveType fromPrimitiveString(String typeString) { Type type = fromTypeName(typeString); if (type.isPrimitiveType()) { - return (PrimitiveType) type; + return type.asPrimitiveType(); } - throw new IllegalArgumentException("Cannot parse type string to primitive: " + typeString); + throw new IllegalArgumentException("Cannot parse type string: variant is not a primitive type"); } public static class BooleanType extends PrimitiveType { @@ -445,6 +445,11 @@ public boolean isVariantType() { return true; } + @Override + public VariantType asVariantType() { + return this; + } + @Override public boolean equals(Object o) { if (this == o) { diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java index ce69f5de3277..02465f9cd8e7 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java @@ -665,11 +665,8 @@ private static Stream testTypes() { public void testAssignFreshIdsWithType(Type testType) { Schema schema = new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); - Schema sourceSchema = - new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); - Schema assignedSchema = - TypeUtil.assignFreshIds(sourceSchema, new AtomicInteger(10)::incrementAndGet); + Schema assignedSchema = TypeUtil.assignFreshIds(schema, new AtomicInteger(10)::incrementAndGet); Schema expectedSchema = new Schema(required(11, "v", testType), required(12, "A", Types.IntegerType.get())); assertThat(assignedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); @@ -683,7 +680,7 @@ public void testReassignIdsWithType(Type testType) { Schema sourceSchema = new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); - final Schema reassignedSchema = TypeUtil.reassignIds(schema, sourceSchema); + Schema reassignedSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(reassignedSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); } @@ -692,11 +689,9 @@ public void testReassignIdsWithType(Type testType) { public void testIndexByIdWithType(Type testType) { Schema schema = new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); - Schema sourceSchema = - new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); - Map indexByIds = TypeUtil.indexById(sourceSchema.asStruct()); - assertThat(indexByIds.get(1).type()).isEqualTo(testType); + Map indexByIds = TypeUtil.indexById(schema.asStruct()); + assertThat(indexByIds.get(0).type()).isEqualTo(testType); } @ParameterizedTest @@ -704,24 +699,22 @@ public void testIndexByIdWithType(Type testType) { public void testIndexNameByIdWithType(Type testType) { Schema schema = new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); - Schema sourceSchema = - new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); - Map indexNameByIds = TypeUtil.indexNameById(sourceSchema.asStruct()); - assertThat(indexNameByIds.get(1)).isEqualTo("v"); + Map indexNameByIds = TypeUtil.indexNameById(schema.asStruct()); + assertThat(indexNameByIds.get(0)).isEqualTo("v"); } @ParameterizedTest @MethodSource("testTypes") public void testProjectWithType(Type testType) { - Schema sourceSchema = - new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); + Schema schema = + new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); - Schema expectedSchema = new Schema(Lists.newArrayList(required(1, "v", testType))); - Schema projectedSchema = TypeUtil.project(sourceSchema, Sets.newHashSet(1)); + Schema expectedSchema = new Schema(Lists.newArrayList(required(0, "v", testType))); + Schema projectedSchema = TypeUtil.project(schema, Sets.newHashSet(0)); assertThat(projectedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); - Set projectedIds = TypeUtil.getProjectedIds(sourceSchema); - assertThat(Set.of(1, 2)).isEqualTo(projectedIds); + Set projectedIds = TypeUtil.getProjectedIds(schema); + assertThat(Set.of(0, 1)).isEqualTo(projectedIds); } } diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypes.java b/api/src/test/java/org/apache/iceberg/types/TestTypes.java index 73b071e821cf..b3381d1ff440 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypes.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypes.java @@ -69,10 +69,10 @@ public void fromPrimitiveString() { assertThatExceptionOfType(IllegalArgumentException.class) .isThrownBy(() -> Types.fromPrimitiveString("variant")) - .withMessage("Cannot parse type string to primitive: variant"); + .withMessage("Cannot parse type string: variant is not a primitive type"); assertThatExceptionOfType(IllegalArgumentException.class) .isThrownBy(() -> Types.fromPrimitiveString("Variant")) - .withMessage("Cannot parse type string to primitive: Variant"); + .withMessage("Cannot parse type string: variant is not a primitive type"); assertThatExceptionOfType(IllegalArgumentException.class) .isThrownBy(() -> Types.fromPrimitiveString("abcdefghij")) diff --git a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java index c009eed76e7c..9896110af272 100644 --- a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java +++ b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java @@ -188,7 +188,7 @@ public Schema map(Types.MapType map, Schema keySchema, Schema valueSchema) { } @Override - public Schema variant() { + public Schema variant(Types.VariantType variant) { String recordName = "r" + fieldIds.peek(); return Schema.createRecord( recordName, diff --git a/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java b/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java index 33ddba95f5b6..0d00123a588b 100644 --- a/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java +++ b/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java @@ -107,7 +107,7 @@ public static T visit( return visitor.map(map, partner, keyResult, valueResult); case VARIANT: - return visitor.variant(partner); + return visitor.variant(type.asVariantType(), partner); default: return visitor.primitive(type.asPrimitiveType(), partner); } @@ -161,7 +161,7 @@ public R map(Types.MapType map, P partner, R keyResult, R valueResult) { return null; } - public R variant(P partner) { + public R variant(Types.VariantType variant, P partner) { throw new UnsupportedOperationException("Unsupported type: variant"); } diff --git a/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java b/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java index f97e7ccfcff6..7c4dac9feff1 100644 --- a/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java +++ b/core/src/main/java/org/apache/iceberg/schema/UnionByNameVisitor.java @@ -143,7 +143,7 @@ public Boolean map( } @Override - public Boolean variant(Integer partnerId) { + public Boolean variant(Types.VariantType variant, Integer partnerId) { return partnerId == null; } From 1d54ff457e8140291aa6b273dcec4d89c85e7789 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Thu, 13 Feb 2025 20:50:53 -0800 Subject: [PATCH 09/11] Remove Avro change --- .../org/apache/iceberg/types/TestTypeUtil.java | 9 ++++++++- .../org/apache/iceberg/avro/TypeToSchema.java | 12 ------------ .../schema/SchemaWithPartnerVisitor.java | 2 ++ .../iceberg/avro/TestSchemaConversions.java | 18 ------------------ 4 files changed, 10 insertions(+), 31 deletions(-) diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java index 02465f9cd8e7..6ae7fac1e631 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java @@ -710,9 +710,16 @@ public void testProjectWithType(Type testType) { Schema schema = new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); - Schema expectedSchema = new Schema(Lists.newArrayList(required(0, "v", testType))); + Schema expectedSchema = new Schema(required(0, "v", testType)); Schema projectedSchema = TypeUtil.project(schema, Sets.newHashSet(0)); assertThat(projectedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); + } + + @ParameterizedTest + @MethodSource("testTypes") + public void testGetProjectedIdsWithType(Type testType) { + Schema schema = + new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); Set projectedIds = TypeUtil.getProjectedIds(schema); assertThat(Set.of(0, 1)).isEqualTo(projectedIds); diff --git a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java index 9896110af272..05ce4e618662 100644 --- a/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java +++ b/core/src/main/java/org/apache/iceberg/avro/TypeToSchema.java @@ -187,18 +187,6 @@ public Schema map(Types.MapType map, Schema keySchema, Schema valueSchema) { return mapSchema; } - @Override - public Schema variant(Types.VariantType variant) { - String recordName = "r" + fieldIds.peek(); - return Schema.createRecord( - recordName, - null, - null, - false, - List.of( - new Schema.Field("metadata", BINARY_SCHEMA), new Schema.Field("value", BINARY_SCHEMA))); - } - @Override public Schema primitive(Type.PrimitiveType primitive) { Schema primitiveSchema; diff --git a/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java b/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java index 0d00123a588b..694bfb2f6242 100644 --- a/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java +++ b/core/src/main/java/org/apache/iceberg/schema/SchemaWithPartnerVisitor.java @@ -106,8 +106,10 @@ public static T visit( } return visitor.map(map, partner, keyResult, valueResult); + case VARIANT: return visitor.variant(type.asVariantType(), partner); + default: return visitor.primitive(type.asPrimitiveType(), partner); } diff --git a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java index bf2c911f2d76..d9dc49d17257 100644 --- a/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java +++ b/core/src/test/java/org/apache/iceberg/avro/TestSchemaConversions.java @@ -370,22 +370,4 @@ public void testFieldDocsArePreserved() { Lists.newArrayList(Iterables.transform(origSchema.columns(), Types.NestedField::doc)); assertThat(fieldDocs).isEqualTo(origFieldDocs); } - - @Test - public void testVariantConversion() { - org.apache.iceberg.Schema schema = - new org.apache.iceberg.Schema( - required(1, "variantCol1", Types.VariantType.get()), - required(2, "variantCol2", Types.VariantType.get())); - org.apache.avro.Schema avroSchema = AvroSchemaUtil.convert(schema.asStruct()); - - for (int id : Lists.newArrayList(1, 2)) { - org.apache.avro.Schema variantSchema = avroSchema.getField("variantCol" + id).schema(); - assertThat(variantSchema.getType()).isEqualTo(org.apache.avro.Schema.Type.RECORD); - assertThat(variantSchema.getFields().size()).isEqualTo(2); - assertThat(variantSchema.getField("metadata").schema().getType()) - .isEqualTo(Schema.Type.BYTES); - assertThat(variantSchema.getField("value").schema().getType()).isEqualTo(Schema.Type.BYTES); - } - } } From 9082cabc91f6350187d8c298bf4aef5afb2c93b2 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Thu, 13 Feb 2025 23:39:48 -0800 Subject: [PATCH 10/11] Add remaining visitor implementation in core --- .../org/apache/iceberg/types/AssignIds.java | 5 ++ .../iceberg/types/CheckCompatibility.java | 13 +++++ .../org/apache/iceberg/types/ReassignDoc.java | 5 ++ .../org/apache/iceberg/types/TypeUtil.java | 8 +++ .../iceberg/types/TestReadabilityChecks.java | 30 +++++++++++ .../apache/iceberg/types/TestTypeUtil.java | 54 ++++++++++++------- .../java/org/apache/iceberg/SchemaUpdate.java | 5 ++ .../apache/iceberg/mapping/MappingUtil.java | 5 ++ .../org/apache/iceberg/types/FixupTypes.java | 6 +++ .../iceberg/mapping/TestNameMapping.java | 12 +++++ .../org/apache/iceberg/spark/Spark3Util.java | 5 ++ .../apache/iceberg/spark/TestSpark3Util.java | 5 +- 12 files changed, 133 insertions(+), 20 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/types/AssignIds.java b/api/src/main/java/org/apache/iceberg/types/AssignIds.java index 68588f581adc..b2f72751eb89 100644 --- a/api/src/main/java/org/apache/iceberg/types/AssignIds.java +++ b/api/src/main/java/org/apache/iceberg/types/AssignIds.java @@ -92,6 +92,11 @@ public Type map(Types.MapType map, Supplier keyFuture, Supplier valu } } + @Override + public Type variant(Types.VariantType variant) { + return variant; + } + @Override public Type primitive(Type.PrimitiveType primitive) { return primitive; diff --git a/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java b/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java index 502e52c345e5..e7cc8712172e 100644 --- a/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java +++ b/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java @@ -250,6 +250,19 @@ public List map( } } + @Override + public List variant(Types.VariantType readVariant) { + if (currentType.isVariantType()) { + return NO_ERRORS; + } + + // Currently promotion is not allowed to variant type + return ImmutableList.of( + String.format( + ": %s cannot be read as a %s", + currentType.typeId().toString().toLowerCase(Locale.ENGLISH), readVariant)); + } + @Override public List primitive(Type.PrimitiveType readPrimitive) { if (currentType.equals(readPrimitive)) { diff --git a/api/src/main/java/org/apache/iceberg/types/ReassignDoc.java b/api/src/main/java/org/apache/iceberg/types/ReassignDoc.java index 9ce04a7bd103..328d81c42885 100644 --- a/api/src/main/java/org/apache/iceberg/types/ReassignDoc.java +++ b/api/src/main/java/org/apache/iceberg/types/ReassignDoc.java @@ -96,6 +96,11 @@ public Type map(Types.MapType map, Supplier keyTypeFuture, Supplier } } + @Override + public Type variant(Types.VariantType variant) { + return variant; + } + @Override public Type primitive(Type.PrimitiveType primitive) { return primitive; diff --git a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java index e1cb123d3504..4892696ab450 100644 --- a/api/src/main/java/org/apache/iceberg/types/TypeUtil.java +++ b/api/src/main/java/org/apache/iceberg/types/TypeUtil.java @@ -616,6 +616,14 @@ public T map(Types.MapType map, T keyResult, T valueResult) { return null; } + /** + * @deprecated will be removed in 2.0.0; use {@link #variant(Types.VariantType)} instead. + */ + @Deprecated + public T variant() { + return variant(Types.VariantType.get()); + } + public T variant(Types.VariantType variant) { throw new UnsupportedOperationException("Unsupported type: variant"); } diff --git a/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java b/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java index 2d02da5346a7..350f17671285 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java +++ b/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java @@ -22,6 +22,8 @@ import static org.apache.iceberg.types.Types.NestedField.required; import static org.assertj.core.api.Assertions.assertThat; +import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import org.apache.iceberg.Schema; import org.apache.iceberg.types.Type.PrimitiveType; @@ -112,6 +114,34 @@ private void testDisallowPrimitiveToStruct(PrimitiveType from, Schema fromSchema .contains("cannot be read as a struct"); } + @Test + public void testVariantType() { + Schema fromSchema = new Schema(required(1, "from_field", Types.VariantType.get())); + List errors = + CheckCompatibility.writeCompatibilityErrors( + new Schema(required(1, "to_field", Types.VariantType.get())), fromSchema); + assertThat(errors).as("Should produce 0 error messages").isEmpty(); + + List incompatibleTypes = new ArrayList<>(); + incompatibleTypes.addAll( + List.of( + Types.StructType.of(required(1, "from", Types.IntegerType.get())), + Types.MapType.ofRequired(1, 2, Types.StringType.get(), Types.IntegerType.get()), + Types.ListType.ofRequired(1, Types.StringType.get()))); + incompatibleTypes.addAll(Arrays.asList(PRIMITIVES)); + + for (Type from : incompatibleTypes) { + fromSchema = new Schema(required(3, "from_field", from)); + errors = + CheckCompatibility.writeCompatibilityErrors( + new Schema(required(3, "to_field", Types.VariantType.get())), fromSchema); + assertThat(errors).hasSize(1); + assertThat(errors.get(0)) + .as("Should complain that other type to variant is not allowed") + .contains("cannot be read as a variant"); + } + } + @Test public void testRequiredSchemaField() { Schema write = new Schema(optional(1, "from_field", Types.IntegerType.get())); diff --git a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java index 6ae7fac1e631..b6556d70bd85 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java +++ b/api/src/test/java/org/apache/iceberg/types/TestTypeUtil.java @@ -660,25 +660,35 @@ private static Stream testTypes() { Arguments.of(Types.TimestampNanoType.withZone())); } + @ParameterizedTest + @MethodSource("testTypes") + public void testAssignIdsWithType(Type testType) { + Types.StructType sourceType = + Types.StructType.of(required(0, "id", IntegerType.get()), required(1, "data", testType)); + Type expectedType = + Types.StructType.of(required(10, "id", IntegerType.get()), required(11, "data", testType)); + + Type assignedType = TypeUtil.assignIds(sourceType, oldId -> oldId + 10); + assertThat(assignedType).isEqualTo(expectedType); + } + @ParameterizedTest @MethodSource("testTypes") public void testAssignFreshIdsWithType(Type testType) { - Schema schema = - new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema schema = new Schema(required(0, "id", IntegerType.get()), required(1, "data", testType)); Schema assignedSchema = TypeUtil.assignFreshIds(schema, new AtomicInteger(10)::incrementAndGet); Schema expectedSchema = - new Schema(required(11, "v", testType), required(12, "A", Types.IntegerType.get())); + new Schema(required(11, "id", IntegerType.get()), required(12, "data", testType)); assertThat(assignedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } @ParameterizedTest @MethodSource("testTypes") public void testReassignIdsWithType(Type testType) { - Schema schema = - new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema schema = new Schema(required(0, "id", IntegerType.get()), required(1, "data", testType)); Schema sourceSchema = - new Schema(required(1, "v", testType), required(2, "A", Types.IntegerType.get())); + new Schema(required(1, "id", IntegerType.get()), required(2, "data", testType)); Schema reassignedSchema = TypeUtil.reassignIds(schema, sourceSchema); assertThat(reassignedSchema.asStruct()).isEqualTo(sourceSchema.asStruct()); @@ -687,41 +697,49 @@ public void testReassignIdsWithType(Type testType) { @ParameterizedTest @MethodSource("testTypes") public void testIndexByIdWithType(Type testType) { - Schema schema = - new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema schema = new Schema(required(0, "id", IntegerType.get()), required(1, "data", testType)); Map indexByIds = TypeUtil.indexById(schema.asStruct()); - assertThat(indexByIds.get(0).type()).isEqualTo(testType); + assertThat(indexByIds.get(1).type()).isEqualTo(testType); } @ParameterizedTest @MethodSource("testTypes") public void testIndexNameByIdWithType(Type testType) { - Schema schema = - new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema schema = new Schema(required(0, "id", IntegerType.get()), required(1, "data", testType)); Map indexNameByIds = TypeUtil.indexNameById(schema.asStruct()); - assertThat(indexNameByIds.get(0)).isEqualTo("v"); + assertThat(indexNameByIds.get(1)).isEqualTo("data"); } @ParameterizedTest @MethodSource("testTypes") public void testProjectWithType(Type testType) { - Schema schema = - new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema schema = new Schema(required(0, "id", IntegerType.get()), required(1, "data", testType)); - Schema expectedSchema = new Schema(required(0, "v", testType)); - Schema projectedSchema = TypeUtil.project(schema, Sets.newHashSet(0)); + Schema expectedSchema = new Schema(required(1, "data", testType)); + Schema projectedSchema = TypeUtil.project(schema, Sets.newHashSet(1)); assertThat(projectedSchema.asStruct()).isEqualTo(expectedSchema.asStruct()); } @ParameterizedTest @MethodSource("testTypes") public void testGetProjectedIdsWithType(Type testType) { - Schema schema = - new Schema(required(0, "v", testType), required(1, "A", Types.IntegerType.get())); + Schema schema = new Schema(required(0, "id", IntegerType.get()), required(1, "data", testType)); Set projectedIds = TypeUtil.getProjectedIds(schema); assertThat(Set.of(0, 1)).isEqualTo(projectedIds); } + + @ParameterizedTest + @MethodSource("testTypes") + public void testReassignDocWithType(Type testType) { + Schema schema = new Schema(required(0, "id", IntegerType.get()), required(1, "data", testType)); + Schema docSourceSchema = + new Schema( + required(0, "id", IntegerType.get(), "id"), required(1, "data", testType, "data")); + + Schema reassignedSchema = TypeUtil.reassignDoc(schema, docSourceSchema); + assertThat(reassignedSchema.asStruct()).isEqualTo(docSourceSchema.asStruct()); + } } diff --git a/core/src/main/java/org/apache/iceberg/SchemaUpdate.java b/core/src/main/java/org/apache/iceberg/SchemaUpdate.java index 2b541080ac72..7726c3a785d0 100644 --- a/core/src/main/java/org/apache/iceberg/SchemaUpdate.java +++ b/core/src/main/java/org/apache/iceberg/SchemaUpdate.java @@ -722,6 +722,11 @@ public Type map(Types.MapType map, Type kResult, Type valueResult) { } } + @Override + public Type variant(Types.VariantType variant) { + return variant; + } + @Override public Type primitive(Type.PrimitiveType primitive) { return primitive; diff --git a/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java b/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java index de6ce2ad0425..e42fcdf9c135 100644 --- a/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java +++ b/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java @@ -302,6 +302,11 @@ public MappedFields map(Types.MapType map, MappedFields keyResult, MappedFields MappedField.of(map.valueId(), "value", valueResult)); } + @Override + public MappedFields variant(Types.VariantType variant) { + return null; // no mapping because variant no nested fields + } + @Override public MappedFields primitive(Type.PrimitiveType primitive) { return null; // no mapping because primitives have no nested fields diff --git a/core/src/main/java/org/apache/iceberg/types/FixupTypes.java b/core/src/main/java/org/apache/iceberg/types/FixupTypes.java index 23fccddda3d9..1e4c0b597a6a 100644 --- a/core/src/main/java/org/apache/iceberg/types/FixupTypes.java +++ b/core/src/main/java/org/apache/iceberg/types/FixupTypes.java @@ -147,6 +147,12 @@ public Type map(Types.MapType map, Supplier keyTypeFuture, Supplier } } + @Override + public Type variant(Types.VariantType variant) { + // nothing to fix up + return variant; + } + @Override public Type primitive(Type.PrimitiveType primitive) { if (sourceType.equals(primitive)) { diff --git a/core/src/test/java/org/apache/iceberg/mapping/TestNameMapping.java b/core/src/test/java/org/apache/iceberg/mapping/TestNameMapping.java index d30a93d50d49..3ebf4d9242ab 100644 --- a/core/src/test/java/org/apache/iceberg/mapping/TestNameMapping.java +++ b/core/src/test/java/org/apache/iceberg/mapping/TestNameMapping.java @@ -289,4 +289,16 @@ public void testMappingFindByName() { "location", MappedFields.of(MappedField.of(11, "latitude"), MappedField.of(12, "longitude")))); } + + @Test + public void testMappingVariantType() { + Schema schema = + new Schema( + required(1, "id", Types.LongType.get()), required(2, "data", Types.VariantType.get())); + + MappedFields expected = MappedFields.of(MappedField.of(1, "id"), MappedField.of(2, "data")); + + NameMapping mapping = MappingUtil.create(schema); + assertThat(mapping.asMappedFields()).isEqualTo(expected); + } } diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/Spark3Util.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/Spark3Util.java index af0fa84f67a1..ad8a4beb55d0 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/Spark3Util.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/Spark3Util.java @@ -565,6 +565,11 @@ public String map(Types.MapType map, String keyResult, String valueResult) { return "map<" + keyResult + ", " + valueResult + ">"; } + @Override + public String variant(Types.VariantType variant) { + return "variant"; + } + @Override public String primitive(Type.PrimitiveType primitive) { switch (primitive.typeId()) { diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSpark3Util.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSpark3Util.java index 6f900ffebb10..e4e66abfefa0 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSpark3Util.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSpark3Util.java @@ -105,12 +105,13 @@ public void testDescribeSchema() { 3, "pairs", Types.MapType.ofOptional(4, 5, Types.StringType.get(), Types.LongType.get())), - required(6, "time", Types.TimestampType.withoutZone())); + required(6, "time", Types.TimestampType.withoutZone()), + required(7, "v", Types.VariantType.get())); assertThat(Spark3Util.describe(schema)) .as("Schema description isn't correct.") .isEqualTo( - "struct not null,pairs: map,time: timestamp not null>"); + "struct not null,pairs: map,time: timestamp not null,v: variant not null>"); } @Test From 1472417cce3e2671de8197164078f6b056c0bd09 Mon Sep 17 00:00:00 2001 From: Aihua Xu Date: Mon, 17 Feb 2025 13:30:35 -0800 Subject: [PATCH 11/11] Update tests --- .../iceberg/types/CheckCompatibility.java | 5 +- .../iceberg/types/TestReadabilityChecks.java | 49 +++++++++++-------- .../iceberg/types/TestSerializableTypes.java | 10 +--- .../apache/iceberg/mapping/MappingUtil.java | 2 +- 4 files changed, 32 insertions(+), 34 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java b/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java index e7cc8712172e..725f7f42562e 100644 --- a/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java +++ b/api/src/main/java/org/apache/iceberg/types/CheckCompatibility.java @@ -257,10 +257,7 @@ public List variant(Types.VariantType readVariant) { } // Currently promotion is not allowed to variant type - return ImmutableList.of( - String.format( - ": %s cannot be read as a %s", - currentType.typeId().toString().toLowerCase(Locale.ENGLISH), readVariant)); + return ImmutableList.of(String.format(": %s cannot be read as a %s", currentType, readVariant)); } @Override diff --git a/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java b/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java index 350f17671285..debb9c9dc1d6 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java +++ b/api/src/test/java/org/apache/iceberg/types/TestReadabilityChecks.java @@ -22,12 +22,15 @@ import static org.apache.iceberg.types.Types.NestedField.required; import static org.assertj.core.api.Assertions.assertThat; -import java.util.ArrayList; import java.util.Arrays; import java.util.List; +import java.util.stream.Stream; import org.apache.iceberg.Schema; import org.apache.iceberg.types.Type.PrimitiveType; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; public class TestReadabilityChecks { private static final Type.PrimitiveType[] PRIMITIVES = @@ -115,31 +118,37 @@ private void testDisallowPrimitiveToStruct(PrimitiveType from, Schema fromSchema } @Test - public void testVariantType() { + public void testVariantToVariant() { Schema fromSchema = new Schema(required(1, "from_field", Types.VariantType.get())); List errors = CheckCompatibility.writeCompatibilityErrors( new Schema(required(1, "to_field", Types.VariantType.get())), fromSchema); assertThat(errors).as("Should produce 0 error messages").isEmpty(); + } - List incompatibleTypes = new ArrayList<>(); - incompatibleTypes.addAll( - List.of( - Types.StructType.of(required(1, "from", Types.IntegerType.get())), - Types.MapType.ofRequired(1, 2, Types.StringType.get(), Types.IntegerType.get()), - Types.ListType.ofRequired(1, Types.StringType.get()))); - incompatibleTypes.addAll(Arrays.asList(PRIMITIVES)); - - for (Type from : incompatibleTypes) { - fromSchema = new Schema(required(3, "from_field", from)); - errors = - CheckCompatibility.writeCompatibilityErrors( - new Schema(required(3, "to_field", Types.VariantType.get())), fromSchema); - assertThat(errors).hasSize(1); - assertThat(errors.get(0)) - .as("Should complain that other type to variant is not allowed") - .contains("cannot be read as a variant"); - } + private static Stream incompatibleTypesToVariant() { + return Stream.of( + Stream.of( + Arguments.of(Types.StructType.of(required(1, "from", Types.IntegerType.get()))), + Arguments.of( + Types.MapType.ofRequired( + 1, 2, Types.StringType.get(), Types.IntegerType.get())), + Arguments.of(Types.ListType.ofRequired(1, Types.StringType.get()))), + Arrays.stream(PRIMITIVES).map(type -> Arguments.of(type))) + .flatMap(s -> s); + } + + @ParameterizedTest + @MethodSource("incompatibleTypesToVariant") + public void testIncompatibleTypesToVariant(Type from) { + Schema fromSchema = new Schema(required(3, "from_field", from)); + List errors = + CheckCompatibility.writeCompatibilityErrors( + new Schema(required(3, "to_field", Types.VariantType.get())), fromSchema); + assertThat(errors).hasSize(1); + assertThat(errors.get(0)) + .as("Should complain that other type to variant is not allowed") + .contains("cannot be read as a variant"); } @Test diff --git a/api/src/test/java/org/apache/iceberg/types/TestSerializableTypes.java b/api/src/test/java/org/apache/iceberg/types/TestSerializableTypes.java index a222e8e66b8e..790f59587c59 100644 --- a/api/src/test/java/org/apache/iceberg/types/TestSerializableTypes.java +++ b/api/src/test/java/org/apache/iceberg/types/TestSerializableTypes.java @@ -46,6 +46,7 @@ public void testIdentityTypes() throws Exception { Types.StringType.get(), Types.UUIDType.get(), Types.BinaryType.get(), + Types.UnknownType.get() }; for (Type type : identityPrimitives) { @@ -136,15 +137,6 @@ public void testVariant() throws Exception { .isEqualTo(variant); } - @Test - public void testUnknown() throws Exception { - Types.UnknownType unknown = Types.UnknownType.get(); - Type copy = TestHelpers.roundTripSerialize(unknown); - assertThat(copy) - .as("Unknown serialization should be equal to starting type") - .isEqualTo(unknown); - } - @Test public void testSchema() throws Exception { Schema schema = diff --git a/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java b/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java index e42fcdf9c135..fc3d920a4069 100644 --- a/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java +++ b/core/src/main/java/org/apache/iceberg/mapping/MappingUtil.java @@ -304,7 +304,7 @@ public MappedFields map(Types.MapType map, MappedFields keyResult, MappedFields @Override public MappedFields variant(Types.VariantType variant) { - return null; // no mapping because variant no nested fields + return null; // no mapping because variant has no nested fields with IDs } @Override