diff --git a/spring-pulsar-dependencies/build.gradle b/spring-pulsar-dependencies/build.gradle index 3b58e0e4..f6b19afa 100644 --- a/spring-pulsar-dependencies/build.gradle +++ b/spring-pulsar-dependencies/build.gradle @@ -11,7 +11,7 @@ ext { protobufJavaVersion = '3.21.5' testcontainersVersion = '1.17.6' pulsarVersion = '2.10.3' - pulsarClientReactiveVersion = '0.1.0' + pulsarClientReactiveVersion = '0.2.0' springBootVersion = '3.0.2' springCloudStreamVersion = '4.0.1' } diff --git a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarOperations.java b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarOperations.java index ac85e0df..1bd5aaf9 100644 --- a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarOperations.java +++ b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarOperations.java @@ -18,6 +18,8 @@ package org.springframework.pulsar.reactive.core; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.reactive.client.api.MessageSendResult; +import org.apache.pulsar.reactive.client.api.MessageSpec; import org.reactivestreams.Publisher; import org.springframework.lang.Nullable; @@ -76,7 +78,7 @@ public interface ReactivePulsarOperations { * @return the ids assigned by the broker to the published messages in the same order * as they were sent */ - Flux send(Publisher messages); + Flux> send(Publisher> messages); /** * Sends multiple messages to the default topic in a reactive manner. @@ -86,7 +88,7 @@ public interface ReactivePulsarOperations { * @return the ids assigned by the broker to the published messages in the same order * as they were sent */ - Flux send(Publisher messages, @Nullable Schema schema); + Flux> send(Publisher> messages, @Nullable Schema schema); /** * Sends multiple messages to the specified topic in a reactive manner. @@ -96,7 +98,7 @@ public interface ReactivePulsarOperations { * @return the ids assigned by the broker to the published messages in the same order * as they were sent */ - Flux send(@Nullable String topic, Publisher messages); + Flux> send(@Nullable String topic, Publisher> messages); /** * Sends multiple messages to the specified topic in a reactive manner. @@ -108,7 +110,8 @@ public interface ReactivePulsarOperations { * @return the ids assigned by the broker to the published messages in the same order * as they were sent */ - Flux send(@Nullable String topic, Publisher messages, @Nullable Schema schema); + Flux> send(@Nullable String topic, Publisher> messages, + @Nullable Schema schema); /** * Create a {@link SendOneMessageBuilder builder} for configuring and sending a @@ -124,7 +127,7 @@ public interface ReactivePulsarOperations { * @param messages the messages to send * @return the builder to configure and send the message */ - SendManyMessageBuilder newMessages(Publisher messages); + SendManyMessageBuilder newMessages(Publisher> messages); /** * Builder that can be used to configure and send a message. Provides more options @@ -184,7 +187,7 @@ public interface ReactivePulsarOperations { * @return the ids assigned by the broker to the published messages in the same * order as they were sent */ - Flux send(); + Flux> send(); } diff --git a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplate.java b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplate.java index ebb2b1ce..0c1b0ee9 100644 --- a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplate.java +++ b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplate.java @@ -20,6 +20,7 @@ import java.util.Optional; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.reactive.client.api.MessageSendResult; import org.apache.pulsar.reactive.client.api.MessageSpec; import org.apache.pulsar.reactive.client.api.MessageSpecBuilder; import org.apache.pulsar.reactive.client.api.ReactiveMessageSender; @@ -62,9 +63,10 @@ public class ReactivePulsarTemplate implements ReactivePulsarOperations { } /** - * Construct a template instance with a custom schema resolver. + * Construct a template instance with a custom schema resolver and a custom topic + * resolver. * @param reactiveMessageSenderFactory the factory used to create the backing Pulsar - * @param schemaResolver the schema resolver to use reactive senders + * @param schemaResolver the schema resolver to use * @param topicResolver the topic resolver to use */ public ReactivePulsarTemplate(ReactivePulsarSenderFactory reactiveMessageSenderFactory, @@ -95,22 +97,23 @@ public class ReactivePulsarTemplate implements ReactivePulsarOperations { } @Override - public Flux send(Publisher messages) { + public Flux> send(Publisher> messages) { return send(null, messages); } @Override - public Flux send(Publisher messages, @Nullable Schema schema) { + public Flux> send(Publisher> messages, @Nullable Schema schema) { return doSendMany(null, Flux.from(messages), schema, null); } @Override - public Flux send(@Nullable String topic, Publisher messages) { + public Flux> send(@Nullable String topic, Publisher> messages) { return doSendMany(topic, Flux.from(messages), null, null); } @Override - public Flux send(@Nullable String topic, Publisher messages, @Nullable Schema schema) { + public Flux> send(@Nullable String topic, Publisher> messages, + @Nullable Schema schema) { return doSendMany(topic, Flux.from(messages), schema, null); } @@ -120,7 +123,7 @@ public class ReactivePulsarTemplate implements ReactivePulsarOperations { } @Override - public SendManyMessageBuilder newMessages(Publisher messages) { + public SendManyMessageBuilder newMessages(Publisher> messages) { return new SendManyMessageBuilderImpl<>(this, messages); } @@ -137,14 +140,15 @@ public class ReactivePulsarTemplate implements ReactivePulsarOperations { // @formatter:on } - private Flux doSendMany(@Nullable String topic, Flux messages, @Nullable Schema schema, - @Nullable ReactiveMessageSenderBuilderCustomizer customizer) { + private Flux> doSendMany(@Nullable String topic, Flux> messages, + @Nullable Schema schema, @Nullable ReactiveMessageSenderBuilderCustomizer customizer) { return messages.switchOnFirst((firstSignal, messageFlux) -> { - T firstMessage = firstSignal.get(); + MessageSpec firstMessage = firstSignal.get(); if (firstMessage != null && firstSignal.isOnNext()) { - String topicName = resolveTopic(topic, firstMessage.getClass()); - ReactiveMessageSender sender = createMessageSender(topicName, firstMessage, schema, customizer); - return messageFlux.map(MessageSpec::of).as(sender::sendMany).doOnError( + String topicName = resolveTopic(topic, firstMessage.getValue().getClass()); + ReactiveMessageSender sender = createMessageSender(topicName, firstMessage.getValue(), schema, + customizer); + return messageFlux.as(sender::sendMany).doOnError( ex -> this.logger.error(ex, () -> "Failed to send messages to '%s' topic".formatted(topicName))) .doOnNext(msgId -> this.logger.trace(() -> "Sent messages to '%s' topic".formatted(topicName))); } @@ -251,15 +255,15 @@ public class ReactivePulsarTemplate implements ReactivePulsarOperations { private static final class SendManyMessageBuilderImpl extends SendMessageBuilderImpl, T> implements SendManyMessageBuilder { - private final Publisher messages; + private final Publisher> messages; - SendManyMessageBuilderImpl(ReactivePulsarTemplate template, Publisher messages) { + SendManyMessageBuilderImpl(ReactivePulsarTemplate template, Publisher> messages) { super(template); this.messages = messages; } @Override - public Flux send() { + public Flux> send() { return this.template.doSendMany(this.topic, Flux.from(this.messages), this.schema, this.senderCustomizer); } diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java index 564ef3b5..3cc15a50 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/core/ReactivePulsarTemplateTests.java @@ -35,6 +35,7 @@ import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.reactive.client.api.MessageSpec; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageSenderSpec; import org.assertj.core.api.InstanceOfAssertFactories; import org.junit.jupiter.params.ParameterizedTest; @@ -78,11 +79,12 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { } if (useSimpleApi) { - pulsarTemplate.send(Flux.fromIterable(foos), Schema.JSON(Foo.class)).subscribe(); + pulsarTemplate.send(Flux.fromIterable(foos).map(MessageSpec::of), Schema.JSON(Foo.class)) + .subscribe(); } else { - pulsarTemplate.newMessages(Flux.fromIterable(foos)).withSchema(Schema.JSON(Foo.class)).send() - .subscribe(); + pulsarTemplate.newMessages(Flux.fromIterable(foos).map(MessageSpec::of)) + .withSchema(Schema.JSON(Foo.class)).send().subscribe(); } for (int i = 0; i < 10; i++) { @@ -118,10 +120,10 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { } if (useSimpleApi) { - pulsarTemplate.send(Flux.fromIterable(foos)).subscribe(); + pulsarTemplate.send(Flux.fromIterable(foos).map(MessageSpec::of)).subscribe(); } else { - pulsarTemplate.newMessages(Flux.fromIterable(foos)).send().subscribe(); + pulsarTemplate.newMessages(Flux.fromIterable(foos).map(MessageSpec::of)).send().subscribe(); } // TODO figure out if expected to not be ordered when schema not set on @@ -174,16 +176,17 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { return Stream.of( arguments("simpleApiNoSchema", (BiConsumer, ReactivePulsarTemplate>) (data, template) -> template - .send(Flux.fromIterable(data)).subscribe()), + .send(Flux.fromIterable(data).map(MessageSpec::of)).subscribe()), arguments("simpleApiWithSchema", (BiConsumer, ReactivePulsarTemplate>) (data, template) -> template - .send(Flux.fromIterable(data), Schema.STRING).subscribe()), + .send(Flux.fromIterable(data).map(MessageSpec::of), Schema.STRING).subscribe()), arguments("fluentApiNoSchema", (BiConsumer, ReactivePulsarTemplate>) (data, template) -> template - .newMessages(Flux.fromIterable(data)).send().subscribe()), + .newMessages(Flux.fromIterable(data).map(MessageSpec::of)).send().subscribe()), arguments("fluentApiWithSchema", (BiConsumer, ReactivePulsarTemplate>) (data, template) -> template - .newMessages(Flux.fromIterable(data)).withSchema(Schema.STRING).send().subscribe())); + .newMessages(Flux.fromIterable(data).map(MessageSpec::of)).withSchema(Schema.STRING) + .send().subscribe())); } @ParameterizedTest(name = "{0}") diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java index 593cd8d9..e20ded2d 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/DefaultReactivePulsarMessageListenerContainerTests.java @@ -31,6 +31,7 @@ import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory; import org.apache.pulsar.reactive.client.adapter.DefaultMessageGroupingFunction; import org.apache.pulsar.reactive.client.api.MessageResult; +import org.apache.pulsar.reactive.client.api.MessageSpec; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageConsumerSpec; import org.apache.pulsar.reactive.client.api.MutableReactiveMessageSenderSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumer; @@ -114,7 +115,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo DefaultReactivePulsarSenderFactory pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>( reactivePulsarClient, prodConfig, null, new DefaultTopicResolver()); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(pulsarProducerFactory); - Flux.range(0, 5).map(i -> "hello john doe" + i).as(pulsarTemplate::send).subscribe(); + Flux.range(0, 5).map(i -> MessageSpec.of("hello john doe" + i)).as(pulsarTemplate::send).subscribe(); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); container.stop(); pulsarClient.close(); @@ -308,7 +309,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo DefaultReactivePulsarSenderFactory pulsarProducerFactory = new DefaultReactivePulsarSenderFactory<>( reactivePulsarClient, prodConfig, null, new DefaultTopicResolver()); ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(pulsarProducerFactory); - Flux.range(0, 5).map(i -> "hello john doe" + i).as(pulsarTemplate::send).subscribe(); + Flux.range(0, 5).map(i -> MessageSpec.of("hello john doe" + i)).as(pulsarTemplate::send).subscribe(); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); CountDownLatch dlqLatch = new CountDownLatch(1); diff --git a/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java b/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java index 96281f1b..961c6abc 100644 --- a/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java +++ b/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java @@ -21,6 +21,7 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionInitialPosition; import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.reactive.client.api.MessageResult; +import org.apache.pulsar.reactive.client.api.MessageSpec; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -61,7 +62,7 @@ public class ReactiveSpringPulsarBootApp { @Override public void onApplicationEvent(ApplicationReadyEvent event) { - Flux.range(0, 10).map((i) -> "sample-message-" + i) + Flux.range(0, 10).map((i) -> MessageSpec.of("sample-message-" + i)) .as(messages -> this.reactivePulsarTemplate.send("sample-reactive-topic1", messages)).subscribe(); } @@ -89,7 +90,7 @@ public class ReactiveSpringPulsarBootApp { @Override public void onApplicationEvent(ApplicationReadyEvent event) { Schema schema = Schema.JSON(Foo.class); - Flux.range(0, 10).map((i) -> new Foo("Foo-" + i, "Bar-" + i)) + Flux.range(0, 10).map((i) -> MessageSpec.of(new Foo("Foo-" + i, "Bar-" + i))) .as(messages -> this.reactivePulsarTemplate.send("sample-reactive-topic2", messages, schema)) .subscribe(); } @@ -125,7 +126,7 @@ public class ReactiveSpringPulsarBootApp { @Bean ApplicationRunner sendSimple(ReactivePulsarTemplate reactivePulsarTemplate) { - return args -> Flux.range(0, 10).map((i) -> "msg-from-sendSimple-" + i) + return args -> Flux.range(0, 10).map((i) -> MessageSpec.of("msg-from-sendSimple-" + i)) .as(messages -> reactivePulsarTemplate.send("sample-reactive-topic3", messages)).subscribe(); }