Reactive template can send multiple messages (#170)

Add reactive template operations to send multiple messages with a r.s.Publisher
This commit is contained in:
Christophe Bornet
2022-10-24 17:45:59 +02:00
committed by Chris Bono
parent 0447b19aa5
commit 911f0de710
3 changed files with 116 additions and 33 deletions

View File

@@ -18,7 +18,9 @@ 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;
import reactor.core.publisher.Mono;
/**
@@ -45,6 +47,24 @@ public interface ReactivePulsarSenderOperations<T> {
*/
Mono<MessageId> send(String topic, T message);
/**
* Sends multiple messages to the default topic in a reactive manner.
* @param messages the messages to send
* @return the ids assigned by the broker to the published messages in the same order
* as they were sent
*/
Flux<MessageId> send(Publisher<T> messages);
/**
* Sends multiple messages to the specified topic in a reactive manner.
* @param topic the topic to send the message to or {@code null} to send to the
* default topic
* @param messages the messages to send
* @return the ids assigned by the broker to the published messages in the same order
* as they were sent
*/
Flux<MessageId> send(String topic, Publisher<T> messages);
/**
* Create a {@link SendMessageBuilder builder} for configuring and sending a message
* reactively.

View File

@@ -24,10 +24,12 @@ import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.reactive.client.api.MessageSpec;
import org.apache.pulsar.reactive.client.api.MessageSpecBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageSender;
import org.reactivestreams.Publisher;
import org.springframework.core.log.LogAccessor;
import org.springframework.pulsar.core.SchemaUtils;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
/**
@@ -55,7 +57,7 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
@Override
public Mono<MessageId> send(T message) {
return doSend(null, message, null, null, null);
return send(null, message);
}
@Override
@@ -63,6 +65,16 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
return doSend(topic, message, null, null, null);
}
@Override
public Flux<MessageId> send(Publisher<T> messages) {
return send(null, messages);
}
@Override
public Flux<MessageId> send(String topic, Publisher<T> messages) {
return doSendMany(topic, messages, null, null, null);
}
@Override
public SendMessageBuilderImpl<T> newMessage(T message) {
return new SendMessageBuilderImpl<>(this, message);
@@ -79,22 +91,49 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
private Mono<MessageId> doSend(String topic, T message,
MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer, MessageRouter messageRouter,
ReactiveMessageSenderBuilderCustomizer<T> customizer) {
return doSendMany(topic, Mono.just(message), messageSpecBuilderCustomizer, messageRouter, customizer).single();
}
private Flux<MessageId> doSendMany(String topic, Publisher<T> messages,
MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer, MessageRouter messageRouter,
ReactiveMessageSenderBuilderCustomizer<T> customizer) {
final String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory);
this.logger.trace(() -> String.format("Sending reative msg to '%s' topic", topicName));
this.logger.trace(() -> String.format("Sending reactive messages to '%s' topic", topicName));
final ReactiveMessageSender<T> sender = createMessageSender(topic, message, messageRouter, customizer);
if (this.schema != null) {
/*
* If the template has a schema, we can create the message sender right away
* and use ReactiveMessageSender::sendMessages to send them as a stream.
* Otherwise we need to wait to get a message to create it and we can't share
* 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);
return Flux.from(messages).map(message -> getMessageSpec(messageSpecBuilderCustomizer, message))
.as(sender::sendMessages)
.doOnError(ex -> this.logger.error(ex,
() -> String.format("Failed to send messages to '%s' topic", topicName)))
.doOnNext(
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);
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(
msgId -> this.logger.trace(() -> String.format("Sent message to '%s' topic", topicName)));
});
}
private static <T> MessageSpec<T> getMessageSpec(MessageSpecBuilderCustomizer<T> messageSpecBuilderCustomizer,
T message) {
MessageSpecBuilder<T> messageSpecBuilder = MessageSpec.builder(message);
if (messageSpecBuilderCustomizer != null) {
messageSpecBuilderCustomizer.customize(messageSpecBuilder);
}
MessageSpec<T> messageSpec = messageSpecBuilder.build();
return sender.sendMessage(Mono.just(messageSpec))
.doOnError(
ex -> this.logger.error(ex, () -> String.format("Failed to send msg to '%s' topic", topicName)))
.doOnSuccess(msgId -> this.logger.trace(() -> String.format("Sent msg to '%s' topic", topicName)));
return messageSpecBuilder.build();
}
private ReactiveMessageSender<T> createMessageSender(String topic, T message, MessageRouter messageRouter,

View File

@@ -26,6 +26,8 @@ import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;
@@ -47,6 +49,7 @@ import org.junit.jupiter.params.provider.MethodSource;
import org.springframework.pulsar.core.PulsarTestContainerSupport;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
/**
@@ -57,7 +60,7 @@ import reactor.core.publisher.Mono;
class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
@Test
void sendMessageWithSpecificSchemaTest() throws Exception {
void sendMessagesWithSpecificSchemaTest() throws Exception {
String topic = "smt-specific-schema-topic-reactive";
try (PulsarClient client = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build()) {
@@ -69,10 +72,17 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
senderSpec, null);
ReactivePulsarSenderTemplate<Foo> pulsarTemplate = new ReactivePulsarSenderTemplate<>(producerFactory);
pulsarTemplate.setSchema(Schema.JSON(Foo.class));
Foo foo = new Foo("Foo-" + UUID.randomUUID(), "Bar-" + UUID.randomUUID());
pulsarTemplate.send(foo).subscribe();
assertThat(consumer.receiveAsync()).succeedsWithin(Duration.ofSeconds(3)).extracting(Message::getValue)
.isEqualTo(foo);
List<Foo> foos = new ArrayList<>();
for (int i = 0; i < 10; i++) {
foos.add(new Foo("Foo-" + UUID.randomUUID(), "Bar-" + UUID.randomUUID()));
}
pulsarTemplate.send(Flux.fromIterable(foos)).subscribe();
for (int i = 0; i < 10; i++) {
assertThat(consumer.receiveAsync()).succeedsWithin(Duration.ofSeconds(3))
.extracting(Message::getValue).isEqualTo(foos.get(i));
}
}
}
}
@@ -116,6 +126,9 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
senderSpec, null);
ReactivePulsarSenderTemplate<String> pulsarTemplate = new ReactivePulsarSenderTemplate<>(senderFactory);
Mono<MessageId> sendResponse;
if (testArgs.useTemplateSchema) {
pulsarTemplate.setSchema(Schema.STRING);
}
if (testArgs.useSimpleApi) {
sendResponse = testArgs.useSpecificTopic ? pulsarTemplate.send(topic, msgPayload)
: pulsarTemplate.send(msgPayload);
@@ -164,32 +177,36 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
private static Stream<Arguments> sendMessageTestProvider() {
return Stream.of(arguments("sendReactiveMessageToDefaultTopic", SendTestArgs.useSpecificTopic(false)),
arguments("sendReactiveMessageToDefaultTopicWithSimpleApi",
SendTestArgs.useSpecificTopic(false).useSimpleApi(true)),
SendTestArgs.useSpecificTopic(false).useSimpleApi()),
arguments("sendReactiveMessageToDefaultTopicWithSimpleApiAndTemplateSchema",
SendTestArgs.useSpecificTopic(false).useSimpleApi().useTemplateSchema()),
arguments("sendReactiveMessageToDefaultTopicWithRouter",
SendTestArgs.useSpecificTopic(false).useCustomRouter(true)),
SendTestArgs.useSpecificTopic(false).useCustomRouter()),
arguments("sendReactiveMessageToDefaultTopicWithMessageCustomizer",
SendTestArgs.useSpecificTopic(false).useMessageCustomizer(true)),
SendTestArgs.useSpecificTopic(false).useMessageCustomizer()),
arguments("sendReactiveMessageToDefaultTopicWithProducerCustomizer",
SendTestArgs.useSpecificTopic(false).useSenderCustomizer(true)),
SendTestArgs.useSpecificTopic(false).useSenderCustomizer()),
arguments("sendReactiveMessageToDefaultTopicWithAllOptions",
SendTestArgs.useSpecificTopic(false).useCustomRouter(true).useMessageCustomizer(true)
.useSenderCustomizer(true)),
SendTestArgs.useSpecificTopic(false).useCustomRouter().useMessageCustomizer()
.useSenderCustomizer()),
arguments("sendReactiveMessageToSpecificTopic", SendTestArgs.useSpecificTopic(true)),
arguments("sendReactiveMessageToSpecificTopicWithSimpleApi",
SendTestArgs.useSpecificTopic(true).useSimpleApi(true)),
SendTestArgs.useSpecificTopic(true).useSimpleApi()),
arguments("sendReactiveMessageToSpecificTopicWithSimpleApiAndTemplateSchema",
SendTestArgs.useSpecificTopic(true).useSimpleApi().useTemplateSchema()),
arguments("sendReactiveMessageToSpecificTopicWithRouter",
SendTestArgs.useSpecificTopic(true).useCustomRouter(true)),
SendTestArgs.useSpecificTopic(true).useCustomRouter()),
arguments("sendReactiveMessageToSpecificTopicWithMessageCustomizer",
SendTestArgs.useSpecificTopic(true).useMessageCustomizer(true)),
SendTestArgs.useSpecificTopic(true).useMessageCustomizer()),
arguments("sendReactiveMessageToSpecificTopicWithProducerCustomizer",
SendTestArgs.useSpecificTopic(true).useSenderCustomizer(true)),
SendTestArgs.useSpecificTopic(true).useSenderCustomizer()),
arguments("sendReactiveMessageToSpecificTopicWithAllOptions", SendTestArgs.useSpecificTopic(true)
.useCustomRouter(true).useMessageCustomizer(true).useSenderCustomizer(true)));
.useCustomRouter().useMessageCustomizer().useSenderCustomizer()));
}
static final class SendTestArgs {
private boolean useSpecificTopic;
private final boolean useSpecificTopic;
private boolean useCustomRouter;
@@ -199,6 +216,8 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
private boolean useSimpleApi;
private boolean useTemplateSchema;
private SendTestArgs(boolean useSpecificTopic) {
this.useSpecificTopic = useSpecificTopic;
}
@@ -207,23 +226,28 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
return new SendTestArgs(useSpecificTopic);
}
SendTestArgs useCustomRouter(boolean useCustomRouter) {
this.useCustomRouter = useCustomRouter;
SendTestArgs useCustomRouter() {
this.useCustomRouter = true;
return this;
}
SendTestArgs useMessageCustomizer(boolean useMessageCustomizer) {
this.useMessageCustomizer = useMessageCustomizer;
SendTestArgs useMessageCustomizer() {
this.useMessageCustomizer = true;
return this;
}
SendTestArgs useSenderCustomizer(boolean useSenderCustomizer) {
this.useSenderCustomizer = useSenderCustomizer;
SendTestArgs useSenderCustomizer() {
this.useSenderCustomizer = true;
return this;
}
SendTestArgs useSimpleApi(boolean useSimpleApi) {
this.useSimpleApi = useSimpleApi;
SendTestArgs useSimpleApi() {
this.useSimpleApi = true;
return this;
}
SendTestArgs useTemplateSchema() {
this.useTemplateSchema = true;
return this;
}