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
This commit is contained in:
@@ -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<Object> writer = avroSchemaServiceManager()
|
||||
.getDatumWriter(payload.getClass(), schema);
|
||||
DatumWriter<Object> 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);
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<? extends SubjectNamingStrategy> subjectNamingStrategy) {
|
||||
public void setSubjectNamingStrategy(Class<? extends SubjectNamingStrategy> subjectNamingStrategy) {
|
||||
Assert.notNull(subjectNamingStrategy, "cannot be null");
|
||||
this.subjectNamingStrategy = subjectNamingStrategy;
|
||||
}
|
||||
|
||||
@@ -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<String, Object> _headers = (Map<String, Object>) dfa
|
||||
.getPropertyValue("headers");
|
||||
_headers.put(MessageHeaders.CONTENT_TYPE,
|
||||
"application/" + this.prefix + "." + schemaReference.getSubject() + ".v"
|
||||
+ schemaReference.getVersion() + "+" + AVRO_FORMAT);
|
||||
Map<String, Object> _headers = (Map<String, Object>) 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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<Object> getDatumReader(Class<?> type, Schema schema, Schema writerSchema) {
|
||||
public DatumReader<Object> getDatumReader(Class<?> type, Schema readerSchema, Schema writerSchema) {
|
||||
DatumReader<Object> 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<? extends Object> clazz, byte[] payload, Schema readerSchema,
|
||||
Schema writerSchema) throws IOException {
|
||||
DatumReader<Object> reader = this.getDatumReader(clazz,
|
||||
readerSchema, writerSchema);
|
||||
public Object readData(Class<? extends Object> clazz, byte[] payload, Schema readerSchema, Schema writerSchema)
|
||||
throws IOException {
|
||||
DatumReader<Object> reader = this.getDatumReader(clazz, readerSchema, writerSchema);
|
||||
Decoder decoder = DecoderFactory.get().binaryDecoder(payload, null);
|
||||
return reader.read(null, decoder);
|
||||
}
|
||||
|
||||
@@ -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<? extends Object> 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]) {
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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 <U1, U2, S extends TestSinkApplication> List<?> compatibilityTest(
|
||||
U1 user1,
|
||||
U2 user2,
|
||||
Class<S> sinkApplicationClass,
|
||||
String user1Schema,
|
||||
String user2Schema,
|
||||
String readerSchema,
|
||||
boolean dynamicSchemaGenerationEnabled) throws Exception {
|
||||
|
||||
List<String> commonSourceArguments = Arrays.asList("--server.port=0", "--spring.jmx.enabled=false",
|
||||
"--spring.cloud.stream.bindings.output.contentType=application/*+avro");
|
||||
|
||||
// Source 1
|
||||
List<String> 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<String> 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<String> 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<T> implements TestSinkApplication {
|
||||
|
||||
public List<T> 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<GenericRecord> {
|
||||
|
||||
}
|
||||
|
||||
public static class AvroSinkApplicationSpecificRecord extends AvroSinkApplication2<org.apache.avro.specific.SpecificRecord> {
|
||||
|
||||
}
|
||||
|
||||
public static class AvroSinkApplicationUser1V1 extends AvroSinkApplication2<org.springframework.cloud.schema.avro.User1> {
|
||||
|
||||
}
|
||||
|
||||
public static class AvroSinkApplicationUser1V2 extends AvroSinkApplication2<org.springframework.cloud.schema.avro.v2.User1> {
|
||||
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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"]}
|
||||
|
||||
]
|
||||
}
|
||||
@@ -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"}
|
||||
]
|
||||
}
|
||||
@@ -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"}
|
||||
]
|
||||
}
|
||||
Reference in New Issue
Block a user