From 0b9baf8b79bb60080d35566aa520b18e4d28f165 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Sat, 19 Aug 2023 18:08:19 -0500 Subject: [PATCH] DefaultPulsarProducerFactory accepts multiple customizers (#434) See #432 --- .../core/CachingPulsarProducerFactory.java | 7 +-- .../core/DefaultPulsarProducerFactory.java | 26 +++++------ .../CachingPulsarProducerFactoryTests.java | 4 +- .../DefaultPulsarProducerFactoryTests.java | 46 ++++++++++++++++++- .../pulsar/core/FailoverConsumerTests.java | 3 +- .../core/PulsarProducerFactoryTests.java | 13 ++++-- .../core/SharedSubscriptionConsumerTests.java | 3 +- .../pulsar/listener/PulsarListenerTests.java | 4 +- ...aultPulsarMessageReaderContainerTests.java | 6 +-- 9 files changed, 82 insertions(+), 30 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java index b51b58b8..a79294f3 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java @@ -67,16 +67,17 @@ public class CachingPulsarProducerFactory extends DefaultPulsarProducerFactor * configuration. * @param pulsarClient the client used to create the producers * @param defaultTopic the default topic to use for the producers - * @param defaultConfigCustomizer the default configuration to apply to the producers + * @param defaultConfigCustomizers the optional list of customizers to apply to the + * created producers * @param topicResolver the topic resolver to use * @param cacheExpireAfterAccess time period to expire unused entries in the cache * @param cacheMaximumSize maximum size of cache (entries) * @param cacheInitialCapacity the initial size of cache */ public CachingPulsarProducerFactory(PulsarClient pulsarClient, @Nullable String defaultTopic, - ProducerBuilderCustomizer defaultConfigCustomizer, TopicResolver topicResolver, + List> defaultConfigCustomizers, TopicResolver topicResolver, Duration cacheExpireAfterAccess, Long cacheMaximumSize, Integer cacheInitialCapacity) { - super(pulsarClient, defaultTopic, defaultConfigCustomizer, topicResolver); + super(pulsarClient, defaultTopic, defaultConfigCustomizers, topicResolver); var cacheFactory = CacheProviderFactory., Producer>load(); this.producerCache = cacheFactory.create(cacheExpireAfterAccess, cacheMaximumSize, cacheInitialCapacity, (key, producer, cause) -> { 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 26f6d473..8e24a901 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 @@ -52,7 +52,7 @@ public class DefaultPulsarProducerFactory implements PulsarProducerFactory private final String defaultTopic; @Nullable - private final ProducerBuilderCustomizer defaultConfigCustomizer; + private final List> defaultConfigCustomizers; private final TopicResolver topicResolver; @@ -61,8 +61,7 @@ public class DefaultPulsarProducerFactory implements PulsarProducerFactory * @param pulsarClient the client used to create the producers */ public DefaultPulsarProducerFactory(PulsarClient pulsarClient) { - this(pulsarClient, null, (pb) -> { - }, new DefaultTopicResolver()); + this(pulsarClient, null, null, new DefaultTopicResolver()); } /** @@ -71,33 +70,34 @@ public class DefaultPulsarProducerFactory implements PulsarProducerFactory * @param defaultTopic the default topic to use for the producers */ public DefaultPulsarProducerFactory(PulsarClient pulsarClient, @Nullable String defaultTopic) { - this(pulsarClient, defaultTopic, (pb) -> { - }, new DefaultTopicResolver()); + this(pulsarClient, defaultTopic, null, new DefaultTopicResolver()); } /** * Construct a producer factory that uses a default topic resolver. * @param pulsarClient the client used to create the producers * @param defaultTopic the default topic to use for the producers - * @param defaultConfigCustomizer the default configuration to apply to the producers + * @param defaultConfigCustomizers the optional list of customizers to apply to the + * created producers */ public DefaultPulsarProducerFactory(PulsarClient pulsarClient, @Nullable String defaultTopic, - @Nullable ProducerBuilderCustomizer defaultConfigCustomizer) { - this(pulsarClient, defaultTopic, defaultConfigCustomizer, new DefaultTopicResolver()); + @Nullable List> defaultConfigCustomizers) { + this(pulsarClient, defaultTopic, defaultConfigCustomizers, new DefaultTopicResolver()); } /** * Construct a producer factory that uses the specified parameters. * @param pulsarClient the client used to create the producers * @param defaultTopic the default topic to use for the producers - * @param defaultConfigCustomizer the default configuration to apply to the producers + * @param defaultConfigCustomizers the optional list of customizers to apply to the + * created producers * @param topicResolver the topic resolver to use */ public DefaultPulsarProducerFactory(PulsarClient pulsarClient, @Nullable String defaultTopic, - @Nullable ProducerBuilderCustomizer defaultConfigCustomizer, TopicResolver topicResolver) { + @Nullable List> defaultConfigCustomizers, TopicResolver topicResolver) { this.pulsarClient = pulsarClient; this.defaultTopic = defaultTopic; - this.defaultConfigCustomizer = defaultConfigCustomizer; + this.defaultConfigCustomizers = defaultConfigCustomizers; this.topicResolver = topicResolver; } @@ -142,8 +142,8 @@ public class DefaultPulsarProducerFactory implements PulsarProducerFactory var producerBuilder = this.pulsarClient.newProducer(schema); // Apply the default config customizer (preserve the topic) - if (this.defaultConfigCustomizer != null) { - this.defaultConfigCustomizer.customize(producerBuilder); + if (this.defaultConfigCustomizers != null) { + this.defaultConfigCustomizers.forEach((customizer) -> customizer.customize(producerBuilder)); } producerBuilder.topic(resolvedTopic); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java index d1fb2483..2e9681a0 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java @@ -228,9 +228,9 @@ class CachingPulsarProducerFactoryTests extends PulsarProducerFactoryTests { @Override protected CachingPulsarProducerFactory producerFactory(PulsarClient pulsarClient, - @Nullable String defaultTopic, @Nullable ProducerBuilderCustomizer defaultConfigCustomizer) { + @Nullable String defaultTopic, @Nullable List> defaultConfigCustomizers) { var producerFactory = new CachingPulsarProducerFactory(pulsarClient, defaultTopic, - defaultConfigCustomizer, new DefaultTopicResolver(), Duration.ofMinutes(5L), 30L, 2); + defaultConfigCustomizers, new DefaultTopicResolver(), Duration.ofMinutes(5L), 30L, 2); producerFactories.add(producerFactory); return producerFactory; } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarProducerFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarProducerFactoryTests.java index 71b157fe..fca5f43c 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarProducerFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarProducerFactoryTests.java @@ -17,11 +17,19 @@ package org.springframework.pulsar.core; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; + +import java.util.List; import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.ProducerBuilder; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; +import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; +import org.mockito.InOrder; import org.springframework.lang.Nullable; @@ -46,8 +54,42 @@ class DefaultPulsarProducerFactoryTests extends PulsarProducerFactoryTests { @Override protected PulsarProducerFactory producerFactory(PulsarClient pulsarClient, @Nullable String defaultTopic, - @Nullable ProducerBuilderCustomizer defaultConfigCustomizer) { - return new DefaultPulsarProducerFactory<>(pulsarClient, defaultTopic, defaultConfigCustomizer); + @Nullable List> defaultConfigCustomizers) { + return new DefaultPulsarProducerFactory<>(pulsarClient, defaultTopic, defaultConfigCustomizers); + } + + @Nested + @SuppressWarnings("unchecked") + class DefaultConfigCustomizerApi { + + private ProducerBuilderCustomizer configCustomizer1 = mock(ProducerBuilderCustomizer.class); + + private ProducerBuilderCustomizer configCustomizer2 = mock(ProducerBuilderCustomizer.class); + + private ProducerBuilderCustomizer createProducerCustomizer = mock(ProducerBuilderCustomizer.class); + + @Test + void singleConfigCustomizer() throws PulsarClientException { + try (var ignored = newProducerFactoryWithDefaultConfigCustomizers(List.of(configCustomizer1)) + .createProducer(schema, "topic0", createProducerCustomizer)) { + InOrder inOrder = inOrder(configCustomizer1, createProducerCustomizer); + inOrder.verify(configCustomizer1).customize(any(ProducerBuilder.class)); + inOrder.verify(createProducerCustomizer).customize(any(ProducerBuilder.class)); + } + } + + @Test + void multipleConfigCustomizers() throws PulsarClientException { + try (var ignored = newProducerFactoryWithDefaultConfigCustomizers( + List.of(configCustomizer2, configCustomizer1)) + .createProducer(schema, "topic0", createProducerCustomizer)) { + InOrder inOrder = inOrder(configCustomizer1, configCustomizer2, createProducerCustomizer); + inOrder.verify(configCustomizer2).customize(any(ProducerBuilder.class)); + inOrder.verify(configCustomizer1).customize(any(ProducerBuilder.class)); + inOrder.verify(createProducerCustomizer).customize(any(ProducerBuilder.class)); + } + } + } } 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 5556e714..c4d01c42 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 @@ -19,6 +19,7 @@ package org.springframework.pulsar.core; import static org.assertj.core.api.Assertions.assertThat; import java.io.Serial; +import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -85,7 +86,7 @@ class FailoverConsumerTests implements PulsarTestContainerSupport { container3.start(); DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, - "my-part-topic-1", (pb) -> pb.messageRoutingMode(MessageRoutingMode.CustomPartition)); + "my-part-topic-1", List.of((pb) -> pb.messageRoutingMode(MessageRoutingMode.CustomPartition))); PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.newMessage("hello john doe") diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarProducerFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarProducerFactoryTests.java index 83d03ecd..28636b0f 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarProducerFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarProducerFactoryTests.java @@ -26,6 +26,7 @@ import static org.mockito.Mockito.verify; import java.util.Arrays; import java.util.Collections; +import java.util.List; import java.util.Set; import org.apache.pulsar.client.api.Producer; @@ -100,7 +101,12 @@ abstract class PulsarProducerFactoryTests implements PulsarTestContainerSupport } private PulsarProducerFactory newProducerFactoryWithDefaultKeys(Set defaultKeys) { - return producerFactory(pulsarClient, null, (pb) -> defaultKeys.forEach(pb::addEncryptionKey)); + return producerFactory(pulsarClient, null, List.of((pb) -> defaultKeys.forEach(pb::addEncryptionKey))); + } + + protected PulsarProducerFactory newProducerFactoryWithDefaultConfigCustomizers( + List> customizers) { + return producerFactory(pulsarClient, null, customizers); } /** @@ -116,11 +122,12 @@ abstract class PulsarProducerFactoryTests implements PulsarTestContainerSupport * Subclasses override to provide concrete {@link PulsarProducerFactory} instance. * @param pulsarClient the Pulsar client * @param defaultTopic the default topic to use for the producers - * @param defaultConfigCustomizer the default configuration to apply to the producers + * @param defaultConfigCustomizers the optional list of customizers to apply to the + * created producers * @return a Pulsar producer factory instance to use for the tests */ protected abstract PulsarProducerFactory producerFactory(PulsarClient pulsarClient, - @Nullable String defaultTopic, @Nullable ProducerBuilderCustomizer defaultConfigCustomizer); + @Nullable String defaultTopic, @Nullable List> defaultConfigCustomizers); @Test @SuppressWarnings("unchecked") 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 9bbc78e1..5c2b0c8b 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 @@ -19,6 +19,7 @@ package org.springframework.pulsar.core; import static org.assertj.core.api.Assertions.assertThat; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -133,7 +134,7 @@ public class SharedSubscriptionConsumerTests implements PulsarTestContainerSuppo Thread.sleep(5_000); DefaultPulsarProducerFactory producerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, - "key-shared-batch-disabled-topic", (pb) -> pb.enableBatching(false)); + "key-shared-batch-disabled-topic", List.of((pb) -> pb.enableBatching(false))); PulsarTemplate pulsarTemplate = new PulsarTemplate<>(producerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.newMessage("alice-" + i) diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java index d324bec5..893e3725 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java @@ -177,7 +177,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { void concurrencyOnPulsarListenerWithFailoverSubscription(@Autowired PulsarListenerEndpointRegistry registry) throws Exception { var pulsarProducerFactory = new DefaultPulsarProducerFactory(pulsarClient, null, - (pb) -> pb.enableBatching(false)); + List.of((pb) -> pb.enableBatching(false))); var customTemplate = new PulsarTemplate<>(pulsarProducerFactory); var bar = (ConcurrentPulsarMessageListenerContainer) registry.getListenerContainer("bar"); @@ -194,7 +194,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { void nonDefaultConcurrencySettingNotAllowedOnExclusiveSubscriptions( @Autowired PulsarListenerEndpointRegistry registry) throws Exception { var pulsarProducerFactory = new DefaultPulsarProducerFactory(pulsarClient, null, - (pb) -> pb.enableBatching(false)); + List.of((pb) -> pb.enableBatching(false))); var customTemplate = new PulsarTemplate<>(pulsarProducerFactory); var bar = (ConcurrentPulsarMessageListenerContainer) registry.getListenerContainer("bar"); 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 3af04a72..1bf91ebc 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 @@ -87,7 +87,7 @@ public class DefaultPulsarMessageReaderContainerTests implements PulsarTestConta container.start(); DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, "dprlct-001", (pb) -> pb.topic("dprlct-001")); + pulsarClient, "dprlct-001", List.of((pb) -> pb.topic("dprlct-001"))); PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.sendAsync("hello john doe"); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -115,7 +115,7 @@ public class DefaultPulsarMessageReaderContainerTests implements PulsarTestConta container = new DefaultPulsarMessageReaderContainer<>(pulsarReaderFactory, containerProps); container.start(); DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, "dprlct-002", (pb) -> pb.topic("dprlct-002")); + pulsarClient, "dprlct-002", List.of((pb) -> pb.topic("dprlct-002"))); PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.sendAsync("hello buzz doe"); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -144,7 +144,7 @@ public class DefaultPulsarMessageReaderContainerTests implements PulsarTestConta var prodConfig = Map.of("topicName", "dprlct-003"); var producerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, "dprlct-003", - (pb) -> pb.topic("dprlct-003")); + List.of((pb) -> pb.topic("dprlct-003"))); var pulsarTemplate = new PulsarTemplate<>(producerFactory); // The following sends will not be received by the reader as we are using the