From 6a9fb6c523289c903bf8a6ae96e011b520fb54a3 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Tue, 6 Feb 2024 11:34:34 -0600 Subject: [PATCH] Polish "Add @PulsarTypeMapping for default topi.." * Introduce a PulsarTypeMappingRegistry to cache annotations * Expand on the topic/schema resolver tests * Add support for KEY_VALUE on annotation --- .../schema-info/custom-schema-mapping.adoc | 17 +++ .../pages/reference/topic-resolution.adoc | 12 ++- .../pulsar/annotation/PulsarTypeMapping.java | 43 ++++++-- .../pulsar/core/DefaultSchemaResolver.java | 56 ++++++++-- .../pulsar/core/DefaultTopicResolver.java | 42 ++++++-- .../core/PulsarTypeMappingRegistry.java | 82 ++++++++++++++ .../core/DefaultSchemaResolverTests.java | 100 +++++++++++++++++- .../core/DefaultTopicResolverTests.java | 44 ++++++++ .../core/PulsarTypeMappingRegistryTests.java | 94 ++++++++++++++++ 9 files changed, 461 insertions(+), 29 deletions(-) create mode 100644 spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTypeMappingRegistry.java create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarTypeMappingRegistryTests.java diff --git a/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/schema-info/custom-schema-mapping.adoc b/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/schema-info/custom-schema-mapping.adoc index c3755707..563b3df2 100644 --- a/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/schema-info/custom-schema-mapping.adoc +++ b/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/schema-info/custom-schema-mapping.adoc @@ -1,3 +1,4 @@ +==== Configuration properties Schema mappings can be configured with the `spring.pulsar.defaults.type-mappings` property. The following example uses `application.yml` to add mappings for the `User` and `Address` complex objects using `AVRO` and `JSON` schemas, respectively: @@ -17,6 +18,7 @@ spring: NOTE: The `message-type` is the fully-qualified name of the message class. +==== Schema resolver customizer The preferred method of adding mappings is via the property mentioned above. However, if more control is needed you can provide a schema resolver customizer to add the mapping(s). @@ -32,3 +34,18 @@ public SchemaResolverCustomizer schemaResolverCustomizer( } } ---- + +==== Type mapping annotation +Another option for specifying default schema information to use for a particular message type is to mark the message class with the `@PulsarTypeMapping` annotation. +The schema info can be specified via the `schemaType` attribute on the annotation. + +The following example configures the system to use JSON as the default schema when producing or consuming messages of type `Foo`: + +[source,java,indent=0,subs="verbatim"] +---- +@PulsarTypeMapping(schemaType = SchemaType.JSON) +record Foo(String value) { +} +---- + +NOTE: The annotations are looked up on-demand and their result is cached. However, there is still a small performance hit on the first lookup. If you want to disable this feature you can invoke the `usePulsarTypeMappingAnnotations(false)` method on the `DefaultSchemaResolver`. diff --git a/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/topic-resolution.adoc b/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/topic-resolution.adoc index f5d82169..b3e77c9d 100644 --- a/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/topic-resolution.adoc +++ b/spring-pulsar-docs/src/main/antora/modules/ROOT/pages/reference/topic-resolution.adoc @@ -35,15 +35,21 @@ NOTE: The `message-type` is the fully-qualified name of the message class. WARNING: If the message (or the first message of a `Publisher` input) is `null`, the framework won't be able to determine the topic from it. Another method shall be used to specify the topic if your application is likely to send `null` messages. === Specified via annotation -When no topic passed into API and no mappings configured, the system looks for `PulsarTopic` annotation. The following example configures topic for `Baz` class using annotation: + +When no topic is passed into the API and there are no custom topic mappings configured, the system looks for a `@PulsarTypeMapping` annotation on the class of the message being produced or consumed. +The default topic can be specified via the `topic` attribute on the annotation. + +The following example configures the default topic to use when producing or consuming messages of type `Foo`: [source,java,indent=0,subs="verbatim"] ---- -@PulsarTopic("baz-topic") -record Baz(String value) { +@PulsarTypeMapping(topic = "foo-topic") +record Foo(String value) { } ---- +NOTE: The annotations are looked up on-demand and their result is cached. However, there is still a small performance hit on the first lookup. If you want to disable this feature you can invoke the `usePulsarTypeMappingAnnotations(false)` method on the `DefaultTopicResolver`. + === Custom topic resolver The preferred method of adding mappings is via the property mentioned above. However, if more control is needed you can replace the default resolver by proving your own implementation, for example: diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarTypeMapping.java b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarTypeMapping.java index 3ec14973..ec189af7 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarTypeMapping.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarTypeMapping.java @@ -16,18 +16,23 @@ package org.springframework.pulsar.annotation; -import org.apache.pulsar.common.schema.SchemaType; - import java.lang.annotation.Documented; import java.lang.annotation.ElementType; import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; +import org.apache.pulsar.common.schema.SchemaType; + /** - * Specifies default topic and schema for class. + * Specifies default topic and schema info for a message class. + *

+ * When a message class is marked with this annotation, the topic/schema resolution + * process will use the specified information to determine a topic/schema to use for the + * message in process. * * @author Aleksei Arsenev + * @author Chris Bono */ @Target(ElementType.TYPE) @Retention(RetentionPolicy.RUNTIME) @@ -35,21 +40,41 @@ import java.lang.annotation.Target; public @interface PulsarTypeMapping { /** - * Default topic for class. - * @return topic + * Default topic for the annotated message class. + * @return default topic for the annotated message class or empty string to indicate + * no default topic is specified */ String topic() default ""; /** - * Default schema type for class. - * @return schema type + * Default schema type to use for the annotated message class. + *

+ * Note that when this is set to {@code KEY_VALUE} you must specify the actual key and + * value information via the {@link #messageKeyType()} and + * {@link #messageValueSchemaType()} attributes, respectively. + * @return schema type to use for the annotated message class or {@code NONE} to + * indicate no default schema is specified */ SchemaType schemaType() default SchemaType.NONE; /** - * Message key type (must be specified when schema type is {@code KEY_VALUE}) - * @return message key type + * The message key type when schema type is set to {@code KEY_VALUE}. + *

+ * When the {@link #schemaType()} is not set to {@code KEY_VALUE} this attribute is + * ignored. + * @return message key type when using {@code KEY_VALUE} schema type */ Class messageKeyType() default Void.class; + /** + * The default schema type to use for the value schema when {@link #schemaType()} is + * set to {@code KEY_VALUE}. + *

+ * When the {@link #schemaType()} is not set to {@code KEY_VALUE} this attribute is + * ignored and the default schema type must be specified via the {@code schemaType} + * attribute. + * @return message value schema type when using {@code KEY_VALUE} schema type + */ + SchemaType messageValueSchemaType() default SchemaType.NONE; + } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultSchemaResolver.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultSchemaResolver.java index 250bbe2d..8271b2d0 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultSchemaResolver.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultSchemaResolver.java @@ -39,10 +39,10 @@ import org.apache.pulsar.common.schema.KeyValueEncodingType; import org.apache.pulsar.common.schema.SchemaType; import org.springframework.core.ResolvableType; -import org.springframework.core.annotation.AnnotationUtils; import org.springframework.core.log.LogAccessor; import org.springframework.lang.Nullable; import org.springframework.pulsar.annotation.PulsarTypeMapping; +import org.springframework.util.Assert; /** * Default schema resolver capable of handling basic message types. @@ -54,6 +54,7 @@ import org.springframework.pulsar.annotation.PulsarTypeMapping; * @author Soby Chacko * @author Alexander Preuß * @author Chris Bono + * @author Aleksei Arsenev */ public class DefaultSchemaResolver implements SchemaResolver { @@ -91,6 +92,20 @@ public class DefaultSchemaResolver implements SchemaResolver { private final Map, Schema> customSchemaMappings = new LinkedHashMap<>(); + private final PulsarTypeMappingRegistry pulsarTypeMappingRegistry = new PulsarTypeMappingRegistry(); + + private boolean usePulsarTypeMappingAnnotations = true; + + /** + * Sets whether to inspect message classes for the + * {@link PulsarTypeMapping @PulsarTypeMapping} annotation during schema resolution. + * @param usePulsarTypeMappingAnnotations whether to inspect messages for the + * annotation + */ + public void usePulsarTypeMappingAnnotations(boolean usePulsarTypeMappingAnnotations) { + this.usePulsarTypeMappingAnnotations = usePulsarTypeMappingAnnotations; + } + /** * Adds a custom mapping from message type to schema. * @param messageType the message type @@ -139,15 +154,18 @@ public class DefaultSchemaResolver implements SchemaResolver { @Nullable protected Schema getCustomSchemaOrMaybeDefault(@Nullable Class messageClass, boolean returnDefault) { + // Check for custom schema mapping Schema schema = this.customSchemaMappings.get(messageClass); - if (schema == null && messageClass != null) { - PulsarTypeMapping annotation = AnnotationUtils.findAnnotation(messageClass, PulsarTypeMapping.class); - if (annotation != null && annotation.schemaType() != SchemaType.NONE) { - var resolvedSchema = resolveSchema(annotation.schemaType(), messageClass, annotation.messageKeyType()); - resolvedSchema.ifResolved(objectSchema -> addCustomSchemaMapping(messageClass, objectSchema)); - schema = resolvedSchema.get().orElse(null); + + // If no custom schema mapping found, look for @PulsarTypeMapping (if enabled) + if (this.usePulsarTypeMappingAnnotations && schema == null && messageClass != null) { + schema = getAnnotatedSchemaType(messageClass); + if (schema != null) { + this.addCustomSchemaMapping(messageClass, schema); } } + + // If still no schema, possibly return a default if (schema == null && returnDefault) { if (messageClass != null) { try { @@ -162,6 +180,30 @@ public class DefaultSchemaResolver implements SchemaResolver { return schema; } + // VisibleForTesting + Schema getAnnotatedSchemaType(Class messageClass) { + PulsarTypeMapping annotation = this.pulsarTypeMappingRegistry.getTypeMappingFor(messageClass).orElse(null); + if (annotation == null || annotation.schemaType() == SchemaType.NONE) { + return null; + } + var schemaType = annotation.schemaType(); + if (schemaType != SchemaType.KEY_VALUE) { + return resolveSchema(annotation.schemaType(), messageClass, null).value().orElse(null); + } + // handle complicated key value + var messageKeyClass = annotation.messageKeyType(); + Assert.state(messageKeyClass != Void.class, + "messageKeyClass can not be Void.class when using KEY_VALUE schema type"); + + var messageValueSchemaType = annotation.messageValueSchemaType(); + Assert.state(messageValueSchemaType != SchemaType.NONE && messageValueSchemaType != SchemaType.KEY_VALUE, + () -> "messageValueSchemaType can not be NONE or KEY_VALUE when using KEY_VALUE schema type"); + + Schema keySchema = this.resolveSchema(messageKeyClass).orElseThrow(); + Schema valueSchema = this.resolveSchema(messageValueSchemaType, messageClass, null).orElseThrow(); + return Schema.KeyValue(keySchema, valueSchema, KeyValueEncodingType.INLINE); + } + @Override @SuppressWarnings("unchecked") public Resolved> resolveSchema(SchemaType schemaType, @Nullable ResolvableType messageType) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultTopicResolver.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultTopicResolver.java index afbacb8a..4c5cdfdc 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultTopicResolver.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultTopicResolver.java @@ -17,11 +17,10 @@ package org.springframework.pulsar.core; import java.util.Collections; +import java.util.LinkedHashMap; import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; import java.util.function.Supplier; -import org.springframework.core.annotation.AnnotationUtils; import org.springframework.lang.Nullable; import org.springframework.pulsar.annotation.PulsarTypeMapping; import org.springframework.util.StringUtils; @@ -34,10 +33,25 @@ import org.springframework.util.StringUtils; * {@link #addCustomTopicMapping(Class, String)}. * * @author Chris Bono + * @author Aleksei Arsenev */ public class DefaultTopicResolver implements TopicResolver { - private final Map, String> customTopicMappings = new ConcurrentHashMap<>(); + private final Map, String> customTopicMappings = new LinkedHashMap<>(); + + private final PulsarTypeMappingRegistry pulsarTypeMappingRegistry = new PulsarTypeMappingRegistry(); + + private boolean usePulsarTypeMappingAnnotations = true; + + /** + * Sets whether to inspect message classes for the + * {@link PulsarTypeMapping @PulsarTypeMapping} annotation during topic resolution. + * @param usePulsarTypeMappingAnnotations whether to inspect messages for the + * annotation + */ + public void usePulsarTypeMappingAnnotations(boolean usePulsarTypeMappingAnnotations) { + this.usePulsarTypeMappingAnnotations = usePulsarTypeMappingAnnotations; + } /** * Adds a custom mapping from message type to topic. @@ -102,16 +116,18 @@ public class DefaultTopicResolver implements TopicResolver { if (messageType == null) { return Resolved.failed("Topic must be specified when the message is null"); } + // Check for custom topic mapping + String topic = this.customTopicMappings.get(messageType); - String topic = this.getCustomTopicMappings().get(messageType); - if (topic == null) { - PulsarTypeMapping annotation = AnnotationUtils.findAnnotation(messageType, PulsarTypeMapping.class); - if (annotation != null && !annotation.topic().isBlank()) { - this.addCustomTopicMapping(messageType, annotation.topic()); - topic = annotation.topic(); + // If no custom topic mapping found, look for @PulsarTypeMapping (if enabled) + if (this.usePulsarTypeMappingAnnotations && topic == null) { + topic = getAnnotatedTopicInfo(messageType); + if (topic != null) { + this.addCustomTopicMapping(messageType, topic); } } + // If still no topic, consult the default topic supplier if (topic == null) { topic = defaultTopicSupplier.get(); } @@ -119,4 +135,12 @@ public class DefaultTopicResolver implements TopicResolver { : Resolved.of(topic); } + // VisibleForTesting + String getAnnotatedTopicInfo(Class messageType) { + return this.pulsarTypeMappingRegistry.getTypeMappingFor(messageType) + .map(PulsarTypeMapping::topic) + .filter(StringUtils::hasText) + .orElse(null); + } + } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTypeMappingRegistry.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTypeMappingRegistry.java new file mode 100644 index 00000000..57465999 --- /dev/null +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTypeMappingRegistry.java @@ -0,0 +1,82 @@ +/* + * Copyright 2023-2024 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.pulsar.core; + +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; + +import org.springframework.core.annotation.AnnotationUtils; +import org.springframework.core.log.LogAccessor; +import org.springframework.pulsar.annotation.PulsarTypeMapping; +import org.springframework.util.Assert; + +/** + * A registry that holds the {@link PulsarTypeMapping @PulsarTypeMapping} annotations and + * each associated class that is marked with the annotation. + *

+ * The annotations are looked up on-demand and the result is cached. + *

+ * Once the cache reaches a {@link #maxNumberOfMappingsCached certain size} (default of + * {@link #DEFAULT_MAX_CACHE_SIZE}) it is cleared and the annotations will be looked up + * again the next time they are requested. + * + * @author Chris Bono + */ +class PulsarTypeMappingRegistry { + + private static final int DEFAULT_MAX_CACHE_SIZE = 1000; + + private final int maxNumberOfMappingsCached; + + private final LogAccessor logger = new LogAccessor(this.getClass()); + + private ConcurrentHashMap, Optional> typeMappingsByClass = new ConcurrentHashMap<>(); + + PulsarTypeMappingRegistry() { + this(DEFAULT_MAX_CACHE_SIZE); + } + + PulsarTypeMappingRegistry(int maxNumberOfMappingsCached) { + Assert.state(maxNumberOfMappingsCached > 0, "maxNumberOfMappingsCached must be > 0"); + this.maxNumberOfMappingsCached = maxNumberOfMappingsCached; + } + + /** + * Gets the {@link PulsarTypeMapping @PulsarTypeMapping} on the specified class or + * empty if the class is not marked with the annotation. + * @param targetClass the class to check for the annotation + * @return an optional containing the annotation or empty if the class is not marked + * with the annotation. + */ + Optional getTypeMappingFor(Class targetClass) { + var optionalTypeMapping = this.typeMappingsByClass.computeIfAbsent(targetClass, this::findTypeMappingOn); + if (this.typeMappingsByClass.size() > this.maxNumberOfMappingsCached) { + this.logger + .info(() -> "Clearing cache - max entries exceeded (%d)".formatted(this.maxNumberOfMappingsCached)); + this.typeMappingsByClass = new ConcurrentHashMap<>(); + } + return optionalTypeMapping; + } + + // VisibleForTesting + protected Optional findTypeMappingOn(Class targetClass) { + this.logger.debug(() -> "Looking for @PulsarTypeMapping on " + targetClass); + PulsarTypeMapping annotation = AnnotationUtils.findAnnotation(targetClass, PulsarTypeMapping.class); + return Optional.ofNullable(annotation); + } + +} diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultSchemaResolverTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultSchemaResolverTests.java index 8c772f8f..311d9142 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultSchemaResolverTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultSchemaResolverTests.java @@ -16,10 +16,15 @@ package org.springframework.pulsar.core; +import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; +import static org.assertj.core.api.Assertions.assertThatIllegalStateException; import static org.assertj.core.api.AssertionsForClassTypes.assertThat; import static org.assertj.core.api.AssertionsForClassTypes.assertThatExceptionOfType; import static org.junit.jupiter.params.provider.Arguments.arguments; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import java.nio.ByteBuffer; import java.sql.Time; @@ -57,6 +62,7 @@ import org.springframework.pulsar.listener.Proto.Person; * Unit tests for {@link DefaultSchemaResolver}. * * @author Chris Bono + * @author Aleksei Arsenev */ class DefaultSchemaResolverTests { @@ -202,11 +208,11 @@ class DefaultSchemaResolverTests { assertThat(resolver.resolveSchema(Bar.class, true).orElseThrow()).isEqualTo(Schema.BYTES); } - @Test void annotatedMessageType() { assertThat(resolver.resolveSchema(Zaz.class, false).orElseThrow()).isEqualTo(Schema.STRING); } + } @Nested @@ -368,6 +374,98 @@ class DefaultSchemaResolverTests { } + @Nested + class SchemaByAnnotatedMessageType { + + @Test + void annotatedMessageType() { + resolver = spy(resolver); + var resolvedSchema = resolver.resolveSchema(JsonMsgType.class, false).orElseThrow(); + assertThat(resolvedSchema).isInstanceOf(JSONSchema.class) + .extracting("schema.fullName") + .asString() + .endsWith(JsonMsgType.class.getSimpleName()); + + // verify added to custom mappings + assertThat(resolver.getCustomSchemaMappings().get(JsonMsgType.class)).isSameAs(resolvedSchema); + + // verify subsequent calls skip resolution again + assertThat(resolver.resolveSchema(JsonMsgType.class, false).orElseThrow()).isSameAs(resolvedSchema); + verify(resolver, times(1)).getAnnotatedSchemaType(JsonMsgType.class); + } + + @Test + void annotatedMessageTypeKeyValue() { + assertThat(resolver.resolveSchema(KeyValueMsgType.class, false).orElseThrow()) + .asInstanceOf(InstanceOfAssertFactories.type(KeyValueSchema.class)) + .satisfies((keyValueSchema -> { + assertThat(keyValueSchema.getKeySchema()).isEqualTo(Schema.STRING); + assertThat(keyValueSchema.getValueSchema()).isInstanceOf(JSONSchema.class) + .extracting("schema.fullName") + .asString() + .endsWith(KeyValueMsgType.class.getSimpleName()); + assertThat(keyValueSchema.getKeyValueEncodingType()).isEqualTo(KeyValueEncodingType.INLINE); + })); + } + + @Test + void annotatedMessageTypeKeyValueMissingKeyInfo() { + assertThatIllegalStateException() + .isThrownBy(() -> resolver.resolveSchema(KeyValueMsgTypeNoKeyInfo.class, false).orElseThrow()) + .withMessage("messageKeyClass can not be Void.class when using KEY_VALUE schema type"); + } + + @Test + void annotatedMessageTypeKeyValueMissingValueInfo() { + assertThatIllegalStateException() + .isThrownBy(() -> resolver.resolveSchema(KeyValueMsgTypeNoValueInfo.class, false).orElseThrow()) + .withMessage("messageValueSchemaType can not be NONE or KEY_VALUE when using KEY_VALUE schema type"); + } + + @Test + void annotatedMessageTypeNoSchemaInfo() { + assertThatIllegalArgumentException() + .isThrownBy(() -> resolver.resolveSchema(NoSchemaInfoMsgType.class, false).orElseThrow()) + .withMessage("Schema not specified and no schema found for " + NoSchemaInfoMsgType.class); + } + + @Test + void annotationMappingIgnoredWhenFeatureDisabled() { + resolver.usePulsarTypeMappingAnnotations(false); + assertThatIllegalArgumentException() + .isThrownBy(() -> resolver.resolveSchema(JsonMsgType.class, false).orElseThrow()) + .withMessage("Schema not specified and no schema found for " + JsonMsgType.class); + } + + @Test + void customMappingTakesPrecedenceOverAnnotationMapping() { + resolver.addCustomSchemaMapping(JsonMsgType.class, Schema.STRING); + assertThat(resolver.resolveSchema(JsonMsgType.class, false).orElseThrow()).isEqualTo(Schema.STRING); + } + + @PulsarTypeMapping(schemaType = SchemaType.JSON) + record JsonMsgType(String value) { + } + + @PulsarTypeMapping(schemaType = SchemaType.KEY_VALUE, messageKeyType = String.class, + messageValueSchemaType = SchemaType.JSON) + record KeyValueMsgType(String key) { + } + + @PulsarTypeMapping(schemaType = SchemaType.KEY_VALUE, messageValueSchemaType = SchemaType.JSON) + record KeyValueMsgTypeNoKeyInfo(String key) { + } + + @PulsarTypeMapping(schemaType = SchemaType.KEY_VALUE, messageKeyType = String.class) + record KeyValueMsgTypeNoValueInfo(String key) { + } + + @PulsarTypeMapping(topic = "ignore-topic") + record NoSchemaInfoMsgType(String value) { + } + + } + record Foo(String value) { } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultTopicResolverTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultTopicResolverTests.java index 1c3ba11f..ad50644f 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultTopicResolverTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultTopicResolverTests.java @@ -18,9 +18,13 @@ package org.springframework.pulsar.core; import static org.assertj.core.api.AssertionsForClassTypes.assertThat; import static org.junit.jupiter.params.provider.Arguments.arguments; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import java.util.stream.Stream; +import org.apache.pulsar.common.schema.SchemaType; import org.assertj.core.api.InstanceOfAssertFactories; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Nested; @@ -36,6 +40,7 @@ import org.springframework.pulsar.annotation.PulsarTypeMapping; * Unit tests for {@link DefaultTopicResolver}. * * @author Chris Bono + * @author Aleksei Arsenev */ class DefaultTopicResolverTests { @@ -116,6 +121,7 @@ class DefaultTopicResolverTests { arguments("nullMessageWithUserTopicAndDefault", userTopic, null, defaultTopic, userTopic), arguments("annotationMessageWithUserTopic", userTopic, Baz.class, defaultTopic, userTopic), arguments("annotationMessageNoUserTopic", null, Baz.class, defaultTopic, bazTopic), + arguments("annotationMessageNoTopicInfo", null, BazNoTopicInfo.class, defaultTopic, defaultTopic), arguments("nullMessageWithDefault", null, null, defaultTopic, null), arguments("noMatchWithUserTopicAndDefault", userTopic, Bar.class, defaultTopic, userTopic), arguments("noMatchWithUserTopic", userTopic, Bar.class, null, userTopic), @@ -125,6 +131,40 @@ class DefaultTopicResolverTests { // @formatter:on } + @Nested + class TopicByAnnotatedMessageType { + + @Test + void customMappingTakesPrecedenceOverAnnotationMapping() { + assertThat(resolver.resolveTopic(null, Baz.class, () -> defaultTopic).value().orElse(null)) + .isEqualTo(bazTopic); + resolver.addCustomTopicMapping(Baz.class, "baz-custom-topic"); + assertThat(resolver.resolveTopic(null, Baz.class, () -> defaultTopic).value().orElse(null)) + .isEqualTo("baz-custom-topic"); + } + + @Test + void annotationMappingIgnoredWhenFeatureDisabled() { + resolver.usePulsarTypeMappingAnnotations(false); + assertThat(resolver.resolveTopic(null, Baz.class, () -> defaultTopic).value().orElse(null)) + .isEqualTo(defaultTopic); + } + + @Test + void annotatedMessageTypeWithTopicInfo() { + resolver = spy(resolver); + assertThat(resolver.resolveTopic(null, Baz.class, () -> defaultTopic).value().orElse(null)) + .isEqualTo(bazTopic); + // verify added to custom mappings + assertThat(resolver.getCustomTopicMappings().get(Baz.class)).isEqualTo(bazTopic); + // verify subsequent calls skip resolution again + assertThat(resolver.resolveTopic(null, Baz.class, () -> defaultTopic).value().orElse(null)) + .isEqualTo(bazTopic); + verify(resolver, times(1)).getAnnotatedTopicInfo(Baz.class); + } + + } + @Nested class TopicMappingsAPI { @@ -175,4 +215,8 @@ class DefaultTopicResolverTests { record Baz(String value) { } + @PulsarTypeMapping(schemaType = SchemaType.STRING) + record BazNoTopicInfo(String value) { + } + } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarTypeMappingRegistryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarTypeMappingRegistryTests.java new file mode 100644 index 00000000..4f04e713 --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarTypeMappingRegistryTests.java @@ -0,0 +1,94 @@ +/* + * Copyright 2023-2024 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.pulsar.core; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatIllegalStateException; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import org.springframework.pulsar.annotation.PulsarTypeMapping; + +/** + * Unit tests for {@link PulsarTypeMappingRegistry}. + * + * @author Chris Bono + */ +class PulsarTypeMappingRegistryTests { + + @Test + void typeMappingFoundAndCached() { + PulsarTypeMappingRegistry registry = spy(new PulsarTypeMappingRegistry(2)); + var typeMapping = registry.getTypeMappingFor(Foo.class); + assertThat(typeMapping).map(PulsarTypeMapping::topic).hasValue("foo-topic"); + // subsequent calls are cached + assertThat(registry.getTypeMappingFor(Foo.class)).isSameAs(typeMapping); + assertThat(registry.getTypeMappingFor(Foo.class)).isSameAs(typeMapping); + verify(registry, times(1)).findTypeMappingOn(Foo.class); + } + + @Test + void typeMappingNotFoundAndCached() { + PulsarTypeMappingRegistry registry = spy(new PulsarTypeMappingRegistry(2)); + var typeMapping = registry.getTypeMappingFor(Bar.class); + assertThat(typeMapping).isEmpty(); + // subsequent calls are cached + assertThat(registry.getTypeMappingFor(Bar.class)).isSameAs(typeMapping); + assertThat(registry.getTypeMappingFor(Bar.class)).isSameAs(typeMapping); + verify(registry, times(1)).findTypeMappingOn(Bar.class); + } + + @Test + void cacheIsClearedOnceMaxNumberReached() { + PulsarTypeMappingRegistry registry = spy(new PulsarTypeMappingRegistry(2)); + registry.getTypeMappingFor(Foo.class); + registry.getTypeMappingFor(Bar.class); + registry.getTypeMappingFor(Zaa.class); + // the 3rd request will force a cache clear - subsequent calls will do lookup + // again + registry.getTypeMappingFor(Foo.class); + registry.getTypeMappingFor(Bar.class); + // now values should be cached again + registry.getTypeMappingFor(Foo.class); + registry.getTypeMappingFor(Bar.class); + verify(registry, times(2)).findTypeMappingOn(Foo.class); + verify(registry, times(2)).findTypeMappingOn(Bar.class); + } + + @ParameterizedTest + @ValueSource(ints = { -1, 0 }) + void maxNumberOfMappingsMustBePositive(int maxNumberOfMappings) { + assertThatIllegalStateException().isThrownBy(() -> new PulsarTypeMappingRegistry(maxNumberOfMappings)) + .withMessage("maxNumberOfMappingsCached must be > 0"); + } + + @PulsarTypeMapping(topic = "foo-topic") + record Foo(String value) { + } + + record Bar(String value) { + } + + record Zaa(String value) { + } + +}