From 986b54f5d09adfb0be33c17f99149c394094f55b Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 24 Feb 2023 17:03:46 -0500 Subject: [PATCH] Schema Registry Testing Improvements - Allow Schema Registry to be optional when running unit tests. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2641 --- .../spring-cloud-stream-schema-registry.adoc | 4 + .../pom.xml | 15 +++ .../sample/consumer/ConsumerApplication.java | 2 +- .../consumer/ConsumerApplicationTests.java | 66 ++++++++++ ...AvroMessageConverterAutoConfiguration.java | 4 + .../avro/AvroMessageConverterProperties.java | 10 ++ ...oSchemaRegistryClientMessageConverter.java | 122 ++++++++++-------- 7 files changed, 166 insertions(+), 57 deletions(-) create mode 100644 samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/test/java/sample/consumer/ConsumerApplicationTests.java diff --git a/docs/src/main/asciidoc/schema-registry/spring-cloud-stream-schema-registry.adoc b/docs/src/main/asciidoc/schema-registry/spring-cloud-stream-schema-registry.adoc index 1ddcf2248..1a2eb6d62 100644 --- a/docs/src/main/asciidoc/schema-registry/spring-cloud-stream-schema-registry.adoc +++ b/docs/src/main/asciidoc/schema-registry/spring-cloud-stream-schema-registry.adoc @@ -127,6 +127,10 @@ Default: `vnd` spring.cloud.stream.schema.avro.subjectNamingStrategy:: Determines the subject name used to register the Avro schema in the schema registry. Two implementations are available, `org.springframework.cloud.stream.schema.avro.DefaultSubjectNamingStrategy`, where the subject is the schema name, and `org.springframework.cloud.stream.schema.avro.QualifiedSubjectNamingStrategy`, which returns a fully qualified subject using the Avro schema namespace and name. Custom strategies can be created by implementing `org.springframework.cloud.stream.schema.avro.SubjectNamingStrategy`. + Default: `org.springframework.cloud.stream.schema.avro.DefaultSubjectNamingStrategy` ++ +spring.cloud.stream.schema.avro.ignoreSchemaRegistryServer:: Ignore any schema registry communication. Useful for testing purposes so that when running a unit test, it does not unnecessarily try to connect to a Schema Registry server. ++ +Default: `false` === Apache Avro Message Converters diff --git a/samples/spring-cloud-stream-schema-registry-integration/pom.xml b/samples/spring-cloud-stream-schema-registry-integration/pom.xml index b75af3400..d6262cc1f 100644 --- a/samples/spring-cloud-stream-schema-registry-integration/pom.xml +++ b/samples/spring-cloud-stream-schema-registry-integration/pom.xml @@ -62,6 +62,21 @@ avro ${avro.version} + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.cloud + spring-cloud-stream-test-binder + test + + + org.awaitility + awaitility + test + diff --git a/samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/main/java/sample/consumer/ConsumerApplication.java b/samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/main/java/sample/consumer/ConsumerApplication.java index e8591d087..5e97d599c 100644 --- a/samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/main/java/sample/consumer/ConsumerApplication.java +++ b/samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/main/java/sample/consumer/ConsumerApplication.java @@ -38,7 +38,7 @@ public class ConsumerApplication { @Bean public Consumer process() { - return input -> logger.info("input: " + input); + return input -> logger.info("[INPUT-RECEIVED]: " + input); } } diff --git a/samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/test/java/sample/consumer/ConsumerApplicationTests.java b/samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/test/java/sample/consumer/ConsumerApplicationTests.java new file mode 100644 index 000000000..c181396c8 --- /dev/null +++ b/samples/spring-cloud-stream-schema-registry-integration/schema-registry-consumer-kafka/src/test/java/sample/consumer/ConsumerApplicationTests.java @@ -0,0 +1,66 @@ +/* + * Copyright 2023 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 sample.consumer; + +import com.example.Sensor; +import org.awaitility.Awaitility; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.boot.test.system.CapturedOutput; +import org.springframework.boot.test.system.OutputCaptureExtension; +import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.cloud.stream.function.StreamBridge; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.util.MimeType; + +import java.time.Duration; +import java.util.Random; +import java.util.UUID; + +@ExtendWith(OutputCaptureExtension.class) +public class ConsumerApplicationTests { + + private static final int AWAIT_DURATION = 10; + + @Test + void test(CapturedOutput output) { + try (ConfigurableApplicationContext context = + new SpringApplicationBuilder( + TestChannelBinderConfiguration + .getCompleteConfiguration(ConsumerApplication.class)) + .web(WebApplicationType.NONE) + .run( + "--spring.cloud.schema.avro.ignore-schema-registry-server=true" + ) + ) { + + StreamBridge streamBridge = context.getBean("streamBridgeUtils", StreamBridge.class); + Random random = new Random(); + Sensor sensor = new Sensor(); + sensor.setId(UUID.randomUUID() + "-v1"); + sensor.setAcceleration(random.nextFloat() * 10); + sensor.setVelocity(random.nextFloat() * 100); + MimeType mimeType = new MimeType("application", "*+avro"); + streamBridge.send("process-in-0", sensor, mimeType); + + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("[INPUT-RECEIVED]: {\"id\":")); + } + } +} diff --git a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterAutoConfiguration.java b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterAutoConfiguration.java index 0ffcb2799..0a72a0273 100644 --- a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterAutoConfiguration.java +++ b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterAutoConfiguration.java @@ -71,6 +71,10 @@ public class AvroMessageConverterAutoConfiguration { } avroSchemaRegistryClientMessageConverter.setPrefix(avroMessageConverterProperties.getPrefix()); + if (avroMessageConverterProperties.isIgnoreSchemaRegistryServer()) { + avroSchemaRegistryClientMessageConverter.setIgnoreSchemaRegistryServer(true); + } + try { Class clazz = avroMessageConverterProperties.getSubjectNamingStrategy(); Constructor constructor = ReflectionUtils.accessibleConstructor(clazz); diff --git a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterProperties.java b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterProperties.java index 237559482..e32d3ac26 100644 --- a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterProperties.java +++ b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroMessageConverterProperties.java @@ -52,6 +52,8 @@ public class AvroMessageConverterProperties { private String subjectNamePrefix; + private boolean ignoreSchemaRegistryServer; + private Class subjectNamingStrategy = DefaultSubjectNamingStrategy.class; public Resource getReaderSchema() { @@ -112,4 +114,12 @@ public class AvroMessageConverterProperties { public void setSubjectNamePrefix(String subjectNamePrefix) { this.subjectNamePrefix = subjectNamePrefix; } + + public boolean isIgnoreSchemaRegistryServer() { + return this.ignoreSchemaRegistryServer; + } + + public void setIgnoreSchemaRegistryServer(boolean ignoreSchemaRegistryServer) { + this.ignoreSchemaRegistryServer = ignoreSchemaRegistryServer; + } } diff --git a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java index 18696db54..ff45b5b35 100644 --- a/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java +++ b/schema-registry/spring-cloud-stream-schema-registry-client/src/main/java/org/springframework/cloud/stream/schema/registry/avro/AvroSchemaRegistryClientMessageConverter.java @@ -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`. - * + *

* 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; * version is the schema version for the given subject; * * - * + *

* 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 _headers = (Map) dfa.getPropertyValue("headers"); - _headers.put(MessageHeaders.CONTENT_TYPE, "application/" + this.prefix + "." + schemaReference.getSubject() + 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); - + } 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; + } }