From b27422d2e50c4c1bfcffa8a7a96efe979f54454f Mon Sep 17 00:00:00 2001 From: Christian Tzolov Date: Fri, 1 May 2020 18:57:35 +0200 Subject: [PATCH] Avro compliant writer schema resolution - Make the readerSchema a fallback option only if the writer schema is not resolvable - Add Forward & Backward compatibility tests. Resolves #12 --- .../avro/AbstractAvroMessageConverter.java | 12 +- ...AvroMessageConverterAutoConfiguration.java | 37 +- .../avro/AvroMessageConverterProperties.java | 6 +- ...oSchemaRegistryClientMessageConverter.java | 111 ++-- .../avro/AvroSchemaServiceManagerImpl.java | 38 +- .../avro/AvroSchemaServiceManagerTests.java | 29 +- .../ForwardAndBackwardCompatibilityTest.java | 509 ++++++++++++++++++ .../cloud/schema/avro/v2/User1.java | 70 +++ .../test/resources/schemas/user1_v1.schema | 10 + .../test/resources/schemas/user1_v2.schema | 10 + .../src/test/resources/schemas/user_v2.avsc | 10 + 11 files changed, 712 insertions(+), 130 deletions(-) create mode 100644 spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/ForwardAndBackwardCompatibilityTest.java create mode 100644 spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/v2/User1.java create mode 100644 spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v1.schema create mode 100644 spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v2.schema create mode 100644 spring-cloud-schema-registry-client/src/test/resources/schemas/user_v2.avsc diff --git a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AbstractAvroMessageConverter.java b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AbstractAvroMessageConverter.java index 167c3a4..459a146 100644 --- a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AbstractAvroMessageConverter.java +++ b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AbstractAvroMessageConverter.java @@ -86,8 +86,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver } @Override - protected Object convertFromInternal(Message message, Class targetClass, - Object conversionHint) { + protected Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { Object result; try { byte[] payload = (byte[]) message.getPayload(); @@ -114,8 +113,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver } @Override - protected Object convertToInternal(Object payload, MessageHeaders headers, - Object conversionHint) { + protected Object convertToInternal(Object payload, MessageHeaders headers, Object conversionHint) { ByteArrayOutputStream baos = new ByteArrayOutputStream(); try { MimeType hintedContentType = null; @@ -124,8 +122,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver } Schema schema = resolveSchemaForWriting(payload, headers, hintedContentType); @SuppressWarnings("unchecked") - DatumWriter writer = avroSchemaServiceManager() - .getDatumWriter(payload.getClass(), schema); + DatumWriter writer = avroSchemaServiceManager().getDatumWriter(payload.getClass(), schema); Encoder encoder = EncoderFactory.get().binaryEncoder(baos, null); writer.write(payload, encoder); encoder.flush(); @@ -136,8 +133,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver return baos.toByteArray(); } - protected abstract Schema resolveSchemaForWriting(Object payload, - MessageHeaders headers, MimeType hintedContentType); + protected abstract Schema resolveSchemaForWriting(Object payload, MessageHeaders headers, MimeType hintedContentType); protected abstract Schema resolveWriterSchemaForDeserialization(MimeType mimeType); diff --git a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterAutoConfiguration.java b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterAutoConfiguration.java index e19bbdc..2f6c45b 100644 --- a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterAutoConfiguration.java +++ b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterAutoConfiguration.java @@ -50,43 +50,36 @@ public class AvroMessageConverterAutoConfiguration { @Bean @ConditionalOnMissingBean(AvroSchemaRegistryClientMessageConverter.class) public AvroSchemaRegistryClientMessageConverter avroSchemaMessageConverter( - SchemaRegistryClient schemaRegistryClient, AvroSchemaServiceManager avroSchemaServiceManager, + SchemaRegistryClient schemaRegistryClient, + AvroSchemaServiceManager avroSchemaServiceManager, AvroMessageConverterProperties avroMessageConverterProperties) { - AvroSchemaRegistryClientMessageConverter avroSchemaRegistryClientMessageConverter; - avroSchemaRegistryClientMessageConverter = new AvroSchemaRegistryClientMessageConverter( - schemaRegistryClient, cacheManager(), avroSchemaServiceManager); + + AvroSchemaRegistryClientMessageConverter avroSchemaRegistryClientMessageConverter = + new AvroSchemaRegistryClientMessageConverter(schemaRegistryClient, cacheManager(), avroSchemaServiceManager); + avroSchemaRegistryClientMessageConverter.setDynamicSchemaGenerationEnabled( avroMessageConverterProperties.isDynamicSchemaGenerationEnabled()); + if (avroMessageConverterProperties.getReaderSchema() != null) { - avroSchemaRegistryClientMessageConverter.setReaderSchema( - avroMessageConverterProperties.getReaderSchema()); + avroSchemaRegistryClientMessageConverter.setReaderSchema(avroMessageConverterProperties.getReaderSchema()); } - if (!ObjectUtils - .isEmpty(avroMessageConverterProperties.getSchemaLocations())) { - avroSchemaRegistryClientMessageConverter.setSchemaLocations( - avroMessageConverterProperties.getSchemaLocations()); + if (!ObjectUtils.isEmpty(avroMessageConverterProperties.getSchemaLocations())) { + avroSchemaRegistryClientMessageConverter.setSchemaLocations(avroMessageConverterProperties.getSchemaLocations()); } - if (!ObjectUtils - .isEmpty(avroMessageConverterProperties.getSchemaImports())) { - avroSchemaRegistryClientMessageConverter.setSchemaImports( - avroMessageConverterProperties.getSchemaImports()); + if (!ObjectUtils.isEmpty(avroMessageConverterProperties.getSchemaImports())) { + avroSchemaRegistryClientMessageConverter.setSchemaImports(avroMessageConverterProperties.getSchemaImports()); } - avroSchemaRegistryClientMessageConverter - .setPrefix(avroMessageConverterProperties.getPrefix()); + avroSchemaRegistryClientMessageConverter.setPrefix(avroMessageConverterProperties.getPrefix()); try { - Class clazz = avroMessageConverterProperties - .getSubjectNamingStrategy(); + Class clazz = avroMessageConverterProperties.getSubjectNamingStrategy(); Constructor constructor = ReflectionUtils.accessibleConstructor(clazz); - avroSchemaRegistryClientMessageConverter.setSubjectNamingStrategy( (SubjectNamingStrategy) constructor.newInstance()); } catch (Exception ex) { throw new IllegalStateException("Unable to create SubjectNamingStrategy " - + avroMessageConverterProperties.getSubjectNamingStrategy() - .toString(), - ex); + + avroMessageConverterProperties.getSubjectNamingStrategy().toString(), ex); } return avroSchemaRegistryClientMessageConverter; diff --git a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterProperties.java b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterProperties.java index 48101e6..1e7ee36 100644 --- a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterProperties.java +++ b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroMessageConverterProperties.java @@ -74,8 +74,7 @@ public class AvroMessageConverterProperties { return this.dynamicSchemaGenerationEnabled; } - public void setDynamicSchemaGenerationEnabled( - boolean dynamicSchemaGenerationEnabled) { + public void setDynamicSchemaGenerationEnabled(boolean dynamicSchemaGenerationEnabled) { this.dynamicSchemaGenerationEnabled = dynamicSchemaGenerationEnabled; } @@ -91,8 +90,7 @@ public class AvroMessageConverterProperties { return this.subjectNamingStrategy; } - public void setSubjectNamingStrategy( - Class subjectNamingStrategy) { + public void setSubjectNamingStrategy(Class subjectNamingStrategy) { Assert.notNull(subjectNamingStrategy, "cannot be null"); this.subjectNamingStrategy = subjectNamingStrategy; } diff --git a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java index 134e3f6..375ff90 100644 --- a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java +++ b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java @@ -86,8 +86,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag /** * Pattern for validating the prefix to be used in the publised subtype. */ - public static final Pattern PREFIX_VALIDATION_PATTERN = Pattern - .compile("[\\p{Alnum}]"); + public static final Pattern PREFIX_VALIDATION_PATTERN = Pattern.compile("[\\p{Alnum}]"); /** * Spring Cloud Stream schema property prefix. @@ -112,11 +111,10 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag /** * Default Mime type for Avro. */ - public static final MimeType DEFAULT_AVRO_MIME_TYPE = new MimeType("application", - "*+" + AVRO_FORMAT); + public static final MimeType DEFAULT_AVRO_MIME_TYPE = new MimeType("application", "*+" + AVRO_FORMAT); private static final AvroSchemaServiceManager defaultAvroSchemaServiceManager = - new AvroSchemaServiceManagerImpl(); + new AvroSchemaServiceManagerImpl(); private final CacheManager cacheManager; @@ -146,7 +144,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag */ @Deprecated public AvroSchemaRegistryClientMessageConverter( - SchemaRegistryClient schemaRegistryClient, CacheManager cacheManager) { + SchemaRegistryClient schemaRegistryClient, CacheManager cacheManager) { super(Collections.singletonList(DEFAULT_AVRO_MIME_TYPE), defaultAvroSchemaServiceManager); Assert.notNull(schemaRegistryClient, "cannot be null"); Assert.notNull(cacheManager, "'cacheManager' cannot be null"); @@ -182,8 +180,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag * false, it only allows the converter to use pre-registered schemas. Default 'true'. * @param dynamicSchemaGenerationEnabled true if dynamic schema generation is enabled */ - public void setDynamicSchemaGenerationEnabled( - boolean dynamicSchemaGenerationEnabled) { + public void setDynamicSchemaGenerationEnabled(boolean dynamicSchemaGenerationEnabled) { this.dynamicSchemaGenerationEnabled = dynamicSchemaGenerationEnabled; } @@ -212,8 +209,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag */ public void setPrefix(String prefix) { Assert.hasText(prefix, "Prefix cannot be empty"); - Assert.isTrue(!PREFIX_VALIDATION_PATTERN.matcher(this.prefix).matches(), - "Invalid prefix:" + this.prefix); + Assert.isTrue(!PREFIX_VALIDATION_PATTERN.matcher(this.prefix).matches(), "Invalid prefix:" + this.prefix); this.prefix = prefix; } @@ -232,22 +228,25 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag } @Override - public void afterPropertiesSet() throws Exception { + public void afterPropertiesSet() { this.versionedSchema = Pattern.compile("application/" + this.prefix + "\\.([\\p{Alnum}\\$\\.]+)\\.v(\\p{Digit}+)\\+" + AVRO_FORMAT); Stream.of(this.schemaImports, this.schemaLocations) - .filter(arr -> !ObjectUtils.isEmpty(arr)).distinct().peek(resources -> { + .filter(arr -> !ObjectUtils.isEmpty(arr)) + .distinct() + .peek(resources -> { if (this.logger.isInfoEnabled()) { this.logger.info("Scanning avro schema resources on classpath"); this.logger.info("Parsing " + this.schemaImports.length + " schemas"); } - }).flatMap(Arrays::stream).forEach(resource -> { + }) + .flatMap(Arrays::stream) + .forEach(resource -> { try { Schema schema = parseSchema(resource); if (schema.getType().equals(Schema.Type.UNION)) { - schema.getTypes().forEach( - innerSchema -> registerSchema(resource, innerSchema)); + schema.getTypes().forEach(innerSchema -> registerSchema(resource, innerSchema)); } else { registerSchema(resource, schema); @@ -255,9 +254,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag } catch (IOException e) { if (this.logger.isWarnEnabled()) { - this.logger.warn( - "Failed to parse schema at " + resource.getFilename(), - e); + this.logger.warn("Failed to parse schema at " + resource.getFilename(), e); } } }); @@ -295,57 +292,47 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag Schema schema; schema = extractSchemaForWriting(payload); - ParsedSchema parsedSchema = this.getCache(REFERENCE_CACHE_NAME) - .get(schema, ParsedSchema.class); + ParsedSchema parsedSchema = this.getCache(REFERENCE_CACHE_NAME).get(schema, ParsedSchema.class); if (parsedSchema == null) { parsedSchema = new ParsedSchema(schema); - this.getCache(REFERENCE_CACHE_NAME).putIfAbsent(schema, - parsedSchema); + this.getCache(REFERENCE_CACHE_NAME).putIfAbsent(schema, parsedSchema); } if (parsedSchema.getRegistration() == null) { - SchemaRegistrationResponse response = this.schemaRegistryClient.register( - toSubject(schema), AVRO_FORMAT, parsedSchema.getRepresentation()); + SchemaRegistrationResponse response = this.schemaRegistryClient.register(toSubject(schema), + AVRO_FORMAT, parsedSchema.getRepresentation()); parsedSchema.setRegistration(response); } - SchemaReference schemaReference = parsedSchema.getRegistration() - .getSchemaReference(); + SchemaReference schemaReference = parsedSchema.getRegistration().getSchemaReference(); DirectFieldAccessor dfa = new DirectFieldAccessor(headers); @SuppressWarnings("unchecked") - Map _headers = (Map) dfa - .getPropertyValue("headers"); - _headers.put(MessageHeaders.CONTENT_TYPE, - "application/" + this.prefix + "." + schemaReference.getSubject() + ".v" - + schemaReference.getVersion() + "+" + AVRO_FORMAT); + Map _headers = (Map) dfa.getPropertyValue("headers"); + _headers.put(MessageHeaders.CONTENT_TYPE, "application/" + this.prefix + "." + schemaReference.getSubject() + + ".v" + schemaReference.getVersion() + "+" + AVRO_FORMAT); return schema; } @Override protected Schema resolveWriterSchemaForDeserialization(MimeType mimeType) { - if (this.readerSchema == null) { - SchemaReference schemaReference = extractSchemaReference(mimeType); - if (schemaReference != null) { - ParsedSchema parsedSchema = this.getCache(REFERENCE_CACHE_NAME) - .get(schemaReference, ParsedSchema.class); - if (parsedSchema == null) { - String schemaContent = this.schemaRegistryClient - .fetch(schemaReference); - if (schemaContent != null) { - Schema schema = new Schema.Parser().parse(schemaContent); - parsedSchema = new ParsedSchema(schema); - this.getCache(REFERENCE_CACHE_NAME) - .putIfAbsent(schemaReference, parsedSchema); - } - } - if (parsedSchema != null) { - return parsedSchema.getSchema(); + SchemaReference schemaReference = extractSchemaReference(mimeType); + if (schemaReference != null) { + ParsedSchema parsedSchema = this.getCache(REFERENCE_CACHE_NAME).get(schemaReference, ParsedSchema.class); + if (parsedSchema == null) { + String schemaContent = this.schemaRegistryClient.fetch(schemaReference); + if (schemaContent != null) { + Schema schema = new Schema.Parser().parse(schemaContent); + parsedSchema = new ParsedSchema(schema); + this.getCache(REFERENCE_CACHE_NAME).putIfAbsent(schemaReference, parsedSchema); } } + if (parsedSchema != null) { + return parsedSchema.getSchema(); + } } return this.readerSchema; } @@ -367,20 +354,17 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag } } else { - schema = this.getCache(REFLECTION_CACHE_NAME) - .get(payload.getClass().getName(), Schema.class); + schema = this.getCache(REFLECTION_CACHE_NAME).get(payload.getClass().getName(), Schema.class); if (schema == null) { if (!isDynamicSchemaGenerationEnabled()) { throw new SchemaNotFoundException(String.format( - "No schema found in the local cache for %s, and dynamic schema generation " - + "is not enabled", + "No schema found in the local cache for %s, and dynamic schema generation is not enabled", payload.getClass())); } else { schema = super.avroSchemaServiceManager().getSchema(payload.getClass()); } - this.getCache(REFLECTION_CACHE_NAME) - .put(payload.getClass().getName(), schema); + this.getCache(REFLECTION_CACHE_NAME).put(payload.getClass().getName(), schema); } } return schema; @@ -388,18 +372,17 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag private void registerSchema(Resource schemaLocation, Schema schema) { if (this.logger.isInfoEnabled()) { - this.logger.info( - "Resource " + schemaLocation.getFilename() + " parsed into schema " - + schema.getNamespace() + "." + schema.getName()); + this.logger.info("Resource " + schemaLocation.getFilename() + " parsed into schema " + + schema.getNamespace() + "." + schema.getName()); } - this.schemaRegistryClient.register(toSubject(schema), AVRO_FORMAT, - schema.toString()); + + this.schemaRegistryClient.register(toSubject(schema), AVRO_FORMAT, schema.toString()); + if (this.logger.isInfoEnabled()) { - this.logger - .info("Schema " + schema.getName() + " registered with id " + schema); + this.logger.info("Schema " + schema.getName() + " registered with id " + schema); } - this.getCache(REFLECTION_CACHE_NAME) - .put(schema.getNamespace() + "." + schema.getName(), schema); + + this.getCache(REFLECTION_CACHE_NAME).put(schema.getNamespace() + "." + schema.getName(), schema); } private SchemaReference extractSchemaReference(MimeType mimeType) { @@ -417,7 +400,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag Cache cache = this.cacheManager.getCache(name); Assert.notNull(cache, "Cache by the name '" + name + "' is not present in this CacheManager - '" + this.cacheManager + "'. Typically caches are auto-created by the CacheManagers. " - + "Consider reporting it as an issue to the developer of this CacheManager."); + + "Consider reporting it as an issue to the developer of this CacheManager."); return cache; } diff --git a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaServiceManagerImpl.java b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaServiceManagerImpl.java index ca9ff49..bce4bb7 100644 --- a/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaServiceManagerImpl.java +++ b/spring-cloud-schema-registry-client/src/main/java/org/springframework/cloud/schema/registry/avro/AvroSchemaServiceManagerImpl.java @@ -41,11 +41,11 @@ import org.springframework.stereotype.Component; /** * Default Concrete implementation of {@link AvroSchemaServiceManager}. * - * Helps to substitute the default implementation of {@link org.apache.avro.Schema} - * Generation using Custom Avro schema generator + * Helps to substitute the default implementation of {@link org.apache.avro.Schema} Generation using Custom Avro + * schema generator * - * Provide a custom bean definition of {@link AvroSchemaServiceManager} and mark - * it as @Primary to override this default implementation + * Provide a custom bean definition of {@link AvroSchemaServiceManager} and mark it as @Primary to override this + * default implementation * * @author Ish Mahajan * @@ -58,8 +58,7 @@ public class AvroSchemaServiceManagerImpl implements AvroSchemaServiceManager { /** * get {@link Schema}. - * @param clazz {@link Class} for which schema generation - * is required + * @param clazz {@link Class} for which schema generation is required * @return returns avro schema for given class */ @Override @@ -102,21 +101,21 @@ public class AvroSchemaServiceManagerImpl implements AvroSchemaServiceManager { /** * get {@link DatumReader}. * @param type {@link Class} of java object which needs to be serialized - * @param schema {@link Schema} default schema of object which needs to be de-serialized + * @param readerSchema {@link Schema} default schema of object which needs to be de-serialized * @param writerSchema {@link Schema} writerSchema provided at run time * @return datum reader which can be used to read Avro payload */ - @SuppressWarnings({"unchecked", "rawtypes"}) + @SuppressWarnings({ "unchecked", "rawtypes" }) @Override - public DatumReader getDatumReader(Class type, Schema schema, Schema writerSchema) { + public DatumReader getDatumReader(Class type, Schema readerSchema, Schema writerSchema) { DatumReader reader = null; if (SpecificRecord.class.isAssignableFrom(type)) { - if (schema != null) { + if (readerSchema != null) { if (writerSchema != null) { - reader = new SpecificDatumReader<>(writerSchema, schema); + reader = new SpecificDatumReader<>(writerSchema, readerSchema); } else { - reader = new SpecificDatumReader<>(schema); + reader = new SpecificDatumReader<>(readerSchema); } } else { @@ -127,12 +126,12 @@ public class AvroSchemaServiceManagerImpl implements AvroSchemaServiceManager { } } else if (GenericRecord.class.isAssignableFrom(type)) { - if (schema != null) { + if (readerSchema != null) { if (writerSchema != null) { - reader = new GenericDatumReader<>(writerSchema, schema); + reader = new GenericDatumReader<>(writerSchema, readerSchema); } else { - reader = new GenericDatumReader<>(schema); + reader = new GenericDatumReader<>(readerSchema); } } else { @@ -149,7 +148,7 @@ public class AvroSchemaServiceManagerImpl implements AvroSchemaServiceManager { } if (reader == null) { throw new MessageConversionException("No schema can be inferred from type " - + type.getName() + " and no schema has been explicitly configured."); + + type.getName() + " and no schema has been explicitly configured."); } return reader; } @@ -164,10 +163,9 @@ public class AvroSchemaServiceManagerImpl implements AvroSchemaServiceManager { * @throws IOException is thrown in case of error */ @Override - public Object readData(Class clazz, byte[] payload, Schema readerSchema, - Schema writerSchema) throws IOException { - DatumReader reader = this.getDatumReader(clazz, - readerSchema, writerSchema); + public Object readData(Class clazz, byte[] payload, Schema readerSchema, Schema writerSchema) + throws IOException { + DatumReader reader = this.getDatumReader(clazz, readerSchema, writerSchema); Decoder decoder = DecoderFactory.get().binaryDecoder(payload, null); return reader.read(null, decoder); } diff --git a/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/AvroSchemaServiceManagerTests.java b/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/AvroSchemaServiceManagerTests.java index e8a1f39..0c6c630 100644 --- a/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/AvroSchemaServiceManagerTests.java +++ b/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/AvroSchemaServiceManagerTests.java @@ -32,6 +32,8 @@ import org.apache.avro.file.DataFileReader; import org.apache.avro.file.DataFileWriter; import org.apache.avro.io.DatumReader; import org.apache.avro.io.DatumWriter; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.assertj.core.util.Lists; import org.junit.Test; @@ -44,11 +46,14 @@ import org.springframework.util.MimeType; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.fail; + /** * @author Ish Mahajan */ public class AvroSchemaServiceManagerTests { + private final Log logger = LogFactory.getLog(AvroSchemaServiceManagerTests.class); + @SuppressWarnings({ "rawtypes", "unchecked", "resource" }) @Test(expected = DataFileWriter.AppendWriteException.class) public void testWithDefaultImplementation() throws IOException { @@ -77,7 +82,7 @@ public class AvroSchemaServiceManagerTests { // allocating and garbage collecting many objects for files with // many items. foodOrderDeserialized = dataFileReader.next(foodOrderDeserialized); - System.out.println("De-serialised Successfully : " + foodOrderDeserialized); + logger.info("De-serialised Successfully : " + foodOrderDeserialized); } } @@ -110,7 +115,7 @@ public class AvroSchemaServiceManagerTests { @Override public Object readData(Class targetClass, byte[] payload, Schema readerSchema, - Schema writerSchema) throws IOException { + Schema writerSchema) throws IOException { ObjectMapper mapper = new ObjectMapper(new AvroFactory()); AvroSchemaGenerator gen = new AvroSchemaGenerator(); try { @@ -120,8 +125,8 @@ public class AvroSchemaServiceManagerTests { fail("Error while setting acceptJsonFormatVisitor {}", e); } return mapper.readerFor(targetClass) - .with(new AvroSchema(readerSchema)) - .readValue(payload); + .with(new AvroSchema(readerSchema)) + .readValue(payload); } }; @@ -148,19 +153,19 @@ public class AvroSchemaServiceManagerTests { MimeType mimeType = new MimeType("application", "avro"); assertThat(mimeType).isEqualTo(converter.getSupportedMimeTypes().get(0)); - AvroSchemaMessageConverter converter2 = new AvroSchemaMessageConverter(mimeType); + AvroSchemaMessageConverter converter2 = new AvroSchemaMessageConverter(mimeType); assertThat(mimeType).isEqualTo(converter2.getSupportedMimeTypes().get(0)); - AvroSchemaMessageConverter converter3 = - new AvroSchemaMessageConverter(Lists.newArrayList(mimeType)); + AvroSchemaMessageConverter converter3 = + new AvroSchemaMessageConverter(Lists.newArrayList(mimeType)); assertThat(mimeType).isEqualTo(converter3.getSupportedMimeTypes().get(0)); AvroSchemaServiceManager manager = new AvroSchemaServiceManagerImpl(); - AvroSchemaMessageConverter converter4 = new AvroSchemaMessageConverter(manager); + AvroSchemaMessageConverter converter4 = new AvroSchemaMessageConverter(manager); assertThat(mimeType).isEqualTo(converter4.getSupportedMimeTypes().get(0)); - AvroSchemaMessageConverter converter5 = - new AvroSchemaMessageConverter(Lists.newArrayList(mimeType), manager); + AvroSchemaMessageConverter converter5 = + new AvroSchemaMessageConverter(Lists.newArrayList(mimeType), manager); Schema schema = manager.getSchema(FoodOrder.class); converter5.setSchema(schema); assertThat(mimeType).isEqualTo(converter5.getSupportedMimeTypes().get(0)); @@ -171,8 +176,8 @@ public class AvroSchemaServiceManagerTests { public void testAvroSchemaMessageConverterException() { MimeType mimeType = new MimeType("application", "avro"); AvroSchemaServiceManager manager = new AvroSchemaServiceManagerImpl(); - AvroSchemaMessageConverter converter = - new AvroSchemaMessageConverter(Lists.newArrayList(mimeType), manager); + AvroSchemaMessageConverter converter = + new AvroSchemaMessageConverter(Lists.newArrayList(mimeType), manager); converter.setSchemaLocation(new ByteArrayResource(new byte[2]) { }); } diff --git a/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/ForwardAndBackwardCompatibilityTest.java b/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/ForwardAndBackwardCompatibilityTest.java new file mode 100644 index 0000000..fad176e --- /dev/null +++ b/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/ForwardAndBackwardCompatibilityTest.java @@ -0,0 +1,509 @@ +/* + * Copyright 2020-2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.schema.avro; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.TimeUnit; + +import example.avro.v2.User; +import org.apache.avro.Schema; +import org.apache.avro.generic.GenericData; +import org.apache.avro.generic.GenericRecord; +import org.apache.avro.generic.GenericRecordBuilder; +import org.assertj.core.api.Assertions; +import org.junit.Test; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.cloud.schema.registry.client.SchemaRegistryClient; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.cloud.stream.messaging.Source; +import org.springframework.cloud.stream.test.binder.MessageCollector; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.core.io.DefaultResourceLoader; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageDeliveryException; +import org.springframework.messaging.MessagingException; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.util.StringUtils; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * @author Christian Tzolov + */ +public class ForwardAndBackwardCompatibilityTest { + + static SchemaRegistryClient stubSchemaRegistryClient = new StubSchemaRegistryClient(); + + public static final String NO_EXPLICIT_V1_SCHEMA = null; + + public static final String NO_EXPLICIT_V2_SCHEMA = null; + + public static final String NO_READER_SCHEMA = null; + + public static final boolean NO_DYNAMIC_SCHEMA_GENERATION = false; + + public static final boolean ENABLE_DYNAMIC_SCHEMA_GENERATION = true; + + @Test + public void genericRecordBackwardCompatibility() throws Exception { + + Schema s1 = new Schema.Parser().parse( + new DefaultResourceLoader().getResource("classpath:schemas/user.avsc").getInputStream()); + + GenericRecord user1 = new GenericRecordBuilder(s1) + .set("name", "foo" + UUID.randomUUID().toString()) + .set("favoriteColor", "foo" + UUID.randomUUID().toString()) + .set("favoriteNumber", 12) + .build(); + + Schema s2 = new Schema.Parser().parse( + new DefaultResourceLoader().getResource("classpath:schemas/user_v2.avsc").getInputStream()); + + GenericRecord user2 = new GenericRecordBuilder(s2) + .set("name", "foo" + UUID.randomUUID().toString()) + .set("favoriteColor", "foo" + UUID.randomUUID().toString()) + .set("favoriteNumber", 13) + .set("favoritePlace", "Amsterdam") + .build(); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationGenericRecord.class, + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + "classpath:schemas/user_v2.avsc", + NO_DYNAMIC_SCHEMA_GENERATION); + + GenericData.Record resultUser1 = (GenericData.Record) result.get(0); + GenericData.Record resultUser2 = (GenericData.Record) result.get(1); + + assertThat(resultUser1.getSchema()).isEqualTo(s2); + assertThat(resultUser2.getSchema()).isEqualTo(s2); + + assertThat(resultUser1.get("favoriteColor").toString()).isEqualTo(user1.get("favoriteColor").toString()); + assertThat(resultUser1.get("name").toString()).isEqualTo(user1.get("name").toString()); + assertThat(resultUser1.get("favoritePlace").toString()).isEqualTo("NYC"); + assertThat(resultUser1.get("favoriteNumber")).isEqualTo(user1.get("favoriteNumber")); + + assertThat(resultUser2.get("favoriteColor").toString()).isEqualTo(user2.get("favoriteColor").toString()); + assertThat(resultUser2.get("name").toString()).isEqualTo(user2.get("name").toString()); + assertThat(resultUser2.get("favoritePlace").toString()).isEqualTo(user2.get("favoritePlace").toString()); + assertThat(resultUser2.get("favoriteNumber")).isEqualTo(user2.get("favoriteNumber")); + } + + @Test + public void genericRecordForwardCompatibility() throws Exception { + + Schema s1 = new Schema.Parser().parse( + new DefaultResourceLoader().getResource("classpath:schemas/user.avsc").getInputStream()); + + GenericRecord user1 = new GenericRecordBuilder(s1) + .set("name", "foo" + UUID.randomUUID().toString()) + .set("favoriteColor", "foo" + UUID.randomUUID().toString()) + .set("favoriteNumber", 12) + .build(); + + Schema s2 = new Schema.Parser().parse( + new DefaultResourceLoader().getResource("classpath:schemas/user_v2.avsc").getInputStream()); + + GenericRecord user2 = new GenericRecordBuilder(s2) + .set("name", "foo" + UUID.randomUUID().toString()) + .set("favoriteColor", "foo" + UUID.randomUUID().toString()) + .set("favoriteNumber", 13) + .set("favoritePlace", "Amsterdam") + .build(); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationGenericRecord.class, + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + "classpath:schemas/user.avsc", + NO_DYNAMIC_SCHEMA_GENERATION); + + GenericData.Record resultUser1 = (GenericData.Record) result.get(0); + GenericData.Record resultUser2 = (GenericData.Record) result.get(1); + + assertThat(resultUser1.getSchema()).isEqualTo(s1); + assertThat(resultUser2.getSchema()).isEqualTo(s1); + + assertThat(resultUser1.get("favoriteColor").toString()).isEqualTo(user1.get("favoriteColor").toString()); + assertThat(resultUser1.get("name").toString()).isEqualTo(user1.get("name").toString()); + assertThat(resultUser1.get("favoriteNumber")).isEqualTo(user1.get("favoriteNumber")); + + assertThat(resultUser2.get("favoriteColor").toString()).isEqualTo(user2.get("favoriteColor").toString()); + assertThat(resultUser2.get("name").toString()).isEqualTo(user2.get("name").toString()); + assertThat(resultUser2.get("favoriteNumber")).isEqualTo(user2.get("favoriteNumber")); + } + + @Test + public void genericRecordNoReaderSchema() throws Exception { + + Schema s1 = new Schema.Parser().parse( + new DefaultResourceLoader().getResource("classpath:schemas/user.avsc").getInputStream()); + + GenericRecord user1 = new GenericRecordBuilder(s1) + .set("name", "foo" + UUID.randomUUID().toString()) + .set("favoriteColor", "foo" + UUID.randomUUID().toString()) + .set("favoriteNumber", 12) + .build(); + + Schema s2 = new Schema.Parser().parse( + new DefaultResourceLoader().getResource("classpath:schemas/user_v2.avsc").getInputStream()); + + GenericRecord user2 = new GenericRecordBuilder(s2) + .set("name", "foo" + UUID.randomUUID().toString()) + .set("favoriteColor", "foo" + UUID.randomUUID().toString()) + .set("favoriteNumber", 13) + .set("favoritePlace", "Amsterdam") + .build(); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationGenericRecord.class, + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + NO_READER_SCHEMA, + NO_DYNAMIC_SCHEMA_GENERATION); + + GenericData.Record resultUser1 = (GenericData.Record) result.get(0); + GenericData.Record resultUser2 = (GenericData.Record) result.get(1); + + assertThat(resultUser1.getSchema()).isEqualTo(s1); + assertThat(resultUser2.getSchema()).isEqualTo(s2); + + assertThat(resultUser1.get("favoriteColor").toString()).isEqualTo(user1.get("favoriteColor").toString()); + assertThat(resultUser1.get("name").toString()).isEqualTo(user1.get("name").toString()); + assertThat(resultUser1.get("favoriteNumber")).isEqualTo(user1.get("favoriteNumber")); + + assertThat(resultUser2.get("favoriteColor").toString()).isEqualTo(user2.get("favoriteColor").toString()); + assertThat(resultUser2.get("name").toString()).isEqualTo(user2.get("name").toString()); + assertThat(resultUser2.get("favoriteNumber")).isEqualTo(user2.get("favoriteNumber")); + assertThat(resultUser2.get("favoritePlace").toString()).isEqualTo(user2.get("favoritePlace").toString()); + } + + @Test + public void specificRecordBackwardCompatibility() throws Exception { + + example.avro.User user1 = new example.avro.User(); + user1.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user1.setName("foo" + UUID.randomUUID().toString()); + + example.avro.v2.User user2 = new example.avro.v2.User(); + user2.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user2.setName("foo" + UUID.randomUUID().toString()); + user2.setFavoritePlace("Amsterdam"); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationSpecificRecord.class, + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + "classpath:schemas/user_v2.avsc", + NO_DYNAMIC_SCHEMA_GENERATION); + + example.avro.v2.User resultUser1 = (User) result.get(0); + example.avro.v2.User resultUser2 = (User) result.get(1); + + assertThat(resultUser1.getFavoriteColor().toString()).isEqualTo(user1.getFavoriteColor().toString()); + assertThat(resultUser1.getName().toString()).isEqualTo(user1.getName().toString()); + assertThat(resultUser1.getFavoritePlace().toString()).isEqualTo("NYC"); + + assertThat(resultUser2.getFavoriteColor().toString()).isEqualTo(user2.getFavoriteColor().toString()); + assertThat(resultUser2.getName().toString()).isEqualTo(user2.getName().toString()); + assertThat(resultUser2.getFavoritePlace().toString()).isEqualTo("Amsterdam"); + } + + @Test + public void specificRecordForwardCompatibility() throws Exception { + + example.avro.User user1 = new example.avro.User(); + user1.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user1.setName("foo" + UUID.randomUUID().toString()); + + example.avro.v2.User user2 = new example.avro.v2.User(); + user2.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user2.setName("foo" + UUID.randomUUID().toString()); + user2.setFavoritePlace("Amsterdam"); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationSpecificRecord.class, + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + "classpath:schemas/user.avsc", + NO_DYNAMIC_SCHEMA_GENERATION); + + example.avro.User resultUser1 = (example.avro.User) result.get(0); + example.avro.User resultUser2 = (example.avro.User) result.get(1); + + assertThat(resultUser1.getFavoriteColor().toString()).isEqualTo(user1.getFavoriteColor().toString()); + assertThat(resultUser1.getName().toString()).isEqualTo(user1.getName().toString()); + + assertThat(resultUser2.getFavoriteColor().toString()).isEqualTo(user2.getFavoriteColor().toString()); + assertThat(resultUser2.getName().toString()).isEqualTo(user2.getName().toString()); + } + + @Test(expected = MessagingException.class) + public void specificRecordNoReaderSchema() throws Exception { + + example.avro.User user1 = new example.avro.User(); + user1.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user1.setName("foo" + UUID.randomUUID().toString()); + + example.avro.v2.User user2 = new example.avro.v2.User(); + user2.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user2.setName("foo" + UUID.randomUUID().toString()); + user2.setFavoritePlace("Amsterdam"); + + compatibilityTest(user1, user2, AvroSinkApplicationSpecificRecord.class, + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + NO_READER_SCHEMA, + NO_DYNAMIC_SCHEMA_GENERATION); + } + + @Test + public void javaTypeBackwardCompatibility() throws Exception { + org.springframework.cloud.schema.avro.User1 user1 = new org.springframework.cloud.schema.avro.User1(); + user1.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user1.setName("foo" + UUID.randomUUID().toString()); + + org.springframework.cloud.schema.avro.v2.User1 user2 = new org.springframework.cloud.schema.avro.v2.User1(); + user2.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user2.setName("foo" + UUID.randomUUID().toString()); + user2.setFavoritePlace("Amsterdam"); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationUser1V2.class, // Source with User1 v1 payload type + "classpath:schemas/user1_v1.schema", + "classpath:schemas/user1_v2.schema", + NO_READER_SCHEMA, // the readerSchema is IGNORED for java type pojos + NO_DYNAMIC_SCHEMA_GENERATION); + + org.springframework.cloud.schema.avro.v2.User1 resultUser1 = (org.springframework.cloud.schema.avro.v2.User1) result.get(0); + org.springframework.cloud.schema.avro.v2.User1 resultUser2 = (org.springframework.cloud.schema.avro.v2.User1) result.get(1); + + assertThat(resultUser1.getFavoriteColor()).isEqualTo(user1.getFavoriteColor()); + assertThat(resultUser1.getName()).isEqualTo(user1.getName()); + assertThat(resultUser1.getFavoritePlace()).isEqualTo("NYC"); + + assertThat(resultUser2.getFavoriteColor()).isEqualTo(user2.getFavoriteColor()); + assertThat(resultUser2.getName()).isEqualTo(user2.getName()); + assertThat(resultUser2.getFavoritePlace()).isEqualTo("Amsterdam"); + } + + @Test + public void javaTypeForwardCompatibility() throws Exception { + org.springframework.cloud.schema.avro.User1 user1 = new org.springframework.cloud.schema.avro.User1(); + user1.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user1.setName("foo" + UUID.randomUUID().toString()); + + org.springframework.cloud.schema.avro.v2.User1 user2 = new org.springframework.cloud.schema.avro.v2.User1(); + user2.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user2.setName("foo" + UUID.randomUUID().toString()); + user2.setFavoritePlace("Amsterdam"); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationUser1V1.class, // Source with User1 v1 payload type + "classpath:schemas/user1_v1.schema", + "classpath:schemas/user1_v2.schema", + NO_READER_SCHEMA, // the readerSchema is IGNORED for java type pojos + NO_DYNAMIC_SCHEMA_GENERATION); + + org.springframework.cloud.schema.avro.User1 resultUser1 = (org.springframework.cloud.schema.avro.User1) result.get(0); + org.springframework.cloud.schema.avro.User1 resultUser2 = (org.springframework.cloud.schema.avro.User1) result.get(1); + + assertThat(resultUser1.getFavoriteColor()).isEqualTo(user1.getFavoriteColor()); + assertThat(resultUser1.getName()).isEqualTo(user1.getName()); + + assertThat(resultUser2.getFavoriteColor()).isEqualTo(user2.getFavoriteColor()); + assertThat(resultUser2.getName()).isEqualTo(user2.getName()); + } + + @Test(expected = MessageDeliveryException.class) + public void javaTypeWithoutSourceSchemas() throws Exception { + org.springframework.cloud.schema.avro.User1 user1 = new org.springframework.cloud.schema.avro.User1(); + user1.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user1.setName("foo" + UUID.randomUUID().toString()); + + org.springframework.cloud.schema.avro.v2.User1 user2 = new org.springframework.cloud.schema.avro.v2.User1(); + user2.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user2.setName("foo" + UUID.randomUUID().toString()); + user2.setFavoritePlace("Amsterdam"); + + compatibilityTest(user1, user2, + AvroSinkApplicationUser1V1.class, // Source with User1 v1 payload type + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + NO_READER_SCHEMA, // the readerSchema is IGNORED for java type pojos + NO_DYNAMIC_SCHEMA_GENERATION); + } + + @Test + public void javaTypeWithDynamicSchemaGeneration() throws Exception { + org.springframework.cloud.schema.avro.User1 user1 = new org.springframework.cloud.schema.avro.User1(); + user1.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user1.setName("foo" + UUID.randomUUID().toString()); + + org.springframework.cloud.schema.avro.v2.User1 user2 = new org.springframework.cloud.schema.avro.v2.User1(); + user2.setFavoriteColor("foo" + UUID.randomUUID().toString()); + user2.setName("foo" + UUID.randomUUID().toString()); + user2.setFavoritePlace("Amsterdam"); + + List result = compatibilityTest(user1, user2, + AvroSinkApplicationUser1V1.class, // Source with User1 v1 payload type + NO_EXPLICIT_V1_SCHEMA, + NO_EXPLICIT_V2_SCHEMA, + NO_READER_SCHEMA, // the readerSchema is IGNORED for java type pojos + ENABLE_DYNAMIC_SCHEMA_GENERATION); + + org.springframework.cloud.schema.avro.User1 resultUser1 = (org.springframework.cloud.schema.avro.User1) result.get(0); + org.springframework.cloud.schema.avro.User1 resultUser2 = (org.springframework.cloud.schema.avro.User1) result.get(1); + + assertThat(resultUser1.getFavoriteColor()).isEqualTo(user1.getFavoriteColor()); + assertThat(resultUser1.getName()).isEqualTo(user1.getName()); + + assertThat(resultUser2.getFavoriteColor()).isEqualTo(user2.getFavoriteColor()); + assertThat(resultUser2.getName()).isEqualTo(user2.getName()); + } + + public List compatibilityTest( + U1 user1, + U2 user2, + Class sinkApplicationClass, + String user1Schema, + String user2Schema, + String readerSchema, + boolean dynamicSchemaGenerationEnabled) throws Exception { + + List commonSourceArguments = Arrays.asList("--server.port=0", "--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.output.contentType=application/*+avro"); + + // Source 1 + List source1Args = new ArrayList<>(commonSourceArguments); + if (user1Schema != null) { + source1Args.add("--spring.cloud.schema.avro.schema-locations=" + user1Schema); + } + source1Args.add("--spring.cloud.schema.avro.dynamicSchemaGenerationEnabled=" + dynamicSchemaGenerationEnabled); + + ConfigurableApplicationContext sourceContext1 = SpringApplication.run( + AvroSourceApplication.class, source1Args.toArray(new String[source1Args.size()])); + + Source source1 = sourceContext1.getBean(Source.class); + source1.output().send(MessageBuilder.withPayload(user1).build()); + + MessageCollector sourceMessageCollector = sourceContext1.getBean(MessageCollector.class); + Message outboundMessage = sourceMessageCollector.forChannel(source1.output()).poll(1000, TimeUnit.MILLISECONDS); + + // Source2 2 + List source2Args = new ArrayList<>(commonSourceArguments); + if (user2Schema != null) { + source2Args.add("--spring.cloud.schema.avro.schema-locations=" + user2Schema); + } + source2Args.add("--spring.cloud.schema.avro.dynamicSchemaGenerationEnabled=" + dynamicSchemaGenerationEnabled); + + ConfigurableApplicationContext sourceContext2 = SpringApplication.run( + AvroSourceApplication.class, source2Args.toArray(new String[source2Args.size()])); + Source source2 = sourceContext2.getBean(Source.class); + source2.output().send(MessageBuilder.withPayload(user2).build()); + + MessageCollector barSourceMessageCollector = sourceContext2.getBean(MessageCollector.class); + Message barOutboundMessage = barSourceMessageCollector.forChannel(source2.output()).poll(1000, TimeUnit.MILLISECONDS); + + assertThat(barOutboundMessage).isNotNull(); + + // Sink 1 + List sinkArgs = new ArrayList<>(Arrays.asList("--server.port=0", "--spring.jmx.enabled=false")); + if (StringUtils.hasText(readerSchema)) { + sinkArgs.add("--spring.cloud.schema.avro.reader-schema=" + readerSchema); + } + ConfigurableApplicationContext sinkContext = + SpringApplication.run(sinkApplicationClass, sinkArgs.toArray(new String[sinkArgs.size()])); + + Sink sink = sinkContext.getBean(Sink.class); + sink.input().send(outboundMessage); + sink.input().send(barOutboundMessage); + + List receivedPojos = sinkContext.getBean(sinkApplicationClass).getReceivedPojos(); + + Assertions.assertThat(receivedPojos).hasSize(2); + + sourceContext1.close(); + sourceContext2.close(); + sinkContext.close(); + + return receivedPojos; + } + + interface TestSinkApplication { + List getReceivedPojos(); + } + + @EnableBinding(Source.class) + @EnableAutoConfiguration + public static class AvroSourceApplication { + + @Bean + public SchemaRegistryClient schemaRegistryClient() { + return stubSchemaRegistryClient; + } + + } + + @EnableBinding(Sink.class) + @EnableAutoConfiguration + public static class AvroSinkApplication2 implements TestSinkApplication { + + public List receivedPojos = new ArrayList<>(); + + + @StreamListener(Sink.INPUT) + public void listen(T fooPojo) { + this.receivedPojos.add(fooPojo); + } + + @Bean + public SchemaRegistryClient schemaRegistryClient() { + return stubSchemaRegistryClient; + } + + @Override + public List getReceivedPojos() { + return this.receivedPojos; + } + } + + public static class AvroSinkApplicationGenericRecord extends AvroSinkApplication2 { + + } + + public static class AvroSinkApplicationSpecificRecord extends AvroSinkApplication2 { + + } + + public static class AvroSinkApplicationUser1V1 extends AvroSinkApplication2 { + + } + + public static class AvroSinkApplicationUser1V2 extends AvroSinkApplication2 { + + } +} diff --git a/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/v2/User1.java b/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/v2/User1.java new file mode 100644 index 0000000..56c09ee --- /dev/null +++ b/spring-cloud-schema-registry-client/src/test/java/org/springframework/cloud/schema/avro/v2/User1.java @@ -0,0 +1,70 @@ +/* + * Copyright 2016-2019 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.schema.avro.v2; + +import org.apache.avro.reflect.AvroDefault; +import org.apache.avro.reflect.Nullable; + +/** + * @author Marius Bogoevici + */ +public class User1 { + + @Nullable + private String name; + + private int favoriteNumber; + + @Nullable + private String favoriteColor; + + @AvroDefault("\"NYC\"") + private String favoritePlace = "Boston"; + + public String getName() { + return this.name; + } + + public void setName(String name) { + this.name = name; + } + + public int getFavoriteNumber() { + return this.favoriteNumber; + } + + public void setFavoriteNumber(int favoriteNumber) { + this.favoriteNumber = favoriteNumber; + } + + public String getFavoriteColor() { + return this.favoriteColor; + } + + public void setFavoriteColor(String favoriteColor) { + this.favoriteColor = favoriteColor; + } + + public String getFavoritePlace() { + return this.favoritePlace; + } + + public void setFavoritePlace(String favoritePlace) { + this.favoritePlace = favoritePlace; + } + +} diff --git a/spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v1.schema b/spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v1.schema new file mode 100644 index 0000000..40ff89c --- /dev/null +++ b/spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v1.schema @@ -0,0 +1,10 @@ +{"namespace": "org.springframework.cloud.schema.avro", + "type": "record", + "name": "User1", + "fields": [ + {"name": "name", "type": "string"}, + {"name": "favoriteNumber", "type": ["int", "null"]}, + {"name": "favoriteColor", "type": ["string", "null"]} + + ] +} diff --git a/spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v2.schema b/spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v2.schema new file mode 100644 index 0000000..c07c610 --- /dev/null +++ b/spring-cloud-schema-registry-client/src/test/resources/schemas/user1_v2.schema @@ -0,0 +1,10 @@ +{"namespace": "org.springframework.cloud.schema.avro.v2", + "type": "record", + "name": "User1", + "fields": [ + {"name": "name", "type": "string"}, + {"name": "favoriteNumber", "type": ["int", "null"]}, + {"name": "favoriteColor", "type": ["string", "null"]}, + {"name": "favoritePlace", "type": ["string","null"], "default" : "NYC"} + ] +} diff --git a/spring-cloud-schema-registry-client/src/test/resources/schemas/user_v2.avsc b/spring-cloud-schema-registry-client/src/test/resources/schemas/user_v2.avsc new file mode 100644 index 0000000..02f8ac9 --- /dev/null +++ b/spring-cloud-schema-registry-client/src/test/resources/schemas/user_v2.avsc @@ -0,0 +1,10 @@ +{"namespace": "example.avro.v2", + "type": "record", + "name": "User", + "fields": [ + {"name": "name", "type": "string"}, + {"name": "favoriteNumber", "type": ["int", "null"]}, + {"name": "favoriteColor", "type": ["string", "null"]}, + {"name": "favoritePlace", "type": ["string","null"], "default" : "NYC"} + ] +}