From c6680973a51a07d3e97d90cca84df8228b42a0e8 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 7 Feb 2023 12:06:28 -0500 Subject: [PATCH] Pulsar Reader auto config (imperative) Resolves https://github.com/spring-projects-experimental/spring-pulsar/issues/328 --- .../PulsarAutoConfiguration.java | 8 +++++ .../autoconfigure/PulsarProperties.java | 2 +- .../PulsarAutoConfigurationTests.java | 31 +++++++++++++++++++ .../autoconfigure/PulsarPropertiesTests.java | 2 +- 4 files changed, 41 insertions(+), 2 deletions(-) diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfiguration.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfiguration.java index 9c8745a0..b2db6e0f 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfiguration.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfiguration.java @@ -35,11 +35,13 @@ import org.springframework.pulsar.config.PulsarClientFactoryBean; import org.springframework.pulsar.core.CachingPulsarProducerFactory; import org.springframework.pulsar.core.DefaultPulsarConsumerFactory; import org.springframework.pulsar.core.DefaultPulsarProducerFactory; +import org.springframework.pulsar.core.DefaultPulsarReaderFactory; import org.springframework.pulsar.core.DefaultSchemaResolver; import org.springframework.pulsar.core.DefaultTopicResolver; import org.springframework.pulsar.core.PulsarAdministration; 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.SchemaResolver; import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer; @@ -157,4 +159,10 @@ public class PulsarAutoConfiguration { this.properties.getFunction().getPropagateStopFailures()); } + @Bean + @ConditionalOnMissingBean + public PulsarReaderFactory pulsarReaderFactory(PulsarClient pulsarClient) { + return new DefaultPulsarReaderFactory<>(pulsarClient, this.properties.buildReaderProperties()); + } + } diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java index af444083..7471f8ba 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java @@ -2185,7 +2185,7 @@ public class PulsarProperties { PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull(); - map.from(this::getTopicNames).to(properties.in("topicName")); + map.from(this::getTopicNames).to(properties.in("topicNames")); map.from(this::getReceiverQueueSize).to(properties.in("receiverQueueSize")); map.from(this::getReaderName).to(properties.in("readerName")); map.from(this::getSubscriptionName).to(properties.in("subscriptionName")); diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java index 3342c3d9..f8633043 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java @@ -50,11 +50,13 @@ import org.springframework.pulsar.config.PulsarListenerContainerFactory; import org.springframework.pulsar.config.PulsarListenerEndpointRegistry; import org.springframework.pulsar.core.CachingPulsarProducerFactory; import org.springframework.pulsar.core.DefaultPulsarProducerFactory; +import org.springframework.pulsar.core.DefaultPulsarReaderFactory; import org.springframework.pulsar.core.DefaultSchemaResolver; import org.springframework.pulsar.core.DefaultTopicResolver; import org.springframework.pulsar.core.PulsarAdministration; 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.SchemaResolver; import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer; @@ -490,6 +492,35 @@ class PulsarAutoConfigurationTests { } + @Nested + class ReaderFactoryAutoConfigurationTests { + + @Test + void readerFactoryIsAutoConfiguredByDefault() { + contextRunner.run((context) -> assertThat(context).hasNotFailed().hasSingleBean(PulsarReaderFactory.class) + .getBean(PulsarReaderFactory.class).isExactlyInstanceOf(DefaultPulsarReaderFactory.class)); + } + + @Test + void readerFactoryCanBeConfigured() { + contextRunner.withPropertyValues("spring.pulsar.reader.topic-names=foo", + "spring.pulsar.reader.receiver-queue-size=200", "spring.pulsar.reader.reader-name=test-reader", + "spring.pulsar.reader.subscription-name=test-subscription", + "spring.pulsar.reader.subscription-role-prefix=test-prefix", + "spring.pulsar.reader.read-compacted=true", "spring.pulsar.reader.reset-include-head=true") + .run((context -> assertThat(context).hasNotFailed().getBean(PulsarReaderFactory.class) + .extracting("readerConfig") + .hasFieldOrPropertyWithValue("topicNames", new String[] { "foo" }) + .hasFieldOrPropertyWithValue("receiverQueueSize", 200) + .hasFieldOrPropertyWithValue("readerName", "test-reader") + .hasFieldOrPropertyWithValue("subscriptionName", "test-subscription") + .hasFieldOrPropertyWithValue("subscriptionRolePrefix", "test-prefix") + .hasFieldOrPropertyWithValue("readCompacted", true) + .hasFieldOrPropertyWithValue("resetIncludeHead", true))); + } + + } + @Configuration(proxyBeanMethods = false) static class InterceptorTestConfiguration { diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarPropertiesTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarPropertiesTests.java index e00c2bbb..f091fc01 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarPropertiesTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarPropertiesTests.java @@ -528,7 +528,7 @@ public class PulsarPropertiesTests { new ReaderConfigurationData<>(), ReaderConfigurationData.class)); assertThat(readerProps) - .hasEntrySatisfying("topicName", + .hasEntrySatisfying("topicNames", topics -> assertThat(topics).asInstanceOf(InstanceOfAssertFactories.array(String[].class)) .containsExactly("my-topic")) .containsEntry("receiverQueueSize", 100).containsEntry("readerName", "my-reader")