Allow multiple customizers for several components (#435)

* DefaultPulsarReaderFactory accepts multiple customizers
* DefaultPulsarConsumerFactory accepts multiple customizers
* DefaultReactivePulsarSenderFactory accepts multiple customizers

See #432
This commit is contained in:
Chris Bono
2023-08-20 11:46:26 -05:00
committed by GitHub
parent 0b9baf8b79
commit 495de957f9
17 changed files with 237 additions and 82 deletions

View File

@@ -55,23 +55,49 @@ public class DefaultReactivePulsarSenderFactory<T> implements ReactivePulsarSend
@Nullable
private final ReactiveMessageSenderCache reactiveMessageSenderCache;
@Nullable
private final List<ReactiveMessageSenderBuilderCustomizer<T>> 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<ReactiveMessageSenderBuilderCustomizer<T>> 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<ReactiveMessageSenderBuilderCustomizer<T>> 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<T> implements ReactivePulsarSend
ReactiveMessageSenderBuilder<T> 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);

View File

@@ -37,7 +37,7 @@ import org.junit.jupiter.api.Test;
* @author Christophe Bornet
* @author Chris Bono
*/
class DefaultReactiveMessageConsumerFactoryTests {
class DefaultReactivePulsarConsumerFactoryTests {
private static final Schema<String> SCHEMA = Schema.STRING;

View File

@@ -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<String> schema = Schema.STRING;

View File

@@ -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<String> schema = Schema.STRING;
@@ -66,17 +71,17 @@ class DefaultReactiveMessageSenderFactoryTests {
}
private ReactivePulsarSenderFactory<String> newSenderFactory() {
return new DefaultReactivePulsarSenderFactory<>((PulsarClient) null, null, null);
return new DefaultReactivePulsarSenderFactory<>(null, null, null, null);
}
private ReactivePulsarSenderFactory<String> 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<String> 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<String> configCustomizer1 = mock(
ReactiveMessageSenderBuilderCustomizer.class);
private ReactiveMessageSenderBuilderCustomizer<String> configCustomizer2 = mock(
ReactiveMessageSenderBuilderCustomizer.class);
private ReactiveMessageSenderBuilderCustomizer<String> 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<String> newSenderFactoryWithCustomizers(
List<ReactiveMessageSenderBuilderCustomizer<String>> customizers) {
MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec();
return new DefaultReactivePulsarSenderFactory<>(null, senderSpec, null, customizers);
}
}
}

View File

@@ -182,7 +182,8 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
if (producerFactoryHasDefaultTopic) {
spec.setTopicName("fake-topic");
}
ReactivePulsarSenderFactory<Foo> producerFactory = new DefaultReactivePulsarSenderFactory<>(client, spec, null);
ReactivePulsarSenderFactory<Foo> 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<String> senderFactory = new DefaultReactivePulsarSenderFactory<>(client,
new MutableReactiveMessageSenderSpec(), null);
new MutableReactiveMessageSenderSpec(), null, null);
ReactivePulsarTemplate<String> 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<T> senderFactory = new DefaultReactivePulsarSenderFactory<>(client, senderSpec,
null);
null, null);
ReactivePulsarTemplate<T> pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory);
@@ -260,7 +261,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
MutableReactiveMessageSenderSpec spec = new MutableReactiveMessageSenderSpec();
spec.setTopicName(topic);
ReactivePulsarSenderFactory<Foo> 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<String> senderFactory = new DefaultReactivePulsarSenderFactory<>(client, spec,
null);
null, null);
ReactivePulsarTemplate<String> 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<String> senderFactory = new DefaultReactivePulsarSenderFactory<>(client,
new MutableReactiveMessageSenderSpec(), null);
new MutableReactiveMessageSenderSpec(), null, null);
ReactivePulsarTemplate<String> pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory);
assertThatIllegalArgumentException()
.isThrownBy(() -> pulsarTemplate.send("sendNullWithoutSchemaFails", (String) null, null).subscribe())

View File

@@ -82,7 +82,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo
MutableReactiveMessageSenderSpec prodConfig = new MutableReactiveMessageSenderSpec();
prodConfig.setTopicName(topic);
DefaultReactivePulsarSenderFactory<String> pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>(
reactivePulsarClient, prodConfig, null, new DefaultTopicResolver());
reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver());
ReactivePulsarTemplate<String> 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<String> pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>(
reactivePulsarClient, prodConfig, null, new DefaultTopicResolver());
reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver());
ReactivePulsarTemplate<String> 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<String> pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>(
reactivePulsarClient, prodConfig, null, new DefaultTopicResolver());
reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver());
ReactivePulsarTemplate<String> 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<String> pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>(
reactivePulsarClient, prodConfig, null, new DefaultTopicResolver());
reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver());
ReactivePulsarTemplate<String> 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<String> pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>(
reactivePulsarClient, prodConfig, null, new DefaultTopicResolver());
reactivePulsarClient, prodConfig, null, null, new DefaultTopicResolver());
ReactivePulsarTemplate<String> 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();

View File

@@ -48,18 +48,18 @@ public class DefaultPulsarConsumerFactory<T> implements PulsarConsumerFactory<T>
private final PulsarClient pulsarClient;
@Nullable
private final ConsumerBuilderCustomizer<T> defaultConfigCustomizer;
private final List<ConsumerBuilderCustomizer<T>> 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<T> defaultConfigCustomizer) {
List<ConsumerBuilderCustomizer<T>> defaultConfigCustomizers) {
this.pulsarClient = pulsarClient;
this.defaultConfigCustomizer = defaultConfigCustomizer;
this.defaultConfigCustomizers = defaultConfigCustomizers;
}
@Override
@@ -77,8 +77,8 @@ public class DefaultPulsarConsumerFactory<T> implements PulsarConsumerFactory<T>
ConsumerBuilder<T> 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);

View File

@@ -142,7 +142,7 @@ public class DefaultPulsarProducerFactory<T> implements PulsarProducerFactory<T>
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);

View File

@@ -43,7 +43,7 @@ public class DefaultPulsarReaderFactory<T> implements PulsarReaderFactory<T> {
private final PulsarClient pulsarClient;
@Nullable
private final ReaderBuilderCustomizer<T> defaultConfigCustomizer;
private final List<ReaderBuilderCustomizer<T>> defaultConfigCustomizers;
/**
* Construct a reader factory instance with no default configuration.
@@ -56,13 +56,13 @@ public class DefaultPulsarReaderFactory<T> implements PulsarReaderFactory<T> {
/**
* 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<T> defaultConfigCustomizer) {
@Nullable List<ReaderBuilderCustomizer<T>> defaultConfigCustomizers) {
this.pulsarClient = pulsarClient;
this.defaultConfigCustomizer = defaultConfigCustomizer;
this.defaultConfigCustomizers = defaultConfigCustomizers;
}
@Override
@@ -72,8 +72,8 @@ public class DefaultPulsarReaderFactory<T> implements PulsarReaderFactory<T> {
ReaderBuilder<T> 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)) {

View File

@@ -358,11 +358,11 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
pulsarClient.close();
}
private <T> ConsumerBuilderCustomizer<T> defaultConfig(String topicName, String subscriptionName) {
return (consumerBuilder) -> {
private <T> List<ConsumerBuilderCustomizer<T>> defaultConfig(String topicName, String subscriptionName) {
return List.of((consumerBuilder) -> {
consumerBuilder.topic(topicName);
consumerBuilder.subscriptionName(subscriptionName);
};
});
}
}

View File

@@ -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<String> configCustomizer1 = mock(ConsumerBuilderCustomizer.class);
private ConsumerBuilderCustomizer<String> configCustomizer2 = mock(ConsumerBuilderCustomizer.class);
private ConsumerBuilderCustomizer<String> 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));
}
}
}
}
}

View File

@@ -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<String> message;
try (Reader<String> 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<String> message;
@@ -186,6 +192,43 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo
}
@Nested
@SuppressWarnings("unchecked")
class DefaultConfigCustomizerApi {
private ReaderBuilderCustomizer<String> configCustomizer1 = mock(ReaderBuilderCustomizer.class);
private ReaderBuilderCustomizer<String> configCustomizer2 = mock(ReaderBuilderCustomizer.class);
private ReaderBuilderCustomizer<String> 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 {

View File

@@ -55,10 +55,10 @@ class FailoverConsumerTests implements PulsarTestContainerSupport {
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> 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);

View File

@@ -57,10 +57,10 @@ public class SharedSubscriptionConsumerTests implements PulsarTestContainerSuppo
try {
pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
DefaultPulsarConsumerFactory<String> 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<String> 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);

View File

@@ -58,10 +58,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> 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<String> 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<Integer> 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<Integer> 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<Integer> 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<Integer> 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<Integer> 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<Integer> 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);

View File

@@ -67,10 +67,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> 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<String> 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<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,
(consumerBuilder) -> {
List.of((consumerBuilder) -> {
consumerBuilder.topic("dpmlct-014");
consumerBuilder.subscriptionName("dpmlct-sb-014");
});
}));
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
List<String> messages = new ArrayList<>();
pulsarContainerProperties.setMessageListener(
@@ -236,11 +236,11 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
.maxDelayMs(5 * 1000)
.build();
DefaultPulsarConsumerFactory<String> 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<Integer> 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<Integer> 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);

View File

@@ -69,7 +69,7 @@ public class DefaultPulsarMessageReaderContainerTests implements PulsarTestConta
var latch = new CountDownLatch(1);
DefaultPulsarReaderFactory<String> pulsarReaderFactory = new DefaultPulsarReaderFactory<>(pulsarClient,
(readerBuilder -> {
List.of((readerBuilder) -> {
readerBuilder.topic("dprlct-001");
readerBuilder.subscriptionName("dprlct-sub-001");
}));