diff --git a/internal/venice-client-common/src/main/java/com/linkedin/venice/utils/AvroSupersetSchemaUtils.java b/internal/venice-client-common/src/main/java/com/linkedin/venice/utils/AvroSupersetSchemaUtils.java index 5f5e20baae0..7d0dd176faf 100644 --- a/internal/venice-client-common/src/main/java/com/linkedin/venice/utils/AvroSupersetSchemaUtils.java +++ b/internal/venice-client-common/src/main/java/com/linkedin/venice/utils/AvroSupersetSchemaUtils.java @@ -37,11 +37,20 @@ public static boolean isSupersetSchema(Schema s1, Schema s2) { * C/D could be nested record change as well eg, array/map of records, or record of records. * Prerequisite: The top-level schema are of type RECORD only and each field have default values. ie they are compatible * schemas and the generated schema will pick the default value from new value schema. - * @param existingSchema schema existing in the repo - * @param newSchema schema to be added. - * @return super-set schema of existingSchema abd newSchema + * @param rawExistingSchema schema existing in the repo + * @param rawNewSchema schema to be added. + * @return super-set schema of rawExistingSchema and rawNewSchema */ - public static Schema generateSupersetSchema(Schema existingSchema, Schema newSchema) { + public static Schema generateSupersetSchema(Schema rawExistingSchema, Schema rawNewSchema) { + // Normalize single-element unions [T] to their inner type T so that [T] and T + // are treated equivalently. Skip normalization when both inputs are unions — + // the multi-element union case is handled by unionSchema() below and must not + // be disrupted (e.g. [T] vs ["null", T] must remain a UNION-vs-UNION merge). + final boolean bothUnions = + rawExistingSchema.getType() == Schema.Type.UNION && rawNewSchema.getType() == Schema.Type.UNION; + final Schema existingSchema = bothUnions ? rawExistingSchema : unwrapSingleElementUnion(rawExistingSchema); + final Schema newSchema = bothUnions ? rawNewSchema : unwrapSingleElementUnion(rawNewSchema); + if (existingSchema.getType() != newSchema.getType()) { throw new VeniceException("Incompatible schema"); } @@ -193,6 +202,19 @@ private static Schema unionSchema(Schema existingSchema, Schema newSchema) { return Schema.createUnion(combinedSchema); } + /** + * If the schema is a UNION with exactly one member type, return that inner type. + * A single-element union {@code [T]} is semantically equivalent to {@code T} in Avro, + * but {@link Schema#getType()} returns {@link Schema.Type#UNION} for the wrapper form. + * Unwrapping normalizes both representations so they compare as the same type. + */ + private static Schema unwrapSingleElementUnion(Schema schema) { + if (schema.getType() == Schema.Type.UNION && schema.getTypes().size() == 1) { + return schema.getTypes().get(0); + } + return schema; + } + /** * Returns all custom property names for a {@link Schema} using {@link Schema#getObjectProps()} (available since * Avro 1.8) instead of {@link AvroCompatibilityHelper#getAllPropNames(Schema)}. The latter routes through diff --git a/internal/venice-client-common/src/test/java/com/linkedin/venice/schema/TestAvroSupersetSchemaUtils.java b/internal/venice-client-common/src/test/java/com/linkedin/venice/schema/TestAvroSupersetSchemaUtils.java index 34758165731..1e531dfd031 100644 --- a/internal/venice-client-common/src/test/java/com/linkedin/venice/schema/TestAvroSupersetSchemaUtils.java +++ b/internal/venice-client-common/src/test/java/com/linkedin/venice/schema/TestAvroSupersetSchemaUtils.java @@ -16,11 +16,20 @@ import com.linkedin.venice.controllerapi.MultiSchemaResponse; import com.linkedin.venice.exceptions.VeniceException; import com.linkedin.venice.schema.avro.DirectionalSchemaCompatibilityType; +import com.linkedin.venice.serializer.AvroGenericDeserializer; +import com.linkedin.venice.serializer.AvroSerializer; import com.linkedin.venice.utils.AvroSchemaUtils; import com.linkedin.venice.utils.AvroSupersetSchemaUtils; +import java.util.ArrayList; import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import org.apache.avro.Schema; +import org.apache.avro.generic.GenericData; +import org.apache.avro.generic.GenericRecord; import org.testng.Assert; import org.testng.annotations.Test; @@ -775,6 +784,160 @@ public void testGenerateSupersetSchemaWithAnnotatedPrimitiveTypesAndNewEnumField "Superset schema must retain the 'categoryRef' field from v1"); } + /** + * Validates that generateSupersetSchema() treats a single-element union wrapping a type + * (e.g. {@code [array]}) the same as the bare type ({@code array}), since they are + * semantically equivalent in Avro. + */ + @Test + public void testSchemaMergeSingleElementUnionVsPlainType() { + // v1: opportunityIds as a single-element union wrapping an array + String schemaStr1 = "{\"type\":\"record\",\"name\":\"TestRecord\",\"namespace\":\"com.example\"," + + "\"fields\":[{\"name\":\"opportunityIds\"," + "\"type\":[{\"type\":\"array\",\"items\":\"long\"}]," + + "\"doc\":\"List of opportunityIds.\"}]}"; + // v2: opportunityIds as a plain array (no union wrapper) + String schemaStr2 = "{\"type\":\"record\",\"name\":\"TestRecord\",\"namespace\":\"com.example\"," + + "\"fields\":[{\"name\":\"opportunityIds\"," + "\"type\":{\"type\":\"array\",\"items\":\"long\"}," + + "\"doc\":\"List of opportunityIds.\"}]}"; + Schema s1 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr1); + Schema s2 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr2); + + // Should not throw — these are semantically equivalent + Schema superset12 = AvroSupersetSchemaUtils.generateSupersetSchema(s1, s2); + Assert.assertNotNull(superset12); + Assert.assertNotNull(superset12.getField("opportunityIds")); + + // The result should be an array type (the unwrapped form) + Assert.assertEquals(superset12.getField("opportunityIds").schema().getType(), Schema.Type.ARRAY); + + // Should also work in the reverse direction + Schema superset21 = AvroSupersetSchemaUtils.generateSupersetSchema(s2, s1); + Assert.assertNotNull(superset21); + Assert.assertEquals(superset21.getField("opportunityIds").schema().getType(), Schema.Type.ARRAY); + + // Both orderings should produce equivalent superset schemas + Assert.assertTrue(AvroSchemaUtils.compareSchemaIgnoreFieldOrder(superset12, superset21)); + + // Verify that the superset schema can deserialize records serialized with either version + List ids = Collections.singletonList(42L); + + GenericRecord recordV1 = new GenericData.Record(s1); + recordV1.put("opportunityIds", ids); + byte[] bytesV1 = new AvroSerializer<>(s1).serialize(recordV1); + GenericRecord deserializedV1 = (GenericRecord) new AvroGenericDeserializer<>(s1, superset12).deserialize(bytesV1); + Assert.assertEquals(new ArrayList<>((Collection) deserializedV1.get("opportunityIds")), ids); + + GenericRecord recordV2 = new GenericData.Record(s2); + recordV2.put("opportunityIds", ids); + byte[] bytesV2 = new AvroSerializer<>(s2).serialize(recordV2); + GenericRecord deserializedV2 = (GenericRecord) new AvroGenericDeserializer<>(s2, superset12).deserialize(bytesV2); + Assert.assertEquals(new ArrayList<>((Collection) deserializedV2.get("opportunityIds")), ids); + } + + @Test + public void testSchemaMergeSingleElementUnionVsPlainTypeWithAdditionalFields() { + // Record with additional fields to test that single-element union unwrapping + // doesn't interfere with normal superset field merging + String schemaStr1 = "{\"type\":\"record\",\"name\":\"TestRecord\"," + "\"namespace\":\"com.example\"," + + "\"fields\":[" + "{\"name\":\"ids\",\"type\":[{\"type\":\"array\",\"items\":\"long\"}]}," + + "{\"name\":\"name\",\"type\":\"string\",\"default\":\"default\"}" + "]}"; + + String schemaStr2 = "{\"type\":\"record\",\"name\":\"TestRecord\"," + "\"namespace\":\"com.example\"," + + "\"fields\":[" + "{\"name\":\"ids\",\"type\":{\"type\":\"array\",\"items\":\"long\"}}," + + "{\"name\":\"description\",\"type\":\"string\",\"default\":\"none\"}" + "]}"; + + Schema s1 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr1); + Schema s2 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr2); + + Schema superset = AvroSupersetSchemaUtils.generateSupersetSchema(s1, s2); + Assert.assertNotNull(superset); + Assert.assertNotNull(superset.getField("ids")); + Assert.assertNotNull(superset.getField("name")); + Assert.assertNotNull(superset.getField("description")); + Assert.assertEquals(superset.getField("ids").schema().getType(), Schema.Type.ARRAY); + + List ids = Collections.singletonList(42L); + + GenericRecord recordV1 = new GenericData.Record(s1); + recordV1.put("ids", ids); + recordV1.put("name", "alice"); + byte[] bytesV1 = new AvroSerializer<>(s1).serialize(recordV1); + GenericRecord deserializedV1 = (GenericRecord) new AvroGenericDeserializer<>(s1, superset).deserialize(bytesV1); + // GenericData.Array.equals() only accepts other GenericArray — copy to ArrayList to compare + Assert.assertEquals(new ArrayList<>((Collection) deserializedV1.get("ids")), ids); + + GenericRecord recordV2 = new GenericData.Record(s2); + recordV2.put("ids", ids); + recordV2.put("description", "desc"); + byte[] bytesV2 = new AvroSerializer<>(s2).serialize(recordV2); + GenericRecord deserializedV2 = (GenericRecord) new AvroGenericDeserializer<>(s2, superset).deserialize(bytesV2); + Assert.assertEquals(new ArrayList<>((Collection) deserializedV2.get("ids")), ids); + } + + @Test + public void testSchemaMergeSingleElementUnionMap() { + // Test single-element union wrapping a map type + String schemaStr1 = "{\"type\":\"record\",\"name\":\"TestRecord\"," + "\"fields\":[{\"name\":\"metadata\"," + + "\"type\":[{\"type\":\"map\",\"values\":\"string\"}]}]}"; + + String schemaStr2 = "{\"type\":\"record\",\"name\":\"TestRecord\"," + "\"fields\":[{\"name\":\"metadata\"," + + "\"type\":{\"type\":\"map\",\"values\":\"string\"}}]}"; + + Schema s1 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr1); + Schema s2 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr2); + + Schema superset = AvroSupersetSchemaUtils.generateSupersetSchema(s1, s2); + Assert.assertNotNull(superset); + Assert.assertEquals(superset.getField("metadata").schema().getType(), Schema.Type.MAP); + + Map metadata = Collections.singletonMap("k", "v"); + + GenericRecord recordV1 = new GenericData.Record(s1); + recordV1.put("metadata", metadata); + byte[] bytesV1 = new AvroSerializer<>(s1).serialize(recordV1); + GenericRecord deserializedV1 = (GenericRecord) new AvroGenericDeserializer<>(s1, superset).deserialize(bytesV1); + // Avro deserializes string keys/values as Utf8 — convert to String for comparison + Assert.assertEquals(toStringMap((Map) deserializedV1.get("metadata")), metadata); + + GenericRecord recordV2 = new GenericData.Record(s2); + recordV2.put("metadata", metadata); + byte[] bytesV2 = new AvroSerializer<>(s2).serialize(recordV2); + GenericRecord deserializedV2 = (GenericRecord) new AvroGenericDeserializer<>(s2, superset).deserialize(bytesV2); + Assert.assertEquals(toStringMap((Map) deserializedV2.get("metadata")), metadata); + } + + @Test + public void testSchemaMergeSingleElementUnionVsMultiElementUnion() { + // [T] vs ["null", T] must still merge correctly via the union path. + // Unwrapping [T] to T before comparison would produce T vs UNION and throw. + String schemaStr1 = "{\"type\":\"record\",\"name\":\"TestRecord\"," + "\"fields\":[{\"name\":\"ids\"," + + "\"type\":[{\"type\":\"array\",\"items\":\"long\"}]}]}"; + + String schemaStr2 = "{\"type\":\"record\",\"name\":\"TestRecord\"," + "\"fields\":[{\"name\":\"ids\"," + + "\"type\":[\"null\",{\"type\":\"array\",\"items\":\"long\"}],\"default\":null}]}"; + + Schema s1 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr1); + Schema s2 = AvroSchemaParseUtils.parseSchemaFromJSONStrictValidation(schemaStr2); + + Schema superset = AvroSupersetSchemaUtils.generateSupersetSchema(s1, s2); + Assert.assertNotNull(superset); + Assert.assertEquals(superset.getField("ids").schema().getType(), Schema.Type.UNION); + + List ids = Collections.singletonList(42L); + + GenericRecord recordV1 = new GenericData.Record(s1); + recordV1.put("ids", ids); + byte[] bytesV1 = new AvroSerializer<>(s1).serialize(recordV1); + GenericRecord deserializedV1 = (GenericRecord) new AvroGenericDeserializer<>(s1, superset).deserialize(bytesV1); + Assert.assertEquals(new ArrayList<>((Collection) deserializedV1.get("ids")), ids); + + GenericRecord recordV2 = new GenericData.Record(s2); + recordV2.put("ids", ids); + byte[] bytesV2 = new AvroSerializer<>(s2).serialize(recordV2); + GenericRecord deserializedV2 = (GenericRecord) new AvroGenericDeserializer<>(s2, superset).deserialize(bytesV2); + Assert.assertEquals(new ArrayList<>((Collection) deserializedV2.get("ids")), ids); + } + @Test public void testValidateSubsetSchema() { Assert.assertTrue( @@ -794,4 +957,10 @@ public void testValidateSubsetSchema() { Assert.assertTrue( AvroSupersetSchemaUtils.validateSubsetValueSchema(NAME_RECORD_V6_SCHEMA, supersetSchemaForV5AndV4.toString())); } + + private static Map toStringMap(Map map) { + Map result = new HashMap<>(); + map.forEach((k, v) -> result.put(k.toString(), v.toString())); + return result; + } }