From 1758fe87df5a6af396d33e9f2aa1ba2af2513a86 Mon Sep 17 00:00:00 2001 From: Christophe Bornet Date: Thu, 27 Oct 2022 21:44:20 +0200 Subject: [PATCH] Rework ReactivePulsarSenderFactory (#179) * Ensure DefaultReactivePulsarSenderFactory is not null * Simplify API naming * Remove MessageRouter specific methods as it can be configured in the ReactiveMessageSenderSpec (which is not the case in the imperative Producer conf properties) --- .../DefaultReactivePulsarSenderFactory.java | 31 ++++----- .../reactive/ReactivePulsarSenderFactory.java | 16 +---- .../ReactivePulsarSenderOperations.java | 8 --- .../ReactivePulsarSenderTemplate.java | 30 +++------ ...aultReactiveMessageSenderFactoryTests.java | 66 +++++++------------ .../reactive/ReactivePulsarTemplateTests.java | 45 +------------ 6 files changed, 49 insertions(+), 147 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java index 48e4e719..e9117bf4 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java @@ -18,10 +18,11 @@ package org.springframework.pulsar.core.reactive; import java.util.List; -import org.apache.pulsar.client.api.MessageRouter; 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.ImmutableReactiveMessageSenderSpec; +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; @@ -58,42 +59,32 @@ public class DefaultReactivePulsarSenderFactory implements ReactivePulsarSend ReactiveMessageSenderSpec reactiveMessageSenderSpec, ReactiveMessageSenderCache reactiveMessageSenderCache) { this.reactivePulsarClient = reactivePulsarClient; - this.reactiveMessageSenderSpec = reactiveMessageSenderSpec; + this.reactiveMessageSenderSpec = new ImmutableReactiveMessageSenderSpec( + reactiveMessageSenderSpec != null ? reactiveMessageSenderSpec : new MutableReactiveMessageSenderSpec()); this.reactiveMessageSenderCache = reactiveMessageSenderCache; } @Override - public ReactiveMessageSender createReactiveMessageSender(String topic, Schema schema) { - return doCreateReactiveMessageSender(topic, schema, null, null); + public ReactiveMessageSender createSender(String topic, Schema schema) { + return doCreateReactiveMessageSender(topic, schema, null); } @Override - public ReactiveMessageSender createReactiveMessageSender(String topic, Schema schema, - MessageRouter messageRouter) { - return doCreateReactiveMessageSender(topic, schema, messageRouter, null); - } - - @Override - public ReactiveMessageSender createReactiveMessageSender(String topic, Schema schema, - MessageRouter messageRouter, List> customizers) { - return doCreateReactiveMessageSender(topic, schema, messageRouter, customizers); + public ReactiveMessageSender createSender(String topic, Schema schema, + List> customizers) { + return doCreateReactiveMessageSender(topic, schema, customizers); } private ReactiveMessageSender doCreateReactiveMessageSender(String topic, Schema schema, - MessageRouter messageRouter, List> customizers) { + List> customizers) { final String resolvedTopic = ReactiveMessageSenderUtils.resolveTopicName(topic, this); this.logger.trace(() -> String.format("Creating reactive message sender for '%s' topic", resolvedTopic)); final ReactiveMessageSenderBuilder sender = this.reactivePulsarClient.messageSender(schema); - if (this.reactiveMessageSenderSpec != null) { - sender.applySpec(this.reactiveMessageSenderSpec); - } + sender.applySpec(this.reactiveMessageSenderSpec); sender.topic(resolvedTopic); if (this.reactiveMessageSenderCache != null) { sender.cache(this.reactiveMessageSenderCache); } - if (messageRouter != null) { - sender.messageRouter(messageRouter); - } if (!CollectionUtils.isEmpty(customizers)) { customizers.forEach((c) -> c.customize(sender)); } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderFactory.java index d9c801e8..c11ce923 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderFactory.java @@ -18,7 +18,6 @@ package org.springframework.pulsar.core.reactive; import java.util.List; -import org.apache.pulsar.client.api.MessageRouter; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.reactive.client.api.ReactiveMessageSender; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; @@ -38,29 +37,18 @@ public interface ReactivePulsarSenderFactory { * @param schema the schema of the messages to be sent * @return the reactive message sender */ - ReactiveMessageSender createReactiveMessageSender(String topic, Schema schema); + ReactiveMessageSender createSender(String topic, Schema schema); /** * Create a reactive message sender. * @param topic the topic the reactive message sender will send messages to or * {@code null} to use the default topic * @param schema the schema of the messages to be sent - * @param messageRouter the optional message router to use - * @return the reactive message sender - */ - ReactiveMessageSender createReactiveMessageSender(String topic, Schema schema, MessageRouter messageRouter); - - /** - * Create a reactive message sender. - * @param topic the topic the reactive message sender will send messages to or - * {@code null} to use the default topic - * @param schema the schema of the messages to be sent - * @param messageRouter the optional message router to use * @param customizers the optional list of customizers to apply to the reactive * message sender builder * @return the reactive message sender */ - ReactiveMessageSender createReactiveMessageSender(String topic, Schema schema, MessageRouter messageRouter, + ReactiveMessageSender createSender(String topic, Schema schema, List> customizers); /** diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderOperations.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderOperations.java index 04419b07..4b2a5188 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderOperations.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderOperations.java @@ -17,7 +17,6 @@ package org.springframework.pulsar.core.reactive; import org.apache.pulsar.client.api.MessageId; -import org.apache.pulsar.client.api.MessageRouter; import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; @@ -95,13 +94,6 @@ public interface ReactivePulsarSenderOperations { */ SendMessageBuilder withMessageCustomizer(MessageSpecBuilderCustomizer customizer); - /** - * Specifies the custom message router to use when sending the message. - * @param messageRouter the custom message router - * @return the current builder with the custom message router specified - */ - SendMessageBuilder withCustomRouter(MessageRouter messageRouter); - /** * Specifies the customizer to use to further configure the reactive sender * builder. diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderTemplate.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderTemplate.java index 7f901e28..d4cd5f9a 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderTemplate.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderTemplate.java @@ -19,7 +19,6 @@ package org.springframework.pulsar.core.reactive; import java.util.Collections; import org.apache.pulsar.client.api.MessageId; -import org.apache.pulsar.client.api.MessageRouter; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.reactive.client.api.MessageSpec; import org.apache.pulsar.reactive.client.api.MessageSpecBuilder; @@ -62,7 +61,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper @Override public Mono send(String topic, T message) { - return doSend(topic, message, null, null, null); + return doSend(topic, message, null, null); } @Override @@ -72,7 +71,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper @Override public Flux send(String topic, Publisher messages) { - return doSendMany(topic, messages, null, null, null); + return doSendMany(topic, messages, null, null); } @Override @@ -89,13 +88,13 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper } private Mono doSend(String topic, T message, - MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, MessageRouter messageRouter, + MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, ReactiveMessageSenderBuilderCustomizer customizer) { - return doSendMany(topic, Mono.just(message), messageSpecBuilderCustomizer, messageRouter, customizer).single(); + return doSendMany(topic, Mono.just(message), messageSpecBuilderCustomizer, customizer).single(); } private Flux doSendMany(String topic, Publisher messages, - MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, MessageRouter messageRouter, + MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, ReactiveMessageSenderBuilderCustomizer customizer) { final String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory); this.logger.trace(() -> String.format("Sending reactive messages to '%s' topic", topicName)); @@ -108,7 +107,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper * it between messages. So we create one each time and use * ReactiveMessageSender::sendMessage to send messages individually. */ - ReactiveMessageSender sender = createMessageSender(topic, null, messageRouter, customizer); + ReactiveMessageSender sender = createMessageSender(topic, null, customizer); return Flux.from(messages).map(message -> getMessageSpec(messageSpecBuilderCustomizer, message)) .as(sender::sendMessages) .doOnError(ex -> this.logger.error(ex, @@ -117,7 +116,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper msgId -> this.logger.trace(() -> String.format("Sent messages to '%s' topic", topicName))); } return Flux.from(messages).flatMapSequential(message -> { - ReactiveMessageSender sender = createMessageSender(topic, message, messageRouter, customizer); + ReactiveMessageSender sender = createMessageSender(topic, message, customizer); return Mono.just(getMessageSpec(messageSpecBuilderCustomizer, message)).as(sender::sendMessage).doOnError( ex -> this.logger.error(ex, () -> String.format("Failed to send message to '%s' topic", topicName))) .doOnSuccess( @@ -136,10 +135,10 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper return messageSpecBuilder.build(); } - private ReactiveMessageSender createMessageSender(String topic, T message, MessageRouter messageRouter, + private ReactiveMessageSender createMessageSender(String topic, T message, ReactiveMessageSenderBuilderCustomizer customizer) { Schema schema = this.schema != null ? this.schema : SchemaUtils.getSchema(message); - return this.reactiveMessageSenderFactory.createReactiveMessageSender(topic, schema, messageRouter, + return this.reactiveMessageSenderFactory.createSender(topic, schema, customizer == null ? Collections.emptyList() : Collections.singletonList(customizer)); } @@ -153,8 +152,6 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper private MessageSpecBuilderCustomizer messageCustomizer; - private MessageRouter messageRouter; - private ReactiveMessageSenderBuilderCustomizer senderCustomizer; SendMessageBuilderImpl(ReactivePulsarSenderTemplate template, T message) { @@ -174,12 +171,6 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper return this; } - @Override - public SendMessageBuilderImpl withCustomRouter(MessageRouter messageRouter) { - this.messageRouter = messageRouter; - return this; - } - @Override public SendMessageBuilderImpl withSenderCustomizer( ReactiveMessageSenderBuilderCustomizer senderCustomizer) { @@ -189,8 +180,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper @Override public Mono send() { - return this.template.doSend(this.topic, this.message, this.messageCustomizer, this.messageRouter, - this.senderCustomizer); + return this.template.doSend(this.topic, this.message, this.messageCustomizer, this.senderCustomizer); } } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java index 3cd1352a..670acdfe 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/DefaultReactiveMessageSenderFactoryTests.java @@ -18,13 +18,11 @@ package org.springframework.pulsar.core.reactive; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; -import static org.mockito.Mockito.mock; import java.util.Arrays; import java.util.Collections; import java.util.List; -import org.apache.pulsar.client.api.MessageRouter; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory; @@ -33,11 +31,10 @@ import org.apache.pulsar.reactive.client.api.ReactiveMessageSender; 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.ObjectAssert; import org.junit.jupiter.api.Test; /** - * Common tests for {@link DefaultReactivePulsarSenderFactory} + * Tests for {@link DefaultReactivePulsarSenderFactory} * * @author Christophe Bornet */ @@ -47,14 +44,7 @@ class DefaultReactiveMessageSenderFactoryTests { @Test void createSenderWithSpecificTopic() { - testCreateSender(null, null, "topic1", null, null, "topic1", null); - } - - @Test - void createSenderWithSpecificTopicAndMessageRouter() { - MessageRouter router = mock(MessageRouter.class); - - testCreateSender(null, null, "topic1", router, null, "topic1", router); + testCreateSender(null, null, "topic1", null, "topic1"); } @Test @@ -62,59 +52,49 @@ class DefaultReactiveMessageSenderFactoryTests { MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec(); senderSpec.setTopicName("topic0"); - testCreateSender(senderSpec, null, null, null, null, "topic0", null); - } - - @Test - void createSenderWithDefaultTopicAndMessageRouter() { - MutableReactiveMessageSenderSpec senderSpec = new MutableReactiveMessageSenderSpec(); - senderSpec.setTopicName("topic0"); - MessageRouter router = mock(MessageRouter.class); - - testCreateSender(senderSpec, null, null, router, null, "topic0", router); - + testCreateSender(senderSpec, null, null, null, "topic0"); } @Test void createSenderWithSingleSenderCustomizer() { - testCreateSender(null, null, "topic1", null, Collections.singletonList(builder -> builder.topic("topic1")), - "topic1", null); + testCreateSender(null, null, "topic1", Collections.singletonList(builder -> builder.topic("topic1")), "topic1"); } @Test void createSenderWithMultipleSenderCustomizer() { ReactiveMessageSenderBuilderCustomizer customizer1 = builder -> builder.topic("topic1"); - MessageRouter router = mock(MessageRouter.class); - ReactiveMessageSenderBuilderCustomizer customizer2 = builder -> builder.messageRouter(router); + ReactiveMessageSenderCache cache = AdaptedReactivePulsarClientFactory.createCache(); + ReactiveMessageSenderBuilderCustomizer customizer2 = builder -> builder.cache(cache); - testCreateSender(null, null, "topic0", null, Arrays.asList(customizer1, customizer2), "topic1", router); + ReactiveMessageSender sender = testCreateSender(null, null, "topic0", + Arrays.asList(customizer1, customizer2), "topic1"); + assertThat(sender).extracting("producerCache").isSameAs(cache); } @Test void createSenderWithNoTopic() { ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>( (PulsarClient) null, null, null); - assertThatIllegalArgumentException().isThrownBy(() -> senderFactory.createReactiveMessageSender(null, schema)) + assertThatIllegalArgumentException().isThrownBy(() -> senderFactory.createSender(null, schema)) .withMessageContaining("Topic must be specified when no default topic is configured"); } @Test void createSenderWithCache() { - testCreateSender(null, AdaptedReactivePulsarClientFactory.createCache(), "topic1", null, null, "topic1", null); - } - - private void testCreateSender(ReactiveMessageSenderSpec spec, ReactiveMessageSenderCache cache, String topic, - MessageRouter router, List> customizers, - String expectedTopic, MessageRouter expectedRouter) { - ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>( - (PulsarClient) null, spec, cache); - ReactiveMessageSender sender = senderFactory.createReactiveMessageSender(topic, schema, router, - customizers); - ObjectAssert objectAssert = assertThat(sender).extracting("senderSpec") - .asInstanceOf(InstanceOfAssertFactories.type(ReactiveMessageSenderSpec.class)); - objectAssert.extracting(ReactiveMessageSenderSpec::getTopicName).isEqualTo(expectedTopic); - objectAssert.extracting(ReactiveMessageSenderSpec::getMessageRouter).isSameAs(expectedRouter); + ReactiveMessageSenderCache cache = AdaptedReactivePulsarClientFactory.createCache(); + ReactiveMessageSender sender = testCreateSender(null, cache, "topic1", null, "topic1"); assertThat(sender).extracting("producerCache").isSameAs(cache); } + private ReactiveMessageSender testCreateSender(ReactiveMessageSenderSpec spec, + ReactiveMessageSenderCache cache, String topic, + List> customizers, String expectedTopic) { + ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>( + (PulsarClient) null, spec, cache); + ReactiveMessageSender sender = senderFactory.createSender(topic, schema, customizers); + assertThat(sender).extracting("senderSpec", InstanceOfAssertFactories.type(ReactiveMessageSenderSpec.class)) + .extracting(ReactiveMessageSenderSpec::getTopicName).isEqualTo(expectedTopic); + return sender; + } + } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java index de66f0f7..2b9e9c51 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java @@ -19,11 +19,6 @@ package org.springframework.pulsar.core.reactive; import static org.assertj.core.api.Assertions.assertThat; import static org.awaitility.Awaitility.await; import static org.junit.jupiter.params.provider.Arguments.arguments; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.argThat; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; import java.time.Duration; import java.util.ArrayList; @@ -32,14 +27,11 @@ import java.util.UUID; import java.util.concurrent.TimeUnit; import java.util.stream.Stream; -import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; -import org.apache.pulsar.client.api.MessageRouter; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; -import org.apache.pulsar.client.api.TopicMetadata; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageSenderSpec; import org.assertj.core.api.InstanceOfAssertFactories; import org.junit.jupiter.api.Test; @@ -94,11 +86,6 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { String topic = testName; String subscription = topic + "-sub"; String msgPayload = topic + "-msg"; - MessageRouter router = null; - if (testArgs.useCustomRouter) { - router = mock(MessageRouter.class); - when(router.choosePartition(any(Message.class), any(TopicMetadata.class))).thenReturn(0); - } MessageSpecBuilderCustomizer messageCustomizer = null; if (testArgs.useMessageCustomizer) { messageCustomizer = (mb) -> mb.key("foo-key"); @@ -107,13 +94,6 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { if (testArgs.useSenderCustomizer) { senderCustomizer = (sb) -> sb.producerName("foo-producer"); } - - if (router != null) { - try (PulsarAdmin admin = PulsarAdmin.builder() - .serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()).build()) { - admin.topics().createPartitionedTopic("persistent://public/default/" + topic, 1); - } - } try (PulsarClient client = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build()) { try (Consumer consumer = client.newConsumer(Schema.STRING).topic(topic) @@ -142,9 +122,6 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { if (messageCustomizer != null) { messageBuilder = messageBuilder.withMessageCustomizer(messageCustomizer); } - if (router != null) { - messageBuilder = messageBuilder.withCustomRouter(router); - } if (senderCustomizer != null) { messageBuilder = messageBuilder.withSenderCustomizer(senderCustomizer); } @@ -159,10 +136,6 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { if (messageCustomizer != null) { assertThat(msg.getKey()).isEqualTo("foo-key"); } - if (router != null) { - verify(router).choosePartition(argThat((Message m) -> m.getTopicName().equals(topic)), - any(TopicMetadata.class)); - } if (senderCustomizer != null) { assertThat(msg.getProducerName()).isEqualTo("foo-producer"); } @@ -180,36 +153,29 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { SendTestArgs.useSpecificTopic(false).useSimpleApi()), arguments("sendReactiveMessageToDefaultTopicWithSimpleApiAndTemplateSchema", SendTestArgs.useSpecificTopic(false).useSimpleApi().useTemplateSchema()), - arguments("sendReactiveMessageToDefaultTopicWithRouter", - SendTestArgs.useSpecificTopic(false).useCustomRouter()), arguments("sendReactiveMessageToDefaultTopicWithMessageCustomizer", SendTestArgs.useSpecificTopic(false).useMessageCustomizer()), arguments("sendReactiveMessageToDefaultTopicWithProducerCustomizer", SendTestArgs.useSpecificTopic(false).useSenderCustomizer()), arguments("sendReactiveMessageToDefaultTopicWithAllOptions", - SendTestArgs.useSpecificTopic(false).useCustomRouter().useMessageCustomizer() - .useSenderCustomizer()), + SendTestArgs.useSpecificTopic(false).useMessageCustomizer().useSenderCustomizer()), arguments("sendReactiveMessageToSpecificTopic", SendTestArgs.useSpecificTopic(true)), arguments("sendReactiveMessageToSpecificTopicWithSimpleApi", SendTestArgs.useSpecificTopic(true).useSimpleApi()), arguments("sendReactiveMessageToSpecificTopicWithSimpleApiAndTemplateSchema", SendTestArgs.useSpecificTopic(true).useSimpleApi().useTemplateSchema()), - arguments("sendReactiveMessageToSpecificTopicWithRouter", - SendTestArgs.useSpecificTopic(true).useCustomRouter()), arguments("sendReactiveMessageToSpecificTopicWithMessageCustomizer", SendTestArgs.useSpecificTopic(true).useMessageCustomizer()), arguments("sendReactiveMessageToSpecificTopicWithProducerCustomizer", SendTestArgs.useSpecificTopic(true).useSenderCustomizer()), - arguments("sendReactiveMessageToSpecificTopicWithAllOptions", SendTestArgs.useSpecificTopic(true) - .useCustomRouter().useMessageCustomizer().useSenderCustomizer())); + arguments("sendReactiveMessageToSpecificTopicWithAllOptions", + SendTestArgs.useSpecificTopic(true).useMessageCustomizer().useSenderCustomizer())); } static final class SendTestArgs { private final boolean useSpecificTopic; - private boolean useCustomRouter; - private boolean useMessageCustomizer; private boolean useSenderCustomizer; @@ -226,11 +192,6 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { return new SendTestArgs(useSpecificTopic); } - SendTestArgs useCustomRouter() { - this.useCustomRouter = true; - return this; - } - SendTestArgs useMessageCustomizer() { this.useMessageCustomizer = true; return this;