diff --git a/buildSrc/src/main/java/org/springframework/pulsar/gradle/docs/configprops/DocumentConfigurationProperties.java b/buildSrc/src/main/java/org/springframework/pulsar/gradle/docs/configprops/DocumentConfigurationProperties.java index 84b7200f..ab1a8fb0 100644 --- a/buildSrc/src/main/java/org/springframework/pulsar/gradle/docs/configprops/DocumentConfigurationProperties.java +++ b/buildSrc/src/main/java/org/springframework/pulsar/gradle/docs/configprops/DocumentConfigurationProperties.java @@ -80,6 +80,9 @@ public class DocumentConfigurationProperties extends DefaultTask { snippets.add("application-properties.pulsar-reactive-consumer", "Pulsar Reactive Consumer Properties", (c) -> { c.accept("spring.pulsar.reactive.consumer"); }); + snippets.add("application-properties.pulsar-reactive-reader", "Pulsar Reactive Reader Properties", (c) -> { + c.accept("spring.pulsar.reactive.reader"); + }); snippets.writeTo(this.outputDir.toPath()); } } diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java index 77f8ac21..6c088341 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java @@ -32,8 +32,10 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.pulsar.core.reactive.DefaultReactivePulsarConsumerFactory; +import org.springframework.pulsar.core.reactive.DefaultReactivePulsarReaderFactory; import org.springframework.pulsar.core.reactive.DefaultReactivePulsarSenderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory; +import org.springframework.pulsar.core.reactive.ReactivePulsarReaderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarSenderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate; @@ -100,6 +102,13 @@ public class PulsarReactiveAutoConfiguration { this.properties.buildReactiveMessageConsumerSpec()); } + @Bean + @ConditionalOnMissingBean + public ReactivePulsarReaderFactory reactivePulsarReaderFactory(ReactivePulsarClient pulsarReactivePulsarClient) { + return new DefaultReactivePulsarReaderFactory<>(pulsarReactivePulsarClient, + this.properties.buildReactiveMessageReaderSpec()); + } + @Bean @ConditionalOnMissingBean public ReactivePulsarSenderTemplate pulsarReactiveSenderTemplate( diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java index b64a8135..ccb77c52 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java @@ -31,14 +31,18 @@ import org.apache.pulsar.client.api.KeySharedPolicy; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.ProducerAccessMode; import org.apache.pulsar.client.api.ProducerCryptoFailureAction; +import org.apache.pulsar.client.api.Range; import org.apache.pulsar.client.api.RegexSubscriptionMode; import org.apache.pulsar.client.api.SubscriptionInitialPosition; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.reactive.client.api.ImmutableReactiveMessageConsumerSpec; +import org.apache.pulsar.reactive.client.api.ImmutableReactiveMessageReaderSpec; import org.apache.pulsar.reactive.client.api.ImmutableReactiveMessageSenderSpec; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageConsumerSpec; +import org.apache.pulsar.reactive.client.api.MutableReactiveMessageReaderSpec; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageSenderSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerSpec; +import org.apache.pulsar.reactive.client.api.ReactiveMessageReaderSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; import org.springframework.boot.context.properties.ConfigurationProperties; @@ -61,6 +65,8 @@ public class PulsarReactiveProperties { private final Consumer consumer = new Consumer(); + private final Reader reader = new Reader(); + public Sender getSender() { return this.sender; } @@ -69,10 +75,18 @@ public class PulsarReactiveProperties { return this.consumer; } + public Reader getReader() { + return this.reader; + } + public ReactiveMessageSenderSpec buildReactiveMessageSenderSpec() { return this.sender.buildReactiveMessageSenderSpec(); } + public ReactiveMessageReaderSpec buildReactiveMessageReaderSpec() { + return this.reader.buildReactiveMessageReaderSpec(); + } + public ReactiveMessageConsumerSpec buildReactiveMessageConsumerSpec() { return this.consumer.buildReactiveMessageConsumerSpec(); } @@ -307,6 +321,106 @@ public class PulsarReactiveProperties { } + public static class Reader { + + private String[] topicNames; + + private String readerName; + + private String subscriptionName; + + private String generatedSubscriptionNamePrefix; + + private Integer receiverQueueSize; + + private Boolean readCompacted; + + private Range[] keyHashRanges; + + private ConsumerCryptoFailureAction cryptoFailureAction; + + public String[] getTopicNames() { + return this.topicNames; + } + + public void setTopicNames(String[] topicNames) { + this.topicNames = topicNames; + } + + public String getReaderName() { + return this.readerName; + } + + public void setReaderName(String readerName) { + this.readerName = readerName; + } + + public String getSubscriptionName() { + return this.subscriptionName; + } + + public void setSubscriptionName(String subscriptionName) { + this.subscriptionName = subscriptionName; + } + + public String getGeneratedSubscriptionNamePrefix() { + return this.generatedSubscriptionNamePrefix; + } + + public void setGeneratedSubscriptionNamePrefix(String generatedSubscriptionNamePrefix) { + this.generatedSubscriptionNamePrefix = generatedSubscriptionNamePrefix; + } + + public Integer getReceiverQueueSize() { + return this.receiverQueueSize; + } + + public void setReceiverQueueSize(Integer receiverQueueSize) { + this.receiverQueueSize = receiverQueueSize; + } + + public Boolean getReadCompacted() { + return this.readCompacted; + } + + public void setReadCompacted(Boolean readCompacted) { + this.readCompacted = readCompacted; + } + + public Range[] getKeyHashRanges() { + return this.keyHashRanges; + } + + public void setKeyHashRanges(Range[] keyHashRanges) { + this.keyHashRanges = keyHashRanges; + } + + public ConsumerCryptoFailureAction getCryptoFailureAction() { + return this.cryptoFailureAction; + } + + public void setCryptoFailureAction(ConsumerCryptoFailureAction cryptoFailureAction) { + this.cryptoFailureAction = cryptoFailureAction; + } + + public ReactiveMessageReaderSpec buildReactiveMessageReaderSpec() { + PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull(); + + MutableReactiveMessageReaderSpec spec = new MutableReactiveMessageReaderSpec(); + + map.from(this::getTopicNames).as(List::of).to(spec::setTopicNames); + map.from(this::getReaderName).to(spec::setReaderName); + map.from(this::getSubscriptionName).to(spec::setSubscriptionName); + map.from(this::getGeneratedSubscriptionNamePrefix).to(spec::setGeneratedSubscriptionNamePrefix); + map.from(this::getReceiverQueueSize).to(spec::setReceiverQueueSize); + map.from(this::getReadCompacted).to(spec::setReadCompacted); + map.from(this::getKeyHashRanges).as(List::of).to(spec::setKeyHashRanges); + + return new ImmutableReactiveMessageReaderSpec(spec); + } + + } + public static class Consumer { /** @@ -772,19 +886,22 @@ public class PulsarReactiveProperties { public enum SchedulerType { /** - * Reactor's {@link reactor.core.scheduler.BoundedElasticScheduler}. + * The reactor.core.scheduler.BoundedElasticScheduler. */ boundedElastic, + /** - * Reactor's Reactor's {@link reactor.core.scheduler.ParallelScheduler}. + * The reactor.core.scheduler.ParallelScheduler. */ parallel, + /** - * Reactor's Reactor's {@link reactor.core.scheduler.SingleScheduler}. + * The reactor.core.scheduler.SingleScheduler. */ single, + /** - * Reactor's {@link reactor.core.scheduler.ImmediateScheduler}. + * The reactor.core.scheduler.ImmediateScheduler. */ immediate diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java index cb20f2e8..8fed9309 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java @@ -26,6 +26,7 @@ import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory; import org.apache.pulsar.reactive.client.adapter.ProducerCacheProvider; import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerSpec; +import org.apache.pulsar.reactive.client.api.ReactiveMessageReaderSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderCache; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; import org.apache.pulsar.reactive.client.api.ReactivePulsarClient; @@ -45,8 +46,10 @@ import org.springframework.boot.test.context.assertj.AssertableApplicationContex import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.pulsar.config.PulsarClientFactoryBean; import org.springframework.pulsar.core.reactive.DefaultReactivePulsarConsumerFactory; +import org.springframework.pulsar.core.reactive.DefaultReactivePulsarReaderFactory; import org.springframework.pulsar.core.reactive.DefaultReactivePulsarSenderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory; +import org.springframework.pulsar.core.reactive.ReactivePulsarReaderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarSenderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate; @@ -84,7 +87,7 @@ class PulsarReactiveAutoConfigurationTests { @ParameterizedTest @ValueSource(classes = { ReactivePulsarClient.class, ProducerCacheProvider.class, ReactiveMessageSenderCache.class, - ReactivePulsarSenderFactory.class, ReactivePulsarConsumerFactory.class, + ReactivePulsarSenderFactory.class, ReactivePulsarConsumerFactory.class, ReactivePulsarReaderFactory.class, ReactivePulsarSenderTemplate.class }) void customBeanIsRespected(Class beanClass) { T bean = mock(beanClass); @@ -144,6 +147,22 @@ class PulsarReactiveAutoConfigurationTests { } + @Test + @SuppressWarnings("rawtypes") + void beansAreInjectedInReactivePulsarReaderFactory() { + ReactivePulsarClient client = mock(ReactivePulsarClient.class); + this.contextRunner.withPropertyValues("spring.pulsar.reactive.reader.reader-name=test-reader") + .withBean("customReactivePulsarClient", ReactivePulsarClient.class, () -> client).run((context -> { + AbstractObjectAssert, DefaultReactivePulsarReaderFactory> senderFactory = assertThat( + context).hasNotFailed().getBean(DefaultReactivePulsarReaderFactory.class); + senderFactory + .extracting("readerSpec", InstanceOfAssertFactories.type(ReactiveMessageReaderSpec.class)) + .extracting(ReactiveMessageReaderSpec::getReaderName).isEqualTo("test-reader"); + senderFactory.extracting("reactivePulsarClient", + InstanceOfAssertFactories.type(ReactivePulsarClient.class)).isSameAs(client); + })); + } + @Test void beansAreInjectedInReactiveMessageSenderCache() throws Exception { try (ProducerCacheProvider provider = mock(ProducerCacheProvider.class)) {