From 98da3aa27c0661c9607795cdef2dcf472527e4d3 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Tue, 13 Aug 2024 19:53:21 -0500 Subject: [PATCH 1/2] Add support for Pulsar default tenant/namespace This commit allows Pulsar users to configure a default tenant and/or namespace to be used when producing or consuming messages to topic URLs that are not fully-qualified. See gh-41851 --- .../pulsar/PulsarAutoConfiguration.java | 34 +++++++++----- .../pulsar/PulsarConfiguration.java | 12 +++++ .../pulsar/PulsarProperties.java | 44 +++++++++++++++++++ .../PulsarReactiveAutoConfiguration.java | 21 ++++++--- ...itional-spring-configuration-metadata.json | 6 +++ .../pulsar/PulsarAutoConfigurationTests.java | 12 +++-- .../pulsar/PulsarConfigurationTests.java | 41 +++++++++++++++++ .../pulsar/PulsarPropertiesTests.java | 25 ++++++++++- .../PulsarReactiveAutoConfigurationTests.java | 16 +++++++ 9 files changed, 191 insertions(+), 20 deletions(-) diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java index 92c66e2168..7d466d5719 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java @@ -52,6 +52,7 @@ import org.springframework.pulsar.core.PulsarConsumerFactory; import org.springframework.pulsar.core.PulsarProducerFactory; import org.springframework.pulsar.core.PulsarReaderFactory; import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.PulsarTopicBuilder; import org.springframework.pulsar.core.ReaderBuilderCustomizer; import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.core.TopicResolver; @@ -88,24 +89,29 @@ public class PulsarAutoConfiguration { @ConditionalOnMissingBean(PulsarProducerFactory.class) @ConditionalOnProperty(name = "spring.pulsar.producer.cache.enabled", havingValue = "false") DefaultPulsarProducerFactory pulsarProducerFactory(PulsarClient pulsarClient, TopicResolver topicResolver, - ObjectProvider> customizersProvider) { + ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { List> lambdaSafeCustomizers = lambdaSafeProducerBuilderCustomizers( customizersProvider); - return new DefaultPulsarProducerFactory<>(pulsarClient, this.properties.getProducer().getTopicName(), - lambdaSafeCustomizers, topicResolver); + DefaultPulsarProducerFactory producerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + this.properties.getProducer().getTopicName(), lambdaSafeCustomizers, topicResolver); + producerFactory.setTopicBuilder(topicBuilder); + return producerFactory; } @Bean @ConditionalOnMissingBean(PulsarProducerFactory.class) @ConditionalOnProperty(name = "spring.pulsar.producer.cache.enabled", havingValue = "true", matchIfMissing = true) CachingPulsarProducerFactory cachingPulsarProducerFactory(PulsarClient pulsarClient, TopicResolver topicResolver, - ObjectProvider> customizersProvider) { + ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { PulsarProperties.Producer.Cache cacheProperties = this.properties.getProducer().getCache(); List> lambdaSafeCustomizers = lambdaSafeProducerBuilderCustomizers( customizersProvider); - return new CachingPulsarProducerFactory<>(pulsarClient, this.properties.getProducer().getTopicName(), - lambdaSafeCustomizers, topicResolver, cacheProperties.getExpireAfterAccess(), - cacheProperties.getMaximumSize(), cacheProperties.getInitialCapacity()); + CachingPulsarProducerFactory producerFactory = new CachingPulsarProducerFactory<>(pulsarClient, + this.properties.getProducer().getTopicName(), lambdaSafeCustomizers, topicResolver, + cacheProperties.getExpireAfterAccess(), cacheProperties.getMaximumSize(), + cacheProperties.getInitialCapacity()); + producerFactory.setTopicBuilder(topicBuilder); + return producerFactory; } private List> lambdaSafeProducerBuilderCustomizers( @@ -138,13 +144,16 @@ public class PulsarAutoConfiguration { @Bean @ConditionalOnMissingBean(PulsarConsumerFactory.class) DefaultPulsarConsumerFactory pulsarConsumerFactory(PulsarClient pulsarClient, - ObjectProvider> customizersProvider) { + ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeConsumerBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); List> lambdaSafeCustomizers = List .of((builder) -> applyConsumerBuilderCustomizers(customizers, builder)); - return new DefaultPulsarConsumerFactory<>(pulsarClient, lambdaSafeCustomizers); + DefaultPulsarConsumerFactory consumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + lambdaSafeCustomizers); + consumerFactory.setTopicBuilder(topicBuilder); + return consumerFactory; } @Bean @@ -181,13 +190,16 @@ public class PulsarAutoConfiguration { @Bean @ConditionalOnMissingBean(PulsarReaderFactory.class) DefaultPulsarReaderFactory pulsarReaderFactory(PulsarClient pulsarClient, - ObjectProvider> customizersProvider) { + ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeReaderBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); List> lambdaSafeCustomizers = List .of((builder) -> applyReaderBuilderCustomizers(customizers, builder)); - return new DefaultPulsarReaderFactory<>(pulsarClient, lambdaSafeCustomizers); + DefaultPulsarReaderFactory readerFactory = new DefaultPulsarReaderFactory<>(pulsarClient, + lambdaSafeCustomizers); + readerFactory.setTopicBuilder(topicBuilder); + return readerFactory; } @SuppressWarnings("unchecked") diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfiguration.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfiguration.java index 0f5b860490..ea60717f7b 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfiguration.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfiguration.java @@ -23,6 +23,7 @@ import org.apache.pulsar.client.admin.PulsarAdminBuilder; import org.apache.pulsar.client.api.ClientBuilder; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.common.naming.TopicDomain; import org.apache.pulsar.common.schema.SchemaType; import org.springframework.beans.factory.ObjectProvider; @@ -34,6 +35,7 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.boot.util.LambdaSafe; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; import org.springframework.pulsar.core.DefaultPulsarClientFactory; import org.springframework.pulsar.core.DefaultSchemaResolver; import org.springframework.pulsar.core.DefaultTopicResolver; @@ -41,6 +43,7 @@ import org.springframework.pulsar.core.PulsarAdminBuilderCustomizer; import org.springframework.pulsar.core.PulsarAdministration; import org.springframework.pulsar.core.PulsarClientBuilderCustomizer; import org.springframework.pulsar.core.PulsarClientFactory; +import org.springframework.pulsar.core.PulsarTopicBuilder; import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer; import org.springframework.pulsar.core.TopicResolver; @@ -176,4 +179,13 @@ class PulsarConfiguration { properties.isFailFast(), properties.isPropagateFailures(), properties.isPropagateStopFailures()); } + @Bean + @Scope("prototype") + @ConditionalOnMissingBean + @ConditionalOnProperty(name = "spring.pulsar.defaults.topic.enabled", havingValue = "true", matchIfMissing = true) + PulsarTopicBuilder pulsarTopicBuilder() { + return new PulsarTopicBuilder(TopicDomain.persistent, this.properties.getDefaults().getTopic().getTenant(), + this.properties.getDefaults().getTopic().getNamespace()); + } + } diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java index 458aebb814..e4f6897f0b 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java @@ -257,6 +257,16 @@ public class PulsarProperties { */ private List typeMappings = new ArrayList<>(); + private Topic topic = new Topic(); + + public Topic getTopic() { + return this.topic; + } + + public void setTopic(Topic topic) { + this.topic = topic; + } + public List getTypeMappings() { return this.typeMappings; } @@ -301,6 +311,40 @@ public class PulsarProperties { } + public static class Topic { + + /** + * Default tenant to use when producing or consuming messages against a + * non-fully-qualified topic URL. When not specified Pulsar uses a default + * tenant of 'public'. + */ + private String tenant; + + /** + * Default namespace to use when producing or consuming messages against a + * non-fully-qualified topic URL. When not specified Pulsar uses a default + * namespace of 'default'. + */ + private String namespace; + + public String getTenant() { + return this.tenant; + } + + public void setTenant(String tenant) { + this.tenant = tenant; + } + + public String getNamespace() { + return this.namespace; + } + + public void setNamespace(String namespace) { + this.namespace = namespace; + } + + } + } public static class Function { diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java index 4c2aeb172d..6b983e9329 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java @@ -41,6 +41,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.pulsar.config.PulsarAnnotationSupportBeanNames; +import org.springframework.pulsar.core.PulsarTopicBuilder; import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.core.TopicResolver; import org.springframework.pulsar.reactive.config.DefaultReactivePulsarListenerContainerFactory; @@ -112,7 +113,8 @@ public class PulsarReactiveAutoConfiguration { @ConditionalOnMissingBean(ReactivePulsarSenderFactory.class) DefaultReactivePulsarSenderFactory reactivePulsarSenderFactory(ReactivePulsarClient reactivePulsarClient, ObjectProvider reactiveMessageSenderCache, TopicResolver topicResolver, - ObjectProvider> customizersProvider) { + ObjectProvider> customizersProvider, + PulsarTopicBuilder topicBuilder) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeMessageSenderBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); @@ -122,6 +124,7 @@ public class PulsarReactiveAutoConfiguration { .withDefaultConfigCustomizers(lambdaSafeCustomizers) .withMessageSenderCache(reactiveMessageSenderCache.getIfAvailable()) .withTopicResolver(topicResolver) + .withTopicBuilder(topicBuilder) .build(); } @@ -136,13 +139,17 @@ public class PulsarReactiveAutoConfiguration { @ConditionalOnMissingBean(ReactivePulsarConsumerFactory.class) DefaultReactivePulsarConsumerFactory reactivePulsarConsumerFactory( ReactivePulsarClient pulsarReactivePulsarClient, - ObjectProvider> customizersProvider) { + ObjectProvider> customizersProvider, + PulsarTopicBuilder topicBuilder) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeMessageConsumerBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); List> lambdaSafeCustomizers = List .of((builder) -> applyMessageConsumerBuilderCustomizers(customizers, builder)); - return new DefaultReactivePulsarConsumerFactory<>(pulsarReactivePulsarClient, lambdaSafeCustomizers); + DefaultReactivePulsarConsumerFactory consumerFactory = new DefaultReactivePulsarConsumerFactory<>( + pulsarReactivePulsarClient, lambdaSafeCustomizers); + consumerFactory.setTopicBuilder(topicBuilder); + return consumerFactory; } @SuppressWarnings("unchecked") @@ -167,13 +174,17 @@ public class PulsarReactiveAutoConfiguration { @Bean @ConditionalOnMissingBean(ReactivePulsarReaderFactory.class) DefaultReactivePulsarReaderFactory reactivePulsarReaderFactory(ReactivePulsarClient reactivePulsarClient, - ObjectProvider> customizersProvider) { + ObjectProvider> customizersProvider, + PulsarTopicBuilder topicBuilder) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeMessageReaderBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); List> lambdaSafeCustomizers = List .of((builder) -> applyMessageReaderBuilderCustomizers(customizers, builder)); - return new DefaultReactivePulsarReaderFactory<>(reactivePulsarClient, lambdaSafeCustomizers); + DefaultReactivePulsarReaderFactory readerFactory = new DefaultReactivePulsarReaderFactory<>( + reactivePulsarClient, lambdaSafeCustomizers); + readerFactory.setTopicBuilder(topicBuilder); + return readerFactory; } @SuppressWarnings("unchecked") diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/spring-boot-project/spring-boot-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json index 8683a4ce05..63dae453db 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -2068,6 +2068,12 @@ "name": "spring.neo4j.uri", "defaultValue": "bolt://localhost:7687" }, + { + "name": "spring.pulsar.defaults.topic.enabled", + "type": "java.lang.Boolean", + "description": "Whether to enable default tenant and namespace support for topics.", + "defaultValue": true + }, { "name": "spring.pulsar.function.enabled", "type": "java.lang.Boolean", diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java index 1b5e3fed0e..b61f10fdae 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java @@ -219,7 +219,9 @@ class PulsarAutoConfigurationTests { "spring.pulsar.producer.cache.enabled=false") .run((context) -> assertThat(context).getBean(DefaultPulsarProducerFactory.class) .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class)) - .hasFieldOrPropertyWithValue("topicResolver", context.getBean(TopicResolver.class))); + .hasFieldOrPropertyWithValue("topicResolver", context.getBean(TopicResolver.class)) + .extracting("topicBuilder") + .isNotNull()); // prototype so only check not-null } @ParameterizedTest @@ -375,7 +377,9 @@ class PulsarAutoConfigurationTests { @Test void injectsExpectedBeans() { this.contextRunner.run((context) -> assertThat(context).getBean(DefaultPulsarConsumerFactory.class) - .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class))); + .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class)) + .extracting("topicBuilder") + .isNotNull()); // prototype so only check not-null } @Test @@ -574,7 +578,9 @@ class PulsarAutoConfigurationTests { @Test void injectsExpectedBeans() { this.contextRunner.run((context) -> assertThat(context).getBean(DefaultPulsarReaderFactory.class) - .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class))); + .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class)) + .extracting("topicBuilder") + .isNotNull()); // prototype so only check not-null } @Test diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfigurationTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfigurationTests.java index ef775cab63..c61777862e 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfigurationTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarConfigurationTests.java @@ -49,6 +49,7 @@ import org.springframework.pulsar.core.PulsarAdminBuilderCustomizer; import org.springframework.pulsar.core.PulsarAdministration; import org.springframework.pulsar.core.PulsarClientBuilderCustomizer; import org.springframework.pulsar.core.PulsarClientFactory; +import org.springframework.pulsar.core.PulsarTopicBuilder; import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer; import org.springframework.pulsar.core.TopicResolver; @@ -320,6 +321,46 @@ class PulsarConfigurationTests { } + @Nested + class TopicBuilderTests { + + private final ApplicationContextRunner contextRunner = PulsarConfigurationTests.this.contextRunner; + + @Test + void whenHasUserDefinedBeanDoesNotAutoConfigureBean() { + PulsarTopicBuilder topicBuilder = mock(PulsarTopicBuilder.class); + this.contextRunner.withBean("customPulsarTopicBuilder", PulsarTopicBuilder.class, () -> topicBuilder) + .run((context) -> assertThat(context).getBean(PulsarTopicBuilder.class).isSameAs(topicBuilder)); + } + + @Test + void whenHasDefaultsTopicDisabledPropertyDoesNotCreateBean() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat(context).doesNotHaveBean(PulsarTopicBuilder.class)); + } + + @Test + void whenHasDefaultsTenantAndNamespaceAppliedToTopicBuilder() { + List properties = new ArrayList<>(); + properties.add("spring.pulsar.defaults.topic.tenant=my-tenant"); + properties.add("spring.pulsar.defaults.topic.namespace=my-namespace"); + this.contextRunner.withPropertyValues(properties.toArray(String[]::new)) + .run((context) -> assertThat(context).getBean(PulsarTopicBuilder.class) + .asInstanceOf(InstanceOfAssertFactories.type(PulsarTopicBuilder.class)) + .satisfies((topicBuilder) -> { + assertThat(topicBuilder).hasFieldOrPropertyWithValue("defaultTenant", "my-tenant"); + assertThat(topicBuilder).hasFieldOrPropertyWithValue("defaultNamespace", "my-namespace"); + })); + } + + @Test + void beanHasScopePrototype() { + this.contextRunner.run((context) -> assertThat(context.getBean(PulsarTopicBuilder.class)) + .isNotSameAs(context.getBean(PulsarTopicBuilder.class))); + } + + } + @Nested class FunctionAdministrationTests { diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarPropertiesTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarPropertiesTests.java index 53e90a2b95..f7dd1b7e5a 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarPropertiesTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarPropertiesTests.java @@ -152,7 +152,7 @@ class PulsarPropertiesTests { } @Nested - class DefaultsProperties { + class DefaultsTypeMappingProperties { @Test void bindWhenNoTypeMappings() { @@ -242,6 +242,29 @@ class PulsarPropertiesTests { } + @Nested + class DefaultsTenantNamespaceProperties { + + @Test + void bindWhenValuesNotSpecified() { + assertThat(new PulsarProperties().getDefaults().getTopic()).satisfies((defaults) -> { + assertThat(defaults.getTenant()).isNull(); + assertThat(defaults.getNamespace()).isNull(); + }); + } + + @Test + void bindWhenValuesSpecified() { + Map map = new HashMap<>(); + map.put("spring.pulsar.defaults.topic.tenant", "my-tenant"); + map.put("spring.pulsar.defaults.topic.namespace", "my-namespace"); + PulsarProperties.Defaults.Topic properties = bindProperties(map).getDefaults().getTopic(); + assertThat(properties.getTenant()).isEqualTo("my-tenant"); + assertThat(properties.getNamespace()).isEqualTo("my-namespace"); + } + + } + @Nested class FunctionProperties { diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java index 4f3ab011ea..86fc67c9a0 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java @@ -48,6 +48,7 @@ import org.springframework.core.annotation.Order; import org.springframework.pulsar.core.DefaultSchemaResolver; import org.springframework.pulsar.core.DefaultTopicResolver; import org.springframework.pulsar.core.PulsarAdministration; +import org.springframework.pulsar.core.PulsarTopicBuilder; import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.core.TopicResolver; import org.springframework.pulsar.reactive.config.DefaultReactivePulsarListenerContainerFactory; @@ -177,6 +178,11 @@ class PulsarReactiveAutoConfigurationTests { assertThat(senderFactory) .extracting("topicResolver", InstanceOfAssertFactories.type(TopicResolver.class)) .isSameAs(context.getBean(TopicResolver.class)); + assertThat(senderFactory).extracting("topicBuilder").isNotNull(); // prototype + // so + // only + // check + // not-null }); } @@ -252,13 +258,18 @@ class PulsarReactiveAutoConfigurationTests { @Test void injectsExpectedBeans() { ReactivePulsarClient client = mock(ReactivePulsarClient.class); + PulsarTopicBuilder topicBuilder = mock(PulsarTopicBuilder.class); this.contextRunner.withBean("customReactivePulsarClient", ReactivePulsarClient.class, () -> client) + .withBean("customTopicBuilder", PulsarTopicBuilder.class, () -> topicBuilder) .run((context) -> { ReactivePulsarConsumerFactory consumerFactory = context .getBean(DefaultReactivePulsarConsumerFactory.class); assertThat(consumerFactory) .extracting("reactivePulsarClient", InstanceOfAssertFactories.type(ReactivePulsarClient.class)) .isSameAs(client); + assertThat(consumerFactory) + .extracting("topicBuilder", InstanceOfAssertFactories.type(PulsarTopicBuilder.class)) + .isSameAs(topicBuilder); }); } @@ -362,14 +373,19 @@ class PulsarReactiveAutoConfigurationTests { @Test void injectsExpectedBeans() { ReactivePulsarClient client = mock(ReactivePulsarClient.class); + PulsarTopicBuilder topicBuilder = mock(PulsarTopicBuilder.class); this.contextRunner.withPropertyValues("spring.pulsar.reader.name=test-reader") .withBean("customReactivePulsarClient", ReactivePulsarClient.class, () -> client) + .withBean("customPulsarTopicBuilder", PulsarTopicBuilder.class, () -> topicBuilder) .run((context) -> { DefaultReactivePulsarReaderFactory readerFactory = context .getBean(DefaultReactivePulsarReaderFactory.class); assertThat(readerFactory) .extracting("reactivePulsarClient", InstanceOfAssertFactories.type(ReactivePulsarClient.class)) .isSameAs(client); + assertThat(readerFactory) + .extracting("topicBuilder", InstanceOfAssertFactories.type(PulsarTopicBuilder.class)) + .isSameAs(topicBuilder); }); } From 3bbbef78be6c0e7d5f0ccc339c72da05bfef3dd2 Mon Sep 17 00:00:00 2001 From: Andy Wilkinson Date: Tue, 20 Aug 2024 10:27:47 +0100 Subject: [PATCH 2/2] Polish "Add support for Pulsar default tenant/namespace" See gh-41851 --- .../pulsar/PulsarAutoConfiguration.java | 20 ++++++---- .../pulsar/PulsarProperties.java | 14 +++---- .../PulsarReactiveAutoConfiguration.java | 21 +++++----- .../pulsar/PulsarAutoConfigurationTests.java | 39 +++++++++++++++++-- .../PulsarReactiveAutoConfigurationTests.java | 38 +++++++++++++++--- 5 files changed, 96 insertions(+), 36 deletions(-) diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java index 7d466d5719..5b5f9dc412 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfiguration.java @@ -89,12 +89,13 @@ public class PulsarAutoConfiguration { @ConditionalOnMissingBean(PulsarProducerFactory.class) @ConditionalOnProperty(name = "spring.pulsar.producer.cache.enabled", havingValue = "false") DefaultPulsarProducerFactory pulsarProducerFactory(PulsarClient pulsarClient, TopicResolver topicResolver, - ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { + ObjectProvider> customizersProvider, + ObjectProvider topicBuilderProvider) { List> lambdaSafeCustomizers = lambdaSafeProducerBuilderCustomizers( customizersProvider); DefaultPulsarProducerFactory producerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, this.properties.getProducer().getTopicName(), lambdaSafeCustomizers, topicResolver); - producerFactory.setTopicBuilder(topicBuilder); + topicBuilderProvider.ifAvailable(producerFactory::setTopicBuilder); return producerFactory; } @@ -102,7 +103,8 @@ public class PulsarAutoConfiguration { @ConditionalOnMissingBean(PulsarProducerFactory.class) @ConditionalOnProperty(name = "spring.pulsar.producer.cache.enabled", havingValue = "true", matchIfMissing = true) CachingPulsarProducerFactory cachingPulsarProducerFactory(PulsarClient pulsarClient, TopicResolver topicResolver, - ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { + ObjectProvider> customizersProvider, + ObjectProvider topicBuilderProvider) { PulsarProperties.Producer.Cache cacheProperties = this.properties.getProducer().getCache(); List> lambdaSafeCustomizers = lambdaSafeProducerBuilderCustomizers( customizersProvider); @@ -110,7 +112,7 @@ public class PulsarAutoConfiguration { this.properties.getProducer().getTopicName(), lambdaSafeCustomizers, topicResolver, cacheProperties.getExpireAfterAccess(), cacheProperties.getMaximumSize(), cacheProperties.getInitialCapacity()); - producerFactory.setTopicBuilder(topicBuilder); + topicBuilderProvider.ifAvailable(producerFactory::setTopicBuilder); return producerFactory; } @@ -144,7 +146,8 @@ public class PulsarAutoConfiguration { @Bean @ConditionalOnMissingBean(PulsarConsumerFactory.class) DefaultPulsarConsumerFactory pulsarConsumerFactory(PulsarClient pulsarClient, - ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { + ObjectProvider> customizersProvider, + ObjectProvider topicBuilderProvider) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeConsumerBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); @@ -152,7 +155,7 @@ public class PulsarAutoConfiguration { .of((builder) -> applyConsumerBuilderCustomizers(customizers, builder)); DefaultPulsarConsumerFactory consumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, lambdaSafeCustomizers); - consumerFactory.setTopicBuilder(topicBuilder); + topicBuilderProvider.ifAvailable(consumerFactory::setTopicBuilder); return consumerFactory; } @@ -190,7 +193,8 @@ public class PulsarAutoConfiguration { @Bean @ConditionalOnMissingBean(PulsarReaderFactory.class) DefaultPulsarReaderFactory pulsarReaderFactory(PulsarClient pulsarClient, - ObjectProvider> customizersProvider, PulsarTopicBuilder topicBuilder) { + ObjectProvider> customizersProvider, + ObjectProvider topicBuilderProvider) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeReaderBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); @@ -198,7 +202,7 @@ public class PulsarAutoConfiguration { .of((builder) -> applyReaderBuilderCustomizers(customizers, builder)); DefaultPulsarReaderFactory readerFactory = new DefaultPulsarReaderFactory<>(pulsarClient, lambdaSafeCustomizers); - readerFactory.setTopicBuilder(topicBuilder); + topicBuilderProvider.ifAvailable(readerFactory::setTopicBuilder); return readerFactory; } diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java index e4f6897f0b..45aefa0f5d 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarProperties.java @@ -257,15 +257,7 @@ public class PulsarProperties { */ private List typeMappings = new ArrayList<>(); - private Topic topic = new Topic(); - - public Topic getTopic() { - return this.topic; - } - - public void setTopic(Topic topic) { - this.topic = topic; - } + private final Topic topic = new Topic(); public List getTypeMappings() { return this.typeMappings; @@ -275,6 +267,10 @@ public class PulsarProperties { this.typeMappings = typeMappings; } + public Topic getTopic() { + return this.topic; + } + /** * A mapping from message type to topic and/or schema info to use (at least one of * {@code topicName} or {@code schemaInfo} must be specified. diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java index 6b983e9329..5ca96e7057 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfiguration.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2023 the original author or authors. + * Copyright 2012-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. @@ -49,6 +49,7 @@ import org.springframework.pulsar.reactive.config.annotation.EnableReactivePulsa import org.springframework.pulsar.reactive.core.DefaultReactivePulsarConsumerFactory; import org.springframework.pulsar.reactive.core.DefaultReactivePulsarReaderFactory; import org.springframework.pulsar.reactive.core.DefaultReactivePulsarSenderFactory; +import org.springframework.pulsar.reactive.core.DefaultReactivePulsarSenderFactory.Builder; import org.springframework.pulsar.reactive.core.ReactiveMessageConsumerBuilderCustomizer; import org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer; import org.springframework.pulsar.reactive.core.ReactiveMessageSenderBuilderCustomizer; @@ -114,18 +115,18 @@ public class PulsarReactiveAutoConfiguration { DefaultReactivePulsarSenderFactory reactivePulsarSenderFactory(ReactivePulsarClient reactivePulsarClient, ObjectProvider reactiveMessageSenderCache, TopicResolver topicResolver, ObjectProvider> customizersProvider, - PulsarTopicBuilder topicBuilder) { + ObjectProvider topicBuilderProvider) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeMessageSenderBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); List> lambdaSafeCustomizers = List .of((builder) -> applyMessageSenderBuilderCustomizers(customizers, builder)); - return DefaultReactivePulsarSenderFactory.builderFor(reactivePulsarClient) + Builder senderFactoryBuilder = DefaultReactivePulsarSenderFactory.builderFor(reactivePulsarClient) .withDefaultConfigCustomizers(lambdaSafeCustomizers) .withMessageSenderCache(reactiveMessageSenderCache.getIfAvailable()) - .withTopicResolver(topicResolver) - .withTopicBuilder(topicBuilder) - .build(); + .withTopicResolver(topicResolver); + topicBuilderProvider.ifAvailable(senderFactoryBuilder::withTopicBuilder); + return senderFactoryBuilder.build(); } @SuppressWarnings("unchecked") @@ -140,7 +141,7 @@ public class PulsarReactiveAutoConfiguration { DefaultReactivePulsarConsumerFactory reactivePulsarConsumerFactory( ReactivePulsarClient pulsarReactivePulsarClient, ObjectProvider> customizersProvider, - PulsarTopicBuilder topicBuilder) { + ObjectProvider topicBuilderProvider) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeMessageConsumerBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); @@ -148,7 +149,7 @@ public class PulsarReactiveAutoConfiguration { .of((builder) -> applyMessageConsumerBuilderCustomizers(customizers, builder)); DefaultReactivePulsarConsumerFactory consumerFactory = new DefaultReactivePulsarConsumerFactory<>( pulsarReactivePulsarClient, lambdaSafeCustomizers); - consumerFactory.setTopicBuilder(topicBuilder); + topicBuilderProvider.ifAvailable(consumerFactory::setTopicBuilder); return consumerFactory; } @@ -175,7 +176,7 @@ public class PulsarReactiveAutoConfiguration { @ConditionalOnMissingBean(ReactivePulsarReaderFactory.class) DefaultReactivePulsarReaderFactory reactivePulsarReaderFactory(ReactivePulsarClient reactivePulsarClient, ObjectProvider> customizersProvider, - PulsarTopicBuilder topicBuilder) { + ObjectProvider topicBuilderProvider) { List> customizers = new ArrayList<>(); customizers.add(this.propertiesMapper::customizeMessageReaderBuilder); customizers.addAll(customizersProvider.orderedStream().toList()); @@ -183,7 +184,7 @@ public class PulsarReactiveAutoConfiguration { .of((builder) -> applyMessageReaderBuilderCustomizers(customizers, builder)); DefaultReactivePulsarReaderFactory readerFactory = new DefaultReactivePulsarReaderFactory<>( reactivePulsarClient, lambdaSafeCustomizers); - readerFactory.setTopicBuilder(topicBuilder); + topicBuilderProvider.ifAvailable(readerFactory::setTopicBuilder); return readerFactory; } diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java index b61f10fdae..d8d30e942e 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarAutoConfigurationTests.java @@ -67,6 +67,7 @@ import org.springframework.pulsar.core.PulsarConsumerFactory; import org.springframework.pulsar.core.PulsarProducerFactory; import org.springframework.pulsar.core.PulsarReaderFactory; import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.PulsarTopicBuilder; import org.springframework.pulsar.core.ReaderBuilderCustomizer; import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.core.TopicResolver; @@ -126,6 +127,7 @@ class PulsarAutoConfigurationTests { .hasSingleBean(PulsarConnectionDetails.class) .hasSingleBean(DefaultPulsarClientFactory.class) .hasSingleBean(PulsarClient.class) + .hasSingleBean(PulsarTopicBuilder.class) .hasSingleBean(PulsarAdministration.class) .hasSingleBean(DefaultSchemaResolver.class) .hasSingleBean(DefaultTopicResolver.class) @@ -141,6 +143,12 @@ class PulsarAutoConfigurationTests { .hasSingleBean(PulsarReaderEndpointRegistry.class)); } + @Test + void topicDefaultsCanBeDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat(context).doesNotHaveBean(PulsarTopicBuilder.class)); + } + @Nested class ProducerFactoryTests { @@ -221,7 +229,15 @@ class PulsarAutoConfigurationTests { .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class)) .hasFieldOrPropertyWithValue("topicResolver", context.getBean(TopicResolver.class)) .extracting("topicBuilder") - .isNotNull()); // prototype so only check not-null + .isNotNull()); + } + + @Test + void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat(context).getBean(DefaultPulsarProducerFactory.class) + .extracting("topicBuilder") + .isNull()); } @ParameterizedTest @@ -379,7 +395,16 @@ class PulsarAutoConfigurationTests { this.contextRunner.run((context) -> assertThat(context).getBean(DefaultPulsarConsumerFactory.class) .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class)) .extracting("topicBuilder") - .isNotNull()); // prototype so only check not-null + .isNotNull()); + } + + @Test + void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat(context).getBean(DefaultPulsarConsumerFactory.class) + .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class)) + .extracting("topicBuilder") + .isNull()); } @Test @@ -580,7 +605,15 @@ class PulsarAutoConfigurationTests { this.contextRunner.run((context) -> assertThat(context).getBean(DefaultPulsarReaderFactory.class) .hasFieldOrPropertyWithValue("pulsarClient", context.getBean(PulsarClient.class)) .extracting("topicBuilder") - .isNotNull()); // prototype so only check not-null + .isNotNull()); + } + + @Test + void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat(context).getBean(DefaultPulsarReaderFactory.class) + .extracting("topicBuilder") + .isNull()); } @Test diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java index 86fc67c9a0..0ecb9e85e9 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/pulsar/PulsarReactiveAutoConfigurationTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2023 the original author or authors. + * Copyright 2012-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. @@ -115,6 +115,7 @@ class PulsarReactiveAutoConfigurationTests { void autoConfiguresBeans() { this.contextRunner.run((context) -> assertThat(context).hasSingleBean(PulsarConfiguration.class) .hasSingleBean(PulsarClient.class) + .hasSingleBean(PulsarTopicBuilder.class) .hasSingleBean(PulsarAdministration.class) .hasSingleBean(DefaultSchemaResolver.class) .hasSingleBean(DefaultTopicResolver.class) @@ -129,6 +130,12 @@ class PulsarReactiveAutoConfigurationTests { .hasSingleBean(ReactivePulsarListenerEndpointRegistry.class)); } + @Test + void topicDefaultsCanBeDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat(context).doesNotHaveBean(PulsarTopicBuilder.class)); + } + @Test @SuppressWarnings("rawtypes") void injectsExpectedBeansIntoReactivePulsarClient() { @@ -178,14 +185,17 @@ class PulsarReactiveAutoConfigurationTests { assertThat(senderFactory) .extracting("topicResolver", InstanceOfAssertFactories.type(TopicResolver.class)) .isSameAs(context.getBean(TopicResolver.class)); - assertThat(senderFactory).extracting("topicBuilder").isNotNull(); // prototype - // so - // only - // check - // not-null + assertThat(senderFactory).extracting("topicBuilder").isNotNull(); }); } + @Test + void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat((DefaultReactivePulsarSenderFactory) context + .getBean(DefaultReactivePulsarSenderFactory.class)).extracting("topicBuilder").isNull()); + } + @Test void injectsExpectedBeansIntoReactiveMessageSenderCache() { ProducerCacheProvider provider = mock(ProducerCacheProvider.class); @@ -273,6 +283,15 @@ class PulsarReactiveAutoConfigurationTests { }); } + @Test + void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat( + (ReactivePulsarConsumerFactory) context.getBean(DefaultReactivePulsarConsumerFactory.class)) + .extracting("topicBuilder") + .isNull()); + } + @Test void whenHasUserDefinedCustomizersAppliesInCorrectOrder() { this.contextRunner.withPropertyValues("spring.pulsar.consumer.name=fromPropsCustomizer") @@ -389,6 +408,13 @@ class PulsarReactiveAutoConfigurationTests { }); } + @Test + void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() { + this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false") + .run((context) -> assertThat((DefaultReactivePulsarReaderFactory) context + .getBean(DefaultReactivePulsarReaderFactory.class)).extracting("topicBuilder").isNull()); + } + @Test void whenHasUserDefinedCustomizersAppliesInCorrectOrder() { this.contextRunner.withPropertyValues("spring.pulsar.reader.name=fromPropsCustomizer")