diff --git a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarSenderFactory.java b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarSenderFactory.java index 39c118bf..929a6091 100644 --- a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarSenderFactory.java +++ b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarSenderFactory.java @@ -55,23 +55,49 @@ public class DefaultReactivePulsarSenderFactory implements ReactivePulsarSend @Nullable private final ReactiveMessageSenderCache reactiveMessageSenderCache; + @Nullable + private final List> defaultSenderBuilderCustomizers; + private TopicResolver topicResolver; + /** + * Construct an instance. + * @param pulsarClient the pulsar client to adapt into a reactive client + * @param reactiveMessageSenderSpec spec that defines the initial settings on the + * created senders + * @param reactiveMessageSenderCache cache used to cache created senders + * @param defaultSenderBuilderCustomizers optional list of sender builder customizers + * to apply to the created senders + */ public DefaultReactivePulsarSenderFactory(PulsarClient pulsarClient, @Nullable ReactiveMessageSenderSpec reactiveMessageSenderSpec, - @Nullable ReactiveMessageSenderCache reactiveMessageSenderCache) { + @Nullable ReactiveMessageSenderCache reactiveMessageSenderCache, + @Nullable List> defaultSenderBuilderCustomizers) { this(AdaptedReactivePulsarClientFactory.create(pulsarClient), reactiveMessageSenderSpec, - reactiveMessageSenderCache, new DefaultTopicResolver()); + reactiveMessageSenderCache, defaultSenderBuilderCustomizers, new DefaultTopicResolver()); } + /** + * Construct an instance. + * @param reactivePulsarClient the reactive client to use + * @param reactiveMessageSenderSpec spec that defines the initial settings on the + * created senders + * @param reactiveMessageSenderCache cache used to cache created senders + * @param defaultSenderBuilderCustomizers optional list of sender builder customizers + * to apply to the created senders + * @param topicResolver the topic resolver to use + */ public DefaultReactivePulsarSenderFactory(ReactivePulsarClient reactivePulsarClient, @Nullable ReactiveMessageSenderSpec reactiveMessageSenderSpec, - @Nullable ReactiveMessageSenderCache reactiveMessageSenderCache, TopicResolver topicResolver) { + @Nullable ReactiveMessageSenderCache reactiveMessageSenderCache, + @Nullable List> defaultSenderBuilderCustomizers, + TopicResolver topicResolver) { this.reactivePulsarClient = reactivePulsarClient; this.reactiveMessageSenderSpec = new ImmutableReactiveMessageSenderSpec( reactiveMessageSenderSpec != null ? reactiveMessageSenderSpec : new MutableReactiveMessageSenderSpec()); this.reactiveMessageSenderCache = reactiveMessageSenderCache; this.topicResolver = topicResolver; + this.defaultSenderBuilderCustomizers = defaultSenderBuilderCustomizers; } @Override @@ -102,6 +128,11 @@ public class DefaultReactivePulsarSenderFactory implements ReactivePulsarSend ReactiveMessageSenderBuilder sender = this.reactivePulsarClient.messageSender(schema); sender.applySpec(this.reactiveMessageSenderSpec); + + // Apply the default config customizer (preserve the topic) + if (!CollectionUtils.isEmpty(this.defaultSenderBuilderCustomizers)) { + this.defaultSenderBuilderCustomizers.forEach((customizer -> customizer.customize(sender))); + } sender.topic(resolvedTopic); if (this.reactiveMessageSenderCache != null) { sender.cache(this.reactiveMessageSenderCache); diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageConsumerFactoryTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarConsumerFactoryTests.java similarity index 98% rename from spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageConsumerFactoryTests.java rename to spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarConsumerFactoryTests.java index f2b3f572..4afe27dd 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageConsumerFactoryTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarConsumerFactoryTests.java @@ -37,7 +37,7 @@ import org.junit.jupiter.api.Test; * @author Christophe Bornet * @author Chris Bono */ -class DefaultReactiveMessageConsumerFactoryTests { +class DefaultReactivePulsarConsumerFactoryTests { private static final Schema SCHEMA = Schema.STRING; diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageReaderFactoryTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarReaderFactoryTests.java similarity index 97% rename from spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageReaderFactoryTests.java rename to spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarReaderFactoryTests.java index 34cdf228..d0064fb7 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageReaderFactoryTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarReaderFactoryTests.java @@ -33,8 +33,9 @@ import org.junit.jupiter.api.Test; * Tests for {@link DefaultReactivePulsarReaderFactory}. * * @author Christophe Bornet + * @author Chris Bono */ -class DefaultReactiveMessageReaderFactoryTests { +class DefaultReactivePulsarReaderFactoryTests { private static final Schema schema = Schema.STRING; diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageSenderFactoryTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarSenderFactoryTests.java similarity index 69% rename from spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageSenderFactoryTests.java rename to spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarSenderFactoryTests.java index 2fd62b0d..e8ba85e9 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactiveMessageSenderFactoryTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/DefaultReactivePulsarSenderFactoryTests.java @@ -19,22 +19,27 @@ package org.springframework.pulsar.reactive.core; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; import static org.assertj.core.api.Assertions.assertThatNullPointerException; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; import java.util.Arrays; import java.util.Collections; +import java.util.List; import org.apache.pulsar.client.api.CompressionType; -import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageSenderSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageSender; +import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderBuilder; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderCache; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; import org.assertj.core.api.InstanceOfAssertFactories; import org.assertj.core.api.ThrowingConsumer; import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; +import org.mockito.InOrder; /** * Unit tests for {@link DefaultReactivePulsarSenderFactory}. @@ -42,7 +47,7 @@ import org.junit.jupiter.api.Test; * @author Christophe Bornet * @author Chris Bono */ -class DefaultReactiveMessageSenderFactoryTests { +class DefaultReactivePulsarSenderFactoryTests { protected final Schema schema = Schema.STRING; @@ -66,17 +71,17 @@ class DefaultReactiveMessageSenderFactoryTests { } private ReactivePulsarSenderFactory newSenderFactory() { - return new DefaultReactivePulsarSenderFactory<>((PulsarClient) null, null, null); + return new DefaultReactivePulsarSenderFactory<>(null, null, null, null); } private ReactivePulsarSenderFactory newSenderFactoryWithDefaultTopic(String defaultTopic) { MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec(); senderSpec.setTopicName(defaultTopic); - return new DefaultReactivePulsarSenderFactory<>((PulsarClient) null, senderSpec, null); + return new DefaultReactivePulsarSenderFactory<>(null, senderSpec, null, null); } private ReactivePulsarSenderFactory newSenderFactoryWithCache(ReactiveMessageSenderCache cache) { - return new DefaultReactivePulsarSenderFactory<>((PulsarClient) null, null, cache); + return new DefaultReactivePulsarSenderFactory<>(null, null, cache, null); } @Nested @@ -150,4 +155,44 @@ class DefaultReactiveMessageSenderFactoryTests { } + @Nested + @SuppressWarnings("unchecked") + class DefaultConfigCustomizerApi { + + private ReactiveMessageSenderBuilderCustomizer configCustomizer1 = mock( + ReactiveMessageSenderBuilderCustomizer.class); + + private ReactiveMessageSenderBuilderCustomizer configCustomizer2 = mock( + ReactiveMessageSenderBuilderCustomizer.class); + + private ReactiveMessageSenderBuilderCustomizer createSenderCustomizer = mock( + ReactiveMessageSenderBuilderCustomizer.class); + + @Test + void singleConfigCustomizer() { + newSenderFactoryWithCustomizers(List.of(configCustomizer1)).createSender(schema, "topic1", + List.of(createSenderCustomizer)); + InOrder inOrder = inOrder(configCustomizer1, createSenderCustomizer); + inOrder.verify(configCustomizer1).customize(any(ReactiveMessageSenderBuilder.class)); + inOrder.verify(createSenderCustomizer).customize(any(ReactiveMessageSenderBuilder.class)); + } + + @Test + void multipleConfigCustomizers() { + newSenderFactoryWithCustomizers(List.of(configCustomizer2, configCustomizer1)).createSender(schema, + "topic1", List.of(createSenderCustomizer)); + InOrder inOrder = inOrder(configCustomizer1, configCustomizer2, createSenderCustomizer); + inOrder.verify(configCustomizer2).customize(any(ReactiveMessageSenderBuilder.class)); + inOrder.verify(configCustomizer1).customize(any(ReactiveMessageSenderBuilder.class)); + inOrder.verify(createSenderCustomizer).customize(any(ReactiveMessageSenderBuilder.class)); + } + + private ReactivePulsarSenderFactory newSenderFactoryWithCustomizers( + List> customizers) { + MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec(); + return new DefaultReactivePulsarSenderFactory<>(null, senderSpec, null, customizers); + } + + } + } diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java index 85afa23a..96213ad3 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java @@ -182,7 +182,8 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { if (producerFactoryHasDefaultTopic) { spec.setTopicName("fake-topic"); } - ReactivePulsarSenderFactory producerFactory = new DefaultReactivePulsarSenderFactory<>(client, spec, null); + ReactivePulsarSenderFactory producerFactory = new DefaultReactivePulsarSenderFactory<>(client, spec, null, + null); // Topic mappings allows not specifying the topic when sending (nor having // default on producer) DefaultTopicResolver topicResolver = new DefaultTopicResolver(); @@ -198,7 +199,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { @Test void sendMessageWithoutTopicFails() { ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(client, - new MutableReactiveMessageSenderSpec(), null); + new MutableReactiveMessageSenderSpec(), null, null); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory); assertThatIllegalArgumentException().isThrownBy(() -> pulsarTemplate.send("test-message").subscribe()) .withMessage("Topic must be specified when no default topic is configured"); @@ -211,7 +212,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { senderSpec.setTopicName(topic); } ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(client, senderSpec, - null); + null, null); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory); @@ -260,7 +261,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { MutableReactiveMessageSenderSpec spec = new MutableReactiveMessageSenderSpec(); spec.setTopicName(topic); ReactivePulsarSenderFactory producerFactory = new DefaultReactivePulsarSenderFactory<>(client, spec, - null); + null, null); // Custom schema resolver allows not specifying the schema when sending DefaultSchemaResolver schemaResolver = new DefaultSchemaResolver(); schemaResolver.addCustomSchemaMapping(Foo.class, Schema.JSON(Foo.class)); @@ -282,7 +283,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { MutableReactiveMessageSenderSpec spec = new MutableReactiveMessageSenderSpec(); spec.setTopicName("sendNullWithDefaultTopicFails"); ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(client, spec, - null); + null, null); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory); assertThatIllegalArgumentException() .isThrownBy(() -> pulsarTemplate.send((String) null, Schema.STRING).subscribe()) @@ -292,7 +293,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { @Test void sendNullWithoutSchemaFails() { ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(client, - new MutableReactiveMessageSenderSpec(), null); + new MutableReactiveMessageSenderSpec(), null, null); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory); assertThatIllegalArgumentException() .isThrownBy(() -> pulsarTemplate.send("sendNullWithoutSchemaFails", (String) null, null).subscribe()) diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java index 96dee7b0..f745e3c9 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java @@ -82,7 +82,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo MutableReactiveMessageSenderSpec prodConfig = new MutableReactiveMessageSenderSpec(); prodConfig.setTopicName(topic); DefaultReactivePulsarSenderFactory pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>( - reactivePulsarClient, prodConfig, null, new DefaultTopicResolver()); + reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver()); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.send("hello john doe").subscribe(); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -116,7 +116,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo MutableReactiveMessageSenderSpec prodConfig = new MutableReactiveMessageSenderSpec(); prodConfig.setTopicName(topic); DefaultReactivePulsarSenderFactory pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>( - reactivePulsarClient, prodConfig, null, new DefaultTopicResolver()); + reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver()); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(pulsarProducerFactory); Flux.range(0, 5).map(i -> MessageSpec.of("hello john doe" + i)).as(pulsarTemplate::send).subscribe(); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -157,7 +157,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo MutableReactiveMessageSenderSpec prodConfig = new MutableReactiveMessageSenderSpec(); prodConfig.setTopicName(topic); DefaultReactivePulsarSenderFactory pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>( - reactivePulsarClient, prodConfig, null, new DefaultTopicResolver()); + reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver()); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.send("hello john doe").subscribe(); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -271,7 +271,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo MutableReactiveMessageSenderSpec prodConfig = new MutableReactiveMessageSenderSpec(); prodConfig.setTopicName(topic); DefaultReactivePulsarSenderFactory pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>( - reactivePulsarClient, prodConfig, null, new DefaultTopicResolver()); + reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver()); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.send("hello john doe").subscribe(); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -321,7 +321,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo prodConfig.setBatchingEnabled(false); prodConfig.setTopicName(topic); DefaultReactivePulsarSenderFactory pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>( - reactivePulsarClient, prodConfig, null, new DefaultTopicResolver()); + reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver()); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(pulsarProducerFactory); Flux.range(0, 5).map(i -> MessageSpec.of("hello john doe" + i)).as(pulsarTemplate::send).subscribe(); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactory.java index 96c3cd66..74374c6f 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactory.java @@ -48,18 +48,18 @@ public class DefaultPulsarConsumerFactory implements PulsarConsumerFactory private final PulsarClient pulsarClient; @Nullable - private final ConsumerBuilderCustomizer defaultConfigCustomizer; + private final List> defaultConfigCustomizers; /** * Construct a consumer factory instance. * @param pulsarClient the client used to consume - * @param defaultConfigCustomizer the default configuration to apply to the consumers - * or null to use no default configuration + * @param defaultConfigCustomizers the optional list of customizers to apply to the + * created consumers or null to use no default configuration */ public DefaultPulsarConsumerFactory(PulsarClient pulsarClient, - ConsumerBuilderCustomizer defaultConfigCustomizer) { + List> defaultConfigCustomizers) { this.pulsarClient = pulsarClient; - this.defaultConfigCustomizer = defaultConfigCustomizer; + this.defaultConfigCustomizers = defaultConfigCustomizers; } @Override @@ -77,8 +77,8 @@ public class DefaultPulsarConsumerFactory implements PulsarConsumerFactory ConsumerBuilder consumerBuilder = this.pulsarClient.newConsumer(schema); // Apply the default config customizer (preserve the topic) - if (this.defaultConfigCustomizer != null) { - this.defaultConfigCustomizer.customize(consumerBuilder); + if (!CollectionUtils.isEmpty(this.defaultConfigCustomizers)) { + this.defaultConfigCustomizers.forEach((customizer -> customizer.customize(consumerBuilder))); } if (topics != null) { replaceTopicsOnBuilder(consumerBuilder, topics); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarProducerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarProducerFactory.java index 8e24a901..9ce28833 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarProducerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarProducerFactory.java @@ -142,7 +142,7 @@ public class DefaultPulsarProducerFactory implements PulsarProducerFactory var producerBuilder = this.pulsarClient.newProducer(schema); // Apply the default config customizer (preserve the topic) - if (this.defaultConfigCustomizers != null) { + if (!CollectionUtils.isEmpty(this.defaultConfigCustomizers)) { this.defaultConfigCustomizers.forEach((customizer) -> customizer.customize(producerBuilder)); } producerBuilder.topic(resolvedTopic); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java index f07d450a..cb4034bb 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java @@ -43,7 +43,7 @@ public class DefaultPulsarReaderFactory implements PulsarReaderFactory { private final PulsarClient pulsarClient; @Nullable - private final ReaderBuilderCustomizer defaultConfigCustomizer; + private final List> defaultConfigCustomizers; /** * Construct a reader factory instance with no default configuration. @@ -56,13 +56,13 @@ public class DefaultPulsarReaderFactory implements PulsarReaderFactory { /** * Construct a reader factory instance. * @param pulsarClient the client used to consume - * @param defaultConfigCustomizer the default configuration to apply to the readers or - * null to use no default configuration + * @param defaultConfigCustomizers the optional list of customizers to apply to the + * readers or null to use no default configuration */ public DefaultPulsarReaderFactory(PulsarClient pulsarClient, - @Nullable ReaderBuilderCustomizer defaultConfigCustomizer) { + @Nullable List> defaultConfigCustomizers) { this.pulsarClient = pulsarClient; - this.defaultConfigCustomizer = defaultConfigCustomizer; + this.defaultConfigCustomizers = defaultConfigCustomizers; } @Override @@ -72,8 +72,8 @@ public class DefaultPulsarReaderFactory implements PulsarReaderFactory { ReaderBuilder readerBuilder = this.pulsarClient.newReader(schema); // Apply the default config customizer (preserve the topics) - if (this.defaultConfigCustomizer != null) { - this.defaultConfigCustomizer.customize(readerBuilder); + if (!CollectionUtils.isEmpty(this.defaultConfigCustomizers)) { + this.defaultConfigCustomizers.forEach((customizer -> customizer.customize(readerBuilder))); } if (!CollectionUtils.isEmpty(topics)) { diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java index 99b9e303..9427aaa9 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java @@ -358,11 +358,11 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarClient.close(); } - private ConsumerBuilderCustomizer defaultConfig(String topicName, String subscriptionName) { - return (consumerBuilder) -> { + private List> defaultConfig(String topicName, String subscriptionName) { + return List.of((consumerBuilder) -> { consumerBuilder.topic(topicName); consumerBuilder.subscriptionName(subscriptionName); - }; + }); } } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactoryTests.java index 85d5b122..89020e2b 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarConsumerFactoryTests.java @@ -170,11 +170,11 @@ class DefaultPulsarConsumerFactoryTests implements PulsarTestContainerSupport { @BeforeEach void createConsumerFactory() { - consumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, (consumerBuilder) -> { + consumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, List.of((consumerBuilder) -> { consumerBuilder.topic(defaultTopic); consumerBuilder.subscriptionName(defaultSubscription); consumerBuilder.properties(defaultMetadataProperties); - }); + })); } @Test @@ -228,6 +228,40 @@ class DefaultPulsarConsumerFactoryTests implements PulsarTestContainerSupport { } } + @Nested + @SuppressWarnings("unchecked") + class DefaultConfigCustomizerApi { + + private ConsumerBuilderCustomizer configCustomizer1 = mock(ConsumerBuilderCustomizer.class); + + private ConsumerBuilderCustomizer configCustomizer2 = mock(ConsumerBuilderCustomizer.class); + + private ConsumerBuilderCustomizer createConsumerCustomizer = mock(ConsumerBuilderCustomizer.class); + + @Test + void singleConfigCustomizer() throws PulsarClientException { + try (var ignored = new DefaultPulsarConsumerFactory<>(pulsarClient, List.of(configCustomizer1)) + .createConsumer(SCHEMA, List.of("topic0"), "dft-sub", createConsumerCustomizer)) { + InOrder inOrder = inOrder(configCustomizer1, createConsumerCustomizer); + inOrder.verify(configCustomizer1).customize(any(ConsumerBuilder.class)); + inOrder.verify(createConsumerCustomizer).customize(any(ConsumerBuilder.class)); + } + } + + @Test + void multipleConfigCustomizers() throws PulsarClientException { + try (var ignored = new DefaultPulsarConsumerFactory<>(pulsarClient, + List.of(configCustomizer2, configCustomizer1)) + .createConsumer(SCHEMA, List.of("topic0"), "dft-sub", createConsumerCustomizer)) { + InOrder inOrder = inOrder(configCustomizer1, configCustomizer2, createConsumerCustomizer); + inOrder.verify(configCustomizer2).customize(any(ConsumerBuilder.class)); + inOrder.verify(configCustomizer1).customize(any(ConsumerBuilder.class)); + inOrder.verify(createConsumerCustomizer).customize(any(ConsumerBuilder.class)); + } + } + + } + } } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java index 18fd6afc..0fd10089 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java @@ -18,6 +18,9 @@ package org.springframework.pulsar.core; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; import java.util.Collections; import java.util.List; @@ -28,11 +31,13 @@ import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Reader; +import org.apache.pulsar.client.api.ReaderBuilder; import org.apache.pulsar.client.api.Schema; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; +import org.mockito.InOrder; import org.springframework.pulsar.test.support.PulsarTestContainerSupport; @@ -40,6 +45,7 @@ import org.springframework.pulsar.test.support.PulsarTestContainerSupport; * Testing {@link DefaultPulsarReaderFactory}. * * @author Soby Chacko + * @author Chris Bono */ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSupport { @@ -129,10 +135,10 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo @Test void useFactoryDefaults() throws Exception { - pulsarReaderFactory = new DefaultPulsarReaderFactory<>(pulsarClient, (readerBuilder) -> { + pulsarReaderFactory = new DefaultPulsarReaderFactory<>(pulsarClient, List.of((readerBuilder) -> { readerBuilder.topic("basic-pulsar-reader-topic"); readerBuilder.startMessageId(MessageId.earliest); - }); + })); // The following code expects the above topic and startMessageId to be used Message message; try (Reader reader = pulsarReaderFactory.createReader(null, null, Schema.STRING, @@ -148,10 +154,10 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo @Test void overrideFactoryDefaults() throws Exception { - pulsarReaderFactory = new DefaultPulsarReaderFactory<>(pulsarClient, (readerBuilder) -> { + pulsarReaderFactory = new DefaultPulsarReaderFactory<>(pulsarClient, List.of((readerBuilder) -> { readerBuilder.topic("foo-topic"); readerBuilder.startMessageId(MessageId.latest); - }); + })); // The following code expects the above topic and startMessageId to be ignored // (overridden) Message message; @@ -186,6 +192,43 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo } + @Nested + @SuppressWarnings("unchecked") + class DefaultConfigCustomizerApi { + + private ReaderBuilderCustomizer configCustomizer1 = mock(ReaderBuilderCustomizer.class); + + private ReaderBuilderCustomizer configCustomizer2 = mock(ReaderBuilderCustomizer.class); + + private ReaderBuilderCustomizer createReaderCustomizer = mock(ReaderBuilderCustomizer.class); + + @Test + void singleConfigCustomizer() throws Exception { + try (var ignored = new DefaultPulsarReaderFactory<>(pulsarClient, + List.of(configCustomizer2, configCustomizer1)) + .createReader(List.of("basic-pulsar-reader-topic"), MessageId.earliest, Schema.STRING, + List.of(createReaderCustomizer))) { + InOrder inOrder = inOrder(configCustomizer1, createReaderCustomizer); + inOrder.verify(configCustomizer1).customize(any(ReaderBuilder.class)); + inOrder.verify(createReaderCustomizer).customize(any(ReaderBuilder.class)); + } + } + + @Test + void multipleConfigCustomizers() throws Exception { + try (var ignored = new DefaultPulsarReaderFactory<>(pulsarClient, + List.of(configCustomizer2, configCustomizer1)) + .createReader(List.of("basic-pulsar-reader-topic"), MessageId.earliest, Schema.STRING, + List.of(createReaderCustomizer))) { + InOrder inOrder = inOrder(configCustomizer1, configCustomizer2, createReaderCustomizer); + inOrder.verify(configCustomizer2).customize(any(ReaderBuilder.class)); + inOrder.verify(configCustomizer1).customize(any(ReaderBuilder.class)); + inOrder.verify(createReaderCustomizer).customize(any(ReaderBuilder.class)); + } + } + + } + @Nested class MissingConfig { diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java index c4d01c42..ad6ee2db 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java @@ -55,10 +55,10 @@ class FailoverConsumerTests implements PulsarTestContainerSupport { .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("my-part-topic-1"); consumerBuilder.subscriptionName("my-part-subscription-1"); - }); + })); CountDownLatch latch1 = new CountDownLatch(1); CountDownLatch latch2 = new CountDownLatch(1); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/SharedSubscriptionConsumerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/SharedSubscriptionConsumerTests.java index 5c2b0c8b..7bb119c7 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/SharedSubscriptionConsumerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/SharedSubscriptionConsumerTests.java @@ -57,10 +57,10 @@ public class SharedSubscriptionConsumerTests implements PulsarTestContainerSuppo try { pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, (consumerBuilder) -> { + pulsarClient, List.of((consumerBuilder) -> { consumerBuilder.topic("shared-subscription-single-msg-test-topic"); consumerBuilder.subscriptionName("shared-subscription-single-msg-test-sub"); - }); + })); CountDownLatch latch1 = new CountDownLatch(1); CountDownLatch latch2 = new CountDownLatch(1); @@ -114,10 +114,10 @@ public class SharedSubscriptionConsumerTests implements PulsarTestContainerSuppo try { pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); DefaultPulsarConsumerFactory consumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("key-shared-batch-disabled-topic"); consumerBuilder.subscriptionName("key-shared-batch-disabled-sub"); - }); + })); CountDownLatch latch = new CountDownLatch(30); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java index 1547d6ac..285ef285 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java @@ -58,10 +58,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-1"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-1"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); PulsarRecordMessageListener messageListener = mock(PulsarRecordMessageListener.class); @@ -108,10 +108,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-2"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-2"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); PulsarRecordMessageListener messageListener = mock(PulsarRecordMessageListener.class); @@ -155,10 +155,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-3"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-3"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); PulsarRecordMessageListener messageListener = mock(PulsarRecordMessageListener.class); @@ -213,10 +213,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-4"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-4"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); @@ -285,10 +285,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-5"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-5"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); @@ -355,10 +355,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-6"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-6"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); @@ -425,10 +425,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-7"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-7"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); @@ -493,10 +493,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("default-error-handler-tests-8"); consumerBuilder.subscriptionName("default-error-handler-tests-sub-8"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java index db004c9c..5ed3d703 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java @@ -67,10 +67,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("dpmlct-012"); consumerBuilder.subscriptionName("dpmlct-sb-012"); - }); + })); CountDownLatch latch = new CountDownLatch(1); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties @@ -96,10 +96,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("containerPauseResumeWaitNotify-topic"); consumerBuilder.subscriptionName("containerPauseResumeWaitNotify-sub"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> { }); @@ -166,11 +166,11 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("dpmlct-013"); consumerBuilder.subscriptionName("dpmlct-sb-013"); consumerBuilder.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest); - }); + })); CountDownLatch latch = new CountDownLatch(5); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties @@ -198,10 +198,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("dpmlct-014"); consumerBuilder.subscriptionName("dpmlct-sb-014"); - }); + })); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); List messages = new ArrayList<>(); pulsarContainerProperties.setMessageListener( @@ -236,11 +236,11 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .maxDelayMs(5 * 1000) .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( - new DefaultPulsarConsumerFactory<>(pulsarClient, (consumerBuilder) -> { + new DefaultPulsarConsumerFactory<>(pulsarClient, List.of((consumerBuilder) -> { consumerBuilder.topic("dpmlct-015"); consumerBuilder.subscriptionName("dpmlct-sb-015"); consumerBuilder.negativeAckRedeliveryBackoff(redeliveryBackoff); - })); + }))); CountDownLatch latch = new CountDownLatch(10); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> { @@ -286,12 +286,12 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .deadLetterTopic("dpmlct-016-dlq-topic") .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("dpmlct-016"); consumerBuilder.subscriptionName("dpmlct-sb-016"); consumerBuilder.negativeAckRedeliveryDelay(1L, TimeUnit.SECONDS); consumerBuilder.deadLetterPolicy(deadLetterPolicy); - }); + })); CountDownLatch dlqLatch = new CountDownLatch(1); CountDownLatch latch = new CountDownLatch(6); @@ -345,12 +345,12 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .deadLetterTopic("dlq-topic") .build(); DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, - (consumerBuilder) -> { + List.of((consumerBuilder) -> { consumerBuilder.topic("dpmlct-017"); consumerBuilder.subscriptionName("dpmlct-sb-017"); consumerBuilder.negativeAckRedeliveryDelay(1L, TimeUnit.SECONDS); consumerBuilder.deadLetterPolicy(deadLetterPolicy); - }); + })); CountDownLatch dlqLatch = new CountDownLatch(1); CountDownLatch latch = new CountDownLatch(6); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainerTests.java index 1bf91ebc..1e32bad9 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainerTests.java @@ -69,7 +69,7 @@ public class DefaultPulsarMessageReaderContainerTests implements PulsarTestConta var latch = new CountDownLatch(1); DefaultPulsarReaderFactory pulsarReaderFactory = new DefaultPulsarReaderFactory<>(pulsarClient, - (readerBuilder -> { + List.of((readerBuilder) -> { readerBuilder.topic("dprlct-001"); readerBuilder.subscriptionName("dprlct-sub-001"); }));