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)
This commit is contained in:
Christophe Bornet
2022-10-27 21:44:20 +02:00
committed by Chris Bono
parent d50ab042b9
commit 1758fe87df
6 changed files with 49 additions and 147 deletions

View File

@@ -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<T> 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<T> createReactiveMessageSender(String topic, Schema<T> schema) {
return doCreateReactiveMessageSender(topic, schema, null, null);
public ReactiveMessageSender<T> createSender(String topic, Schema<T> schema) {
return doCreateReactiveMessageSender(topic, schema, null);
}
@Override
public ReactiveMessageSender<T> createReactiveMessageSender(String topic, Schema<T> schema,
MessageRouter messageRouter) {
return doCreateReactiveMessageSender(topic, schema, messageRouter, null);
}
@Override
public ReactiveMessageSender<T> createReactiveMessageSender(String topic, Schema<T> schema,
MessageRouter messageRouter, List<ReactiveMessageSenderBuilderCustomizer<T>> customizers) {
return doCreateReactiveMessageSender(topic, schema, messageRouter, customizers);
public ReactiveMessageSender<T> createSender(String topic, Schema<T> schema,
List<ReactiveMessageSenderBuilderCustomizer<T>> customizers) {
return doCreateReactiveMessageSender(topic, schema, customizers);
}
private ReactiveMessageSender<T> doCreateReactiveMessageSender(String topic, Schema<T> schema,
MessageRouter messageRouter, List<ReactiveMessageSenderBuilderCustomizer<T>> customizers) {
List<ReactiveMessageSenderBuilderCustomizer<T>> customizers) {
final String resolvedTopic = ReactiveMessageSenderUtils.resolveTopicName(topic, this);
this.logger.trace(() -> String.format("Creating reactive message sender for '%s' topic", resolvedTopic));
final ReactiveMessageSenderBuilder<T> 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));
}

View File

@@ -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<T> {
* @param schema the schema of the messages to be sent
* @return the reactive message sender
*/
ReactiveMessageSender<T> createReactiveMessageSender(String topic, Schema<T> schema);
ReactiveMessageSender<T> createSender(String topic, Schema<T> 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<T> createReactiveMessageSender(String topic, Schema<T> 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<T> createReactiveMessageSender(String topic, Schema<T> schema, MessageRouter messageRouter,
ReactiveMessageSender<T> createSender(String topic, Schema<T> schema,
List<ReactiveMessageSenderBuilderCustomizer<T>> customizers);
/**

View File

@@ -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<T> {
*/
SendMessageBuilder<T> withMessageCustomizer(MessageSpecBuilderCustomizer<T> 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<T> withCustomRouter(MessageRouter messageRouter);
/**
* Specifies the customizer to use to further configure the reactive sender
* builder.

View File

@@ -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<T> implements ReactivePulsarSenderOper
@Override
public Mono<MessageId> 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<T> implements ReactivePulsarSenderOper
@Override
public Flux<MessageId> send(String topic, Publisher<T> messages) {
return doSendMany(topic, messages, null, null, null);
return doSendMany(topic, messages, null, null);
}
@Override
@@ -89,13 +88,13 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
}
private Mono<MessageId> doSend(String topic, T message,
MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer, MessageRouter messageRouter,
MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer,
ReactiveMessageSenderBuilderCustomizer<T> customizer) {
return doSendMany(topic, Mono.just(message), messageSpecBuilderCustomizer, messageRouter, customizer).single();
return doSendMany(topic, Mono.just(message), messageSpecBuilderCustomizer, customizer).single();
}
private Flux<MessageId> doSendMany(String topic, Publisher<T> messages,
MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer, MessageRouter messageRouter,
MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer,
ReactiveMessageSenderBuilderCustomizer<T> 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<T> implements ReactivePulsarSenderOper
* it between messages. So we create one each time and use
* ReactiveMessageSender::sendMessage to send messages individually.
*/
ReactiveMessageSender<T> sender = createMessageSender(topic, null, messageRouter, customizer);
ReactiveMessageSender<T> 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<T> implements ReactivePulsarSenderOper
msgId -> this.logger.trace(() -> String.format("Sent messages to '%s' topic", topicName)));
}
return Flux.from(messages).flatMapSequential(message -> {
ReactiveMessageSender<T> sender = createMessageSender(topic, message, messageRouter, customizer);
ReactiveMessageSender<T> 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<T> implements ReactivePulsarSenderOper
return messageSpecBuilder.build();
}
private ReactiveMessageSender<T> createMessageSender(String topic, T message, MessageRouter messageRouter,
private ReactiveMessageSender<T> createMessageSender(String topic, T message,
ReactiveMessageSenderBuilderCustomizer<T> customizer) {
Schema<T> 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<T> implements ReactivePulsarSenderOper
private MessageSpecBuilderCustomizer<T> messageCustomizer;
private MessageRouter messageRouter;
private ReactiveMessageSenderBuilderCustomizer<T> senderCustomizer;
SendMessageBuilderImpl(ReactivePulsarSenderTemplate<T> template, T message) {
@@ -174,12 +171,6 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
return this;
}
@Override
public SendMessageBuilderImpl<T> withCustomRouter(MessageRouter messageRouter) {
this.messageRouter = messageRouter;
return this;
}
@Override
public SendMessageBuilderImpl<T> withSenderCustomizer(
ReactiveMessageSenderBuilderCustomizer<T> senderCustomizer) {
@@ -189,8 +180,7 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
@Override
public Mono<MessageId> 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);
}
}

View File

@@ -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<String> customizer1 = builder -> builder.topic("topic1");
MessageRouter router = mock(MessageRouter.class);
ReactiveMessageSenderBuilderCustomizer<String> customizer2 = builder -> builder.messageRouter(router);
ReactiveMessageSenderCache cache = AdaptedReactivePulsarClientFactory.createCache();
ReactiveMessageSenderBuilderCustomizer<String> customizer2 = builder -> builder.cache(cache);
testCreateSender(null, null, "topic0", null, Arrays.asList(customizer1, customizer2), "topic1", router);
ReactiveMessageSender<String> sender = testCreateSender(null, null, "topic0",
Arrays.asList(customizer1, customizer2), "topic1");
assertThat(sender).extracting("producerCache").isSameAs(cache);
}
@Test
void createSenderWithNoTopic() {
ReactivePulsarSenderFactory<String> 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<ReactiveMessageSenderBuilderCustomizer<String>> customizers,
String expectedTopic, MessageRouter expectedRouter) {
ReactivePulsarSenderFactory<String> senderFactory = new DefaultReactivePulsarSenderFactory<>(
(PulsarClient) null, spec, cache);
ReactiveMessageSender<String> sender = senderFactory.createReactiveMessageSender(topic, schema, router,
customizers);
ObjectAssert<ReactiveMessageSenderSpec> 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<String> sender = testCreateSender(null, cache, "topic1", null, "topic1");
assertThat(sender).extracting("producerCache").isSameAs(cache);
}
private ReactiveMessageSender<String> testCreateSender(ReactiveMessageSenderSpec spec,
ReactiveMessageSenderCache cache, String topic,
List<ReactiveMessageSenderBuilderCustomizer<String>> customizers, String expectedTopic) {
ReactivePulsarSenderFactory<String> senderFactory = new DefaultReactivePulsarSenderFactory<>(
(PulsarClient) null, spec, cache);
ReactiveMessageSender<String> sender = senderFactory.createSender(topic, schema, customizers);
assertThat(sender).extracting("senderSpec", InstanceOfAssertFactories.type(ReactiveMessageSenderSpec.class))
.extracting(ReactiveMessageSenderSpec::getTopicName).isEqualTo(expectedTopic);
return sender;
}
}

View File

@@ -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<String> 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<String> 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<String> 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;