From 1055637ff24937d94e628b84a2e22581c5909caf Mon Sep 17 00:00:00 2001 From: Christophe Bornet Date: Tue, 22 Nov 2022 03:05:56 +0100 Subject: [PATCH] Update with new MessageResult::acknowledge(Message) function (#233) --- ...ivePulsarMessageListenerContainerTests.java | 18 ++++++------------ .../listener/ReactivePulsarListenerTests.java | 2 +- .../ReactiveSpringPulsarBootApp.java | 2 +- .../ReactivePulsarListenerTests.java | 2 +- 4 files changed, 9 insertions(+), 15 deletions(-) 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 be480902..4171013c 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 @@ -102,10 +102,8 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo pulsarConsumerFactory.createConsumer(Schema.STRING).consumeNothing().block(Duration.ofSeconds(10)); CountDownLatch latch = new CountDownLatch(5); ReactivePulsarContainerProperties pulsarContainerProperties = new ReactivePulsarContainerProperties<>(); - pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler) (msg) -> msg.map(m -> { - latch.countDown(); - return MessageResult.acknowledge(m.getMessageId()); - })); + pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler) (msg) -> msg + .doOnNext((m) -> latch.countDown()).map(MessageResult::acknowledge)); pulsarContainerProperties.setSchema(Schema.STRING); DefaultReactivePulsarMessageListenerContainer container = new DefaultReactivePulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); @@ -292,13 +290,9 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo CountDownLatch latch = new CountDownLatch(6); ReactivePulsarContainerProperties pulsarContainerProperties = new ReactivePulsarContainerProperties<>(); - pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler) (msg) -> msg.map(m -> { - latch.countDown(); - if (m.getValue().endsWith("4")) { - return MessageResult.negativeAcknowledge(m.getMessageId()); - } - return MessageResult.acknowledge(m.getMessageId()); - })); + pulsarContainerProperties.setMessageHandler((ReactivePulsarStreamingHandler) (msg) -> msg + .doOnNext((m) -> latch.countDown()).map((m) -> m.getValue().endsWith("4") + ? MessageResult.negativeAcknowledge(m) : MessageResult.acknowledge(m))); pulsarContainerProperties.setSchema(Schema.STRING); pulsarContainerProperties.setSubscriptionType(SubscriptionType.Shared); DefaultReactivePulsarMessageListenerContainer container = new DefaultReactivePulsarMessageListenerContainer<>( @@ -320,7 +314,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo if (message.getValue().endsWith("4")) { dlqLatch.countDown(); } - return Mono.just(MessageResult.acknowledge(message.getMessageId())); + return Mono.just(MessageResult.acknowledge(message)); }).block(); assertThat(dlqLatch.await(10, TimeUnit.SECONDS)).isTrue(); diff --git a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/ReactivePulsarListenerTests.java b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/ReactivePulsarListenerTests.java index 944da9c2..1d1ca48c 100644 --- a/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/ReactivePulsarListenerTests.java +++ b/spring-pulsar-reactive/src/test/java/org/springframework/pulsar/reactive/listener/ReactivePulsarListenerTests.java @@ -259,7 +259,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport { @ReactivePulsarListener(topics = "streaming-1", stream = true, consumerCustomizer = "consumerCustomizer") Flux> listen1(Flux> messages) { - return messages.doOnNext(m -> latch1.countDown()).map(m -> MessageResult.acknowledge(m.getMessageId())); + return messages.doOnNext(m -> latch1.countDown()).map(MessageResult::acknowledge); } @ReactivePulsarListener(topics = "streaming-2", stream = true, consumerCustomizer = "consumerCustomizer") 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 750131ef..6c53a51e 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 @@ -98,7 +98,7 @@ public class ReactiveSpringPulsarBootApp { public Flux> listenStreaming(Flux> messages) { return messages .doOnNext((msg) -> this.logger.info("Streaming reactive listener received: {}", msg.getValue())) - .map(m -> MessageResult.acknowledge(m.getMessageId())); + .map(MessageResult::acknowledge); } } diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/ReactivePulsarListenerTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/ReactivePulsarListenerTests.java index 9f304f28..4969e204 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/ReactivePulsarListenerTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/ReactivePulsarListenerTests.java @@ -101,7 +101,7 @@ class ReactivePulsarListenerTests implements PulsarTestContainerSupport { @ReactivePulsarListener(subscriptionName = "rplt-subscription2", topics = "rplt-topic2", stream = true, consumerCustomizer = "consumerCustomizer") public Flux> listen(Flux> messages) { - return messages.doOnNext(t -> LATCH2.countDown()).map(m -> MessageResult.acknowledge(m.getMessageId())); + return messages.doOnNext(t -> LATCH2.countDown()).map(MessageResult::acknowledge); } }