Pulsar Reader auto config (imperative)

Resolves https://github.com/spring-projects-experimental/spring-pulsar/issues/328
This commit is contained in:
Soby Chacko
2023-02-07 12:06:28 -05:00
parent 9fb8879491
commit c6680973a5
4 changed files with 41 additions and 2 deletions

View File

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

View File

@@ -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"));

View File

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

View File

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