diff --git a/spring-cloud-stream-schema/pom.xml b/spring-cloud-stream-schema/pom.xml
index 84f04ba5e..fee7e8100 100644
--- a/spring-cloud-stream-schema/pom.xml
+++ b/spring-cloud-stream-schema/pom.xml
@@ -39,6 +39,10 @@
${avro.version}
true
+
+ org.projectlombok
+ lombok
+
org.springframework.cloud
spring-cloud-stream-test-support
@@ -54,6 +58,11 @@
spring-cloud-stream-schema-server
test
+
+ com.fasterxml.jackson.dataformat
+ jackson-dataformat-avro
+ test
+
diff --git a/spring-cloud-stream-schema/src/main/java/org/springframework/cloud/stream/schema/avro/AbstractAvroMessageConverter.java b/spring-cloud-stream-schema/src/main/java/org/springframework/cloud/stream/schema/avro/AbstractAvroMessageConverter.java
index 2ef6220c8..65cb102fe 100644
--- a/spring-cloud-stream-schema/src/main/java/org/springframework/cloud/stream/schema/avro/AbstractAvroMessageConverter.java
+++ b/spring-cloud-stream-schema/src/main/java/org/springframework/cloud/stream/schema/avro/AbstractAvroMessageConverter.java
@@ -22,20 +22,9 @@ import java.util.Collection;
import java.util.Collections;
import org.apache.avro.Schema;
-import org.apache.avro.generic.GenericDatumReader;
-import org.apache.avro.generic.GenericDatumWriter;
-import org.apache.avro.generic.GenericRecord;
-import org.apache.avro.io.DatumReader;
import org.apache.avro.io.DatumWriter;
-import org.apache.avro.io.Decoder;
-import org.apache.avro.io.DecoderFactory;
import org.apache.avro.io.Encoder;
import org.apache.avro.io.EncoderFactory;
-import org.apache.avro.reflect.ReflectDatumReader;
-import org.apache.avro.reflect.ReflectDatumWriter;
-import org.apache.avro.specific.SpecificDatumReader;
-import org.apache.avro.specific.SpecificDatumWriter;
-import org.apache.avro.specific.SpecificRecord;
import org.springframework.core.io.Resource;
import org.springframework.messaging.Message;
@@ -58,14 +47,31 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver
* common parser will let user to import external schemas.
*/
private Schema.Parser schemaParser = new Schema.Parser();
+ private AvroSchemaServiceManager avroSchemaServiceManager;
+ @Deprecated
protected AbstractAvroMessageConverter(MimeType supportedMimeType) {
- this(Collections.singletonList(supportedMimeType));
+ this(Collections.singletonList(supportedMimeType), new AvroSchemaServiceManagerImpl());
}
+ protected AbstractAvroMessageConverter(MimeType supportedMimeType, AvroSchemaServiceManager avroSchemaServiceManager) {
+ this(Collections.singletonList(supportedMimeType), avroSchemaServiceManager);
+ }
+
+ @Deprecated
protected AbstractAvroMessageConverter(Collection supportedMimeTypes) {
+ this(supportedMimeTypes, new AvroSchemaServiceManagerImpl());
+ setContentTypeResolver(new OriginalContentTypeResolver());
+ }
+
+ protected AbstractAvroMessageConverter(Collection supportedMimeTypes, AvroSchemaServiceManager manager) {
super(supportedMimeTypes);
setContentTypeResolver(new OriginalContentTypeResolver());
+ this.avroSchemaServiceManager = manager;
+ }
+
+ protected AvroSchemaServiceManager avroSchemaServiceManager() {
+ return this.avroSchemaServiceManager;
}
protected Schema parseSchema(Resource r) throws IOException {
@@ -81,7 +87,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver
@Override
protected Object convertFromInternal(Message> message, Class> targetClass,
Object conversionHint) {
- Object result = null;
+ Object result;
try {
byte[] payload = (byte[]) message.getPayload();
@@ -98,11 +104,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver
Schema writerSchema = resolveWriterSchemaForDeserialization(mimeType);
Schema readerSchema = resolveReaderSchemaForDeserialization(targetClass);
- @SuppressWarnings("unchecked")
- DatumReader