Add autoconfiguration for ReactivePulsarReaderFactory

This commit is contained in:
Christophe Bornet
2022-11-01 00:28:38 -05:00
committed by Chris Bono
parent 8c37afeffe
commit 02479e8cb9
4 changed files with 153 additions and 5 deletions

View File

@@ -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());
}
}

View File

@@ -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(

View File

@@ -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

View File

@@ -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 })
<T> void customBeanIsRespected(Class<T> 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<? extends AbstractObjectAssert<?, DefaultReactivePulsarReaderFactory>, 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)) {