Merge pull request #41851 from
* gh-41851: Polish "Add support for Pulsar default tenant/namespace" Add support for Pulsar default tenant/namespace Closes gh-41851
This commit is contained in:
@@ -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,31 @@ public class PulsarAutoConfiguration {
|
||||
@ConditionalOnMissingBean(PulsarProducerFactory.class)
|
||||
@ConditionalOnProperty(name = "spring.pulsar.producer.cache.enabled", havingValue = "false")
|
||||
DefaultPulsarProducerFactory<?> pulsarProducerFactory(PulsarClient pulsarClient, TopicResolver topicResolver,
|
||||
ObjectProvider<ProducerBuilderCustomizer<?>> customizersProvider) {
|
||||
ObjectProvider<ProducerBuilderCustomizer<?>> customizersProvider,
|
||||
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
|
||||
List<ProducerBuilderCustomizer<Object>> lambdaSafeCustomizers = lambdaSafeProducerBuilderCustomizers(
|
||||
customizersProvider);
|
||||
return new DefaultPulsarProducerFactory<>(pulsarClient, this.properties.getProducer().getTopicName(),
|
||||
lambdaSafeCustomizers, topicResolver);
|
||||
DefaultPulsarProducerFactory<?> producerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
|
||||
this.properties.getProducer().getTopicName(), lambdaSafeCustomizers, topicResolver);
|
||||
topicBuilderProvider.ifAvailable(producerFactory::setTopicBuilder);
|
||||
return producerFactory;
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(PulsarProducerFactory.class)
|
||||
@ConditionalOnProperty(name = "spring.pulsar.producer.cache.enabled", havingValue = "true", matchIfMissing = true)
|
||||
CachingPulsarProducerFactory<?> cachingPulsarProducerFactory(PulsarClient pulsarClient, TopicResolver topicResolver,
|
||||
ObjectProvider<ProducerBuilderCustomizer<?>> customizersProvider) {
|
||||
ObjectProvider<ProducerBuilderCustomizer<?>> customizersProvider,
|
||||
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
|
||||
PulsarProperties.Producer.Cache cacheProperties = this.properties.getProducer().getCache();
|
||||
List<ProducerBuilderCustomizer<Object>> 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());
|
||||
topicBuilderProvider.ifAvailable(producerFactory::setTopicBuilder);
|
||||
return producerFactory;
|
||||
}
|
||||
|
||||
private List<ProducerBuilderCustomizer<Object>> lambdaSafeProducerBuilderCustomizers(
|
||||
@@ -138,13 +146,17 @@ public class PulsarAutoConfiguration {
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(PulsarConsumerFactory.class)
|
||||
DefaultPulsarConsumerFactory<?> pulsarConsumerFactory(PulsarClient pulsarClient,
|
||||
ObjectProvider<ConsumerBuilderCustomizer<?>> customizersProvider) {
|
||||
ObjectProvider<ConsumerBuilderCustomizer<?>> customizersProvider,
|
||||
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
|
||||
List<ConsumerBuilderCustomizer<?>> customizers = new ArrayList<>();
|
||||
customizers.add(this.propertiesMapper::customizeConsumerBuilder);
|
||||
customizers.addAll(customizersProvider.orderedStream().toList());
|
||||
List<ConsumerBuilderCustomizer<Object>> lambdaSafeCustomizers = List
|
||||
.of((builder) -> applyConsumerBuilderCustomizers(customizers, builder));
|
||||
return new DefaultPulsarConsumerFactory<>(pulsarClient, lambdaSafeCustomizers);
|
||||
DefaultPulsarConsumerFactory<?> consumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
|
||||
lambdaSafeCustomizers);
|
||||
topicBuilderProvider.ifAvailable(consumerFactory::setTopicBuilder);
|
||||
return consumerFactory;
|
||||
}
|
||||
|
||||
@Bean
|
||||
@@ -181,13 +193,17 @@ public class PulsarAutoConfiguration {
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(PulsarReaderFactory.class)
|
||||
DefaultPulsarReaderFactory<?> pulsarReaderFactory(PulsarClient pulsarClient,
|
||||
ObjectProvider<ReaderBuilderCustomizer<?>> customizersProvider) {
|
||||
ObjectProvider<ReaderBuilderCustomizer<?>> customizersProvider,
|
||||
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
|
||||
List<ReaderBuilderCustomizer<?>> customizers = new ArrayList<>();
|
||||
customizers.add(this.propertiesMapper::customizeReaderBuilder);
|
||||
customizers.addAll(customizersProvider.orderedStream().toList());
|
||||
List<ReaderBuilderCustomizer<Object>> lambdaSafeCustomizers = List
|
||||
.of((builder) -> applyReaderBuilderCustomizers(customizers, builder));
|
||||
return new DefaultPulsarReaderFactory<>(pulsarClient, lambdaSafeCustomizers);
|
||||
DefaultPulsarReaderFactory<?> readerFactory = new DefaultPulsarReaderFactory<>(pulsarClient,
|
||||
lambdaSafeCustomizers);
|
||||
topicBuilderProvider.ifAvailable(readerFactory::setTopicBuilder);
|
||||
return readerFactory;
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -257,6 +257,8 @@ public class PulsarProperties {
|
||||
*/
|
||||
private List<TypeMapping> typeMappings = new ArrayList<>();
|
||||
|
||||
private final Topic topic = new Topic();
|
||||
|
||||
public List<TypeMapping> getTypeMappings() {
|
||||
return this.typeMappings;
|
||||
}
|
||||
@@ -265,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.
|
||||
@@ -301,6 +307,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 {
|
||||
|
||||
@@ -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.
|
||||
@@ -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;
|
||||
@@ -48,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;
|
||||
@@ -112,17 +114,19 @@ public class PulsarReactiveAutoConfiguration {
|
||||
@ConditionalOnMissingBean(ReactivePulsarSenderFactory.class)
|
||||
DefaultReactivePulsarSenderFactory<?> reactivePulsarSenderFactory(ReactivePulsarClient reactivePulsarClient,
|
||||
ObjectProvider<ReactiveMessageSenderCache> reactiveMessageSenderCache, TopicResolver topicResolver,
|
||||
ObjectProvider<ReactiveMessageSenderBuilderCustomizer<?>> customizersProvider) {
|
||||
ObjectProvider<ReactiveMessageSenderBuilderCustomizer<?>> customizersProvider,
|
||||
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
|
||||
List<ReactiveMessageSenderBuilderCustomizer<?>> customizers = new ArrayList<>();
|
||||
customizers.add(this.propertiesMapper::customizeMessageSenderBuilder);
|
||||
customizers.addAll(customizersProvider.orderedStream().toList());
|
||||
List<ReactiveMessageSenderBuilderCustomizer<Object>> lambdaSafeCustomizers = List
|
||||
.of((builder) -> applyMessageSenderBuilderCustomizers(customizers, builder));
|
||||
return DefaultReactivePulsarSenderFactory.builderFor(reactivePulsarClient)
|
||||
Builder<Object> senderFactoryBuilder = DefaultReactivePulsarSenderFactory.builderFor(reactivePulsarClient)
|
||||
.withDefaultConfigCustomizers(lambdaSafeCustomizers)
|
||||
.withMessageSenderCache(reactiveMessageSenderCache.getIfAvailable())
|
||||
.withTopicResolver(topicResolver)
|
||||
.build();
|
||||
.withTopicResolver(topicResolver);
|
||||
topicBuilderProvider.ifAvailable(senderFactoryBuilder::withTopicBuilder);
|
||||
return senderFactoryBuilder.build();
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@@ -136,13 +140,17 @@ public class PulsarReactiveAutoConfiguration {
|
||||
@ConditionalOnMissingBean(ReactivePulsarConsumerFactory.class)
|
||||
DefaultReactivePulsarConsumerFactory<?> reactivePulsarConsumerFactory(
|
||||
ReactivePulsarClient pulsarReactivePulsarClient,
|
||||
ObjectProvider<ReactiveMessageConsumerBuilderCustomizer<?>> customizersProvider) {
|
||||
ObjectProvider<ReactiveMessageConsumerBuilderCustomizer<?>> customizersProvider,
|
||||
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
|
||||
List<ReactiveMessageConsumerBuilderCustomizer<?>> customizers = new ArrayList<>();
|
||||
customizers.add(this.propertiesMapper::customizeMessageConsumerBuilder);
|
||||
customizers.addAll(customizersProvider.orderedStream().toList());
|
||||
List<ReactiveMessageConsumerBuilderCustomizer<Object>> lambdaSafeCustomizers = List
|
||||
.of((builder) -> applyMessageConsumerBuilderCustomizers(customizers, builder));
|
||||
return new DefaultReactivePulsarConsumerFactory<>(pulsarReactivePulsarClient, lambdaSafeCustomizers);
|
||||
DefaultReactivePulsarConsumerFactory<?> consumerFactory = new DefaultReactivePulsarConsumerFactory<>(
|
||||
pulsarReactivePulsarClient, lambdaSafeCustomizers);
|
||||
topicBuilderProvider.ifAvailable(consumerFactory::setTopicBuilder);
|
||||
return consumerFactory;
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@@ -167,13 +175,17 @@ public class PulsarReactiveAutoConfiguration {
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(ReactivePulsarReaderFactory.class)
|
||||
DefaultReactivePulsarReaderFactory<?> reactivePulsarReaderFactory(ReactivePulsarClient reactivePulsarClient,
|
||||
ObjectProvider<ReactiveMessageReaderBuilderCustomizer<?>> customizersProvider) {
|
||||
ObjectProvider<ReactiveMessageReaderBuilderCustomizer<?>> customizersProvider,
|
||||
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
|
||||
List<ReactiveMessageReaderBuilderCustomizer<?>> customizers = new ArrayList<>();
|
||||
customizers.add(this.propertiesMapper::customizeMessageReaderBuilder);
|
||||
customizers.addAll(customizersProvider.orderedStream().toList());
|
||||
List<ReactiveMessageReaderBuilderCustomizer<Object>> lambdaSafeCustomizers = List
|
||||
.of((builder) -> applyMessageReaderBuilderCustomizers(customizers, builder));
|
||||
return new DefaultReactivePulsarReaderFactory<>(reactivePulsarClient, lambdaSafeCustomizers);
|
||||
DefaultReactivePulsarReaderFactory<?> readerFactory = new DefaultReactivePulsarReaderFactory<>(
|
||||
reactivePulsarClient, lambdaSafeCustomizers);
|
||||
topicBuilderProvider.ifAvailable(readerFactory::setTopicBuilder);
|
||||
return readerFactory;
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -219,7 +227,17 @@ 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());
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() {
|
||||
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
|
||||
.run((context) -> assertThat(context).getBean(DefaultPulsarProducerFactory.class)
|
||||
.extracting("topicBuilder")
|
||||
.isNull());
|
||||
}
|
||||
|
||||
@ParameterizedTest
|
||||
@@ -375,7 +393,18 @@ 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());
|
||||
}
|
||||
|
||||
@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
|
||||
@@ -574,7 +603,17 @@ 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());
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() {
|
||||
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
|
||||
.run((context) -> assertThat(context).getBean(DefaultPulsarReaderFactory.class)
|
||||
.extracting("topicBuilder")
|
||||
.isNull());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -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<String> 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 {
|
||||
|
||||
|
||||
@@ -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<String, String> 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 {
|
||||
|
||||
|
||||
@@ -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.
|
||||
@@ -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;
|
||||
@@ -114,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)
|
||||
@@ -128,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() {
|
||||
@@ -177,9 +185,17 @@ class PulsarReactiveAutoConfigurationTests {
|
||||
assertThat(senderFactory)
|
||||
.extracting("topicResolver", InstanceOfAssertFactories.type(TopicResolver.class))
|
||||
.isSameAs(context.getBean(TopicResolver.class));
|
||||
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);
|
||||
@@ -252,16 +268,30 @@ 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);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() {
|
||||
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
|
||||
.run((context) -> assertThat(
|
||||
(ReactivePulsarConsumerFactory<?>) context.getBean(DefaultReactivePulsarConsumerFactory.class))
|
||||
.extracting("topicBuilder")
|
||||
.isNull());
|
||||
}
|
||||
|
||||
@Test
|
||||
<T> void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
|
||||
this.contextRunner.withPropertyValues("spring.pulsar.consumer.name=fromPropsCustomizer")
|
||||
@@ -362,17 +392,29 @@ 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);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() {
|
||||
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
|
||||
.run((context) -> assertThat((DefaultReactivePulsarReaderFactory<?>) context
|
||||
.getBean(DefaultReactivePulsarReaderFactory.class)).extracting("topicBuilder").isNull());
|
||||
}
|
||||
|
||||
@Test
|
||||
<T> void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
|
||||
this.contextRunner.withPropertyValues("spring.pulsar.reader.name=fromPropsCustomizer")
|
||||
|
||||
Reference in New Issue
Block a user