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 1d0d5352..04419b07 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 @@ -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 { */ Mono 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 send(Publisher 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 send(String topic, Publisher messages); + /** * Create a {@link SendMessageBuilder builder} for configuring and sending a message * reactively. 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 c76646f5..7f901e28 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 @@ -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 implements ReactivePulsarSenderOper @Override public Mono send(T message) { - return doSend(null, message, null, null, null); + return send(null, message); } @Override @@ -63,6 +65,16 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper return doSend(topic, message, null, null, null); } + @Override + public Flux send(Publisher messages) { + return send(null, messages); + } + + @Override + public Flux send(String topic, Publisher messages) { + return doSendMany(topic, messages, null, null, null); + } + @Override public SendMessageBuilderImpl newMessage(T message) { return new SendMessageBuilderImpl<>(this, message); @@ -79,22 +91,49 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper private Mono doSend(String topic, T message, MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, MessageRouter messageRouter, ReactiveMessageSenderBuilderCustomizer customizer) { + return doSendMany(topic, Mono.just(message), messageSpecBuilderCustomizer, messageRouter, customizer).single(); + } + + private Flux doSendMany(String topic, Publisher messages, + MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, MessageRouter messageRouter, + ReactiveMessageSenderBuilderCustomizer 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 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 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 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 MessageSpec getMessageSpec(MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, + T message) { MessageSpecBuilder messageSpecBuilder = MessageSpec.builder(message); if (messageSpecBuilderCustomizer != null) { messageSpecBuilderCustomizer.customize(messageSpecBuilder); } - MessageSpec 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 createMessageSender(String topic, T message, MessageRouter messageRouter, 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 3b477083..de66f0f7 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 @@ -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 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 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 pulsarTemplate = new ReactivePulsarSenderTemplate<>(senderFactory); Mono 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 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; }