|
|
|
|
@@ -49,7 +49,7 @@ import org.springframework.util.ObjectUtils;
|
|
|
|
|
* with the ability to publish and retrieve schemas stored in a schema server, allowing
|
|
|
|
|
* for schema evolution in applications. The supported content types are in the form
|
|
|
|
|
* `application/*+avro`.
|
|
|
|
|
*
|
|
|
|
|
* <p>
|
|
|
|
|
* During the conversion to a message, the converter will set the 'contentType' header to
|
|
|
|
|
* 'application/[prefix].[subject].v[version]+avro', where:
|
|
|
|
|
*
|
|
|
|
|
@@ -65,7 +65,7 @@ import org.springframework.util.ObjectUtils;
|
|
|
|
|
* <i>version</i> is the schema version for the given subject;
|
|
|
|
|
* </ul>
|
|
|
|
|
* </li>
|
|
|
|
|
*
|
|
|
|
|
* <p>
|
|
|
|
|
* When converting from a message, the converter will parse the content-type and use it to
|
|
|
|
|
* fetch and cache the writer schema using the provided {@link SchemaRegistryClient}.
|
|
|
|
|
*
|
|
|
|
|
@@ -76,7 +76,7 @@ import org.springframework.util.ObjectUtils;
|
|
|
|
|
* @author Ish Mahajan
|
|
|
|
|
*/
|
|
|
|
|
public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessageConverter
|
|
|
|
|
implements InitializingBean {
|
|
|
|
|
implements InitializingBean {
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Avro format defined in the Mime type.
|
|
|
|
|
@@ -114,11 +114,11 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
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;
|
|
|
|
|
|
|
|
|
|
protected Resource[] schemaImports = new Resource[] {};
|
|
|
|
|
protected Resource[] schemaImports = new Resource[]{};
|
|
|
|
|
|
|
|
|
|
private Pattern versionedSchema;
|
|
|
|
|
|
|
|
|
|
@@ -136,17 +136,20 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
|
|
|
|
|
private SubjectNamingStrategy subjectNamingStrategy;
|
|
|
|
|
|
|
|
|
|
private boolean ignoreSchemaRegistryServer;
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Creates a new instance, configuring it with {@link SchemaRegistryClient} and
|
|
|
|
|
* {@link CacheManager}.
|
|
|
|
|
*
|
|
|
|
|
* @param schemaRegistryClient the {@link SchemaRegistryClient} used to interact with
|
|
|
|
|
* the schema registry server.
|
|
|
|
|
* @param cacheManager instance of {@link CacheManager} to cache parsed schemas. If
|
|
|
|
|
* caching is not required use {@link NoOpCacheManager}
|
|
|
|
|
* the schema registry server.
|
|
|
|
|
* @param cacheManager instance of {@link CacheManager} to cache parsed schemas. If
|
|
|
|
|
* caching is not required use {@link NoOpCacheManager}
|
|
|
|
|
*/
|
|
|
|
|
@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");
|
|
|
|
|
@@ -157,14 +160,15 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
/**
|
|
|
|
|
* Creates a new instance, configuring it with {@link SchemaRegistryClient} and
|
|
|
|
|
* {@link CacheManager}.
|
|
|
|
|
*
|
|
|
|
|
* @param schemaRegistryClient the {@link SchemaRegistryClient} used to interact with
|
|
|
|
|
* the schema registry server.
|
|
|
|
|
* @param cacheManager instance of {@link CacheManager} to cache parsed schemas. If
|
|
|
|
|
* caching is not required use {@link NoOpCacheManager}
|
|
|
|
|
* @param manager instance of {@link AvroSchemaServiceManager} to manage schemas.
|
|
|
|
|
* the schema registry server.
|
|
|
|
|
* @param cacheManager instance of {@link CacheManager} to cache parsed schemas. If
|
|
|
|
|
* caching is not required use {@link NoOpCacheManager}
|
|
|
|
|
* @param manager instance of {@link AvroSchemaServiceManager} to manage schemas.
|
|
|
|
|
*/
|
|
|
|
|
public AvroSchemaRegistryClientMessageConverter(
|
|
|
|
|
SchemaRegistryClient schemaRegistryClient, CacheManager cacheManager, AvroSchemaServiceManager manager) {
|
|
|
|
|
SchemaRegistryClient schemaRegistryClient, CacheManager cacheManager, AvroSchemaServiceManager manager) {
|
|
|
|
|
super(Collections.singletonList(DEFAULT_AVRO_MIME_TYPE), manager);
|
|
|
|
|
Assert.notNull(schemaRegistryClient, "cannot be null");
|
|
|
|
|
Assert.notNull(cacheManager, "'cacheManager' cannot be null");
|
|
|
|
|
@@ -180,6 +184,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
/**
|
|
|
|
|
* Allows the converter to generate and register schemas automatically. If set to
|
|
|
|
|
* 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) {
|
|
|
|
|
@@ -189,6 +194,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
/**
|
|
|
|
|
* A set of locations where the converter can load schemas from. Schemas provided at
|
|
|
|
|
* these locations will be registered automatically.
|
|
|
|
|
*
|
|
|
|
|
* @param schemaLocations array of locations
|
|
|
|
|
*/
|
|
|
|
|
public void setSchemaLocations(Resource[] schemaLocations) {
|
|
|
|
|
@@ -199,6 +205,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
/**
|
|
|
|
|
* A set of schema locations where should be imported first. Schemas provided at these
|
|
|
|
|
* locations will be reference, thus they should not reference each other.
|
|
|
|
|
*
|
|
|
|
|
* @param schemaImports array of schema imports
|
|
|
|
|
*/
|
|
|
|
|
public void setSchemaImports(Resource[] schemaImports) {
|
|
|
|
|
@@ -207,6 +214,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Set the prefix to be used in the published subtype. Default 'vnd'.
|
|
|
|
|
*
|
|
|
|
|
* @param prefix prefix to be set
|
|
|
|
|
*/
|
|
|
|
|
public void setPrefix(String prefix) {
|
|
|
|
|
@@ -236,40 +244,40 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
@Override
|
|
|
|
|
public void afterPropertiesSet() {
|
|
|
|
|
this.versionedSchema = Pattern.compile("application/" + this.prefix
|
|
|
|
|
+ "\\.([\\p{Alnum}\\$\\.]+)\\.v(\\p{Digit}+)\\+" + AVRO_FORMAT);
|
|
|
|
|
+ "\\.([\\p{Alnum}\\$\\.]+)\\.v(\\p{Digit}+)\\+" + AVRO_FORMAT);
|
|
|
|
|
|
|
|
|
|
Stream.of(this.schemaImports, this.schemaLocations)
|
|
|
|
|
.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");
|
|
|
|
|
.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 -> {
|
|
|
|
|
try {
|
|
|
|
|
Schema schema = parseSchema(resource);
|
|
|
|
|
if (schema.getType().equals(Schema.Type.UNION)) {
|
|
|
|
|
schema.getTypes().forEach(innerSchema -> registerSchema(resource, innerSchema));
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
.flatMap(Arrays::stream)
|
|
|
|
|
.forEach(resource -> {
|
|
|
|
|
try {
|
|
|
|
|
Schema schema = parseSchema(resource);
|
|
|
|
|
if (schema.getType().equals(Schema.Type.UNION)) {
|
|
|
|
|
schema.getTypes().forEach(innerSchema -> registerSchema(resource, innerSchema));
|
|
|
|
|
}
|
|
|
|
|
else {
|
|
|
|
|
registerSchema(resource, schema);
|
|
|
|
|
}
|
|
|
|
|
else {
|
|
|
|
|
registerSchema(resource, schema);
|
|
|
|
|
}
|
|
|
|
|
catch (IOException e) {
|
|
|
|
|
if (this.logger.isWarnEnabled()) {
|
|
|
|
|
this.logger.warn("Failed to parse schema at " + resource.getFilename(), e);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
catch (IOException e) {
|
|
|
|
|
if (this.logger.isWarnEnabled()) {
|
|
|
|
|
this.logger.warn("Failed to parse schema at " + resource.getFilename(), e);
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
if (this.cacheManager instanceof NoOpCacheManager) {
|
|
|
|
|
this.logger.warn("Schema caching is effectively disabled "
|
|
|
|
|
+ "since configured cache manager is a NoOpCacheManager. If this was not "
|
|
|
|
|
+ "the intention, please provide the appropriate instance of CacheManager "
|
|
|
|
|
+ "(i.e., ConcurrentMapCacheManager).");
|
|
|
|
|
+ "since configured cache manager is a NoOpCacheManager. If this was not "
|
|
|
|
|
+ "the intention, please provide the appropriate instance of CacheManager "
|
|
|
|
|
+ "(i.e., ConcurrentMapCacheManager).");
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@@ -294,7 +302,7 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
protected Schema resolveSchemaForWriting(Object payload, MessageHeaders headers,
|
|
|
|
|
MimeType hintedContentType) {
|
|
|
|
|
MimeType hintedContentType) {
|
|
|
|
|
|
|
|
|
|
Schema schema;
|
|
|
|
|
schema = extractSchemaForWriting(payload);
|
|
|
|
|
@@ -304,22 +312,21 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
parsedSchema = new ParsedSchema(schema);
|
|
|
|
|
this.getCache(REFERENCE_CACHE_NAME).putIfAbsent(schema, parsedSchema);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (parsedSchema.getRegistration() == null) {
|
|
|
|
|
if (parsedSchema.getRegistration() == null && !this.ignoreSchemaRegistryServer) {
|
|
|
|
|
SchemaRegistrationResponse response = this.schemaRegistryClient.register(toSubject(this.subjectNamePrefix, schema),
|
|
|
|
|
AVRO_FORMAT, parsedSchema.getRepresentation());
|
|
|
|
|
AVRO_FORMAT, parsedSchema.getRepresentation());
|
|
|
|
|
parsedSchema.setRegistration(response);
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
if (!this.ignoreSchemaRegistryServer) {
|
|
|
|
|
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()
|
|
|
|
|
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);
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
return schema;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@@ -364,8 +371,8 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
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",
|
|
|
|
|
payload.getClass()));
|
|
|
|
|
"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());
|
|
|
|
|
@@ -379,7 +386,7 @@ 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());
|
|
|
|
|
+ schema.getNamespace() + "." + schema.getName());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
this.schemaRegistryClient.register(toSubject(this.subjectNamePrefix, schema), AVRO_FORMAT, schema.toString());
|
|
|
|
|
@@ -405,9 +412,12 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
|
|
|
|
|
private Cache getCache(String name) {
|
|
|
|
|
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.");
|
|
|
|
|
+ this.cacheManager + "'. Typically caches are auto-created by the CacheManagers. "
|
|
|
|
|
+ "Consider reporting it as an issue to the developer of this CacheManager.");
|
|
|
|
|
return cache;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public void setIgnoreSchemaRegistryServer(boolean ignoreSchemaRegistryServer) {
|
|
|
|
|
this.ignoreSchemaRegistryServer = ignoreSchemaRegistryServer;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|