From a0de08e09889de5b6a469b49d4e3bd44df64dd79 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 9 Aug 2022 18:29:07 -0400 Subject: [PATCH] Improve acking when there are no errors When no shared subcription types are used and no errors from processing, we can optimize acking by relying on the acknowledgeCumulative on the Pulsar consumer. --- ...DefaultPulsarMessageListenerContainer.java | 26 +++++++++++++++++-- .../listener/PulsarContainerProperties.java | 2 +- .../PulsarMessageListenerContainerTests.java | 4 +-- 3 files changed, 27 insertions(+), 5 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java index c87673fd..2a1d912e 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java @@ -25,6 +25,8 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.stream.Stream; +import java.util.stream.StreamSupport; import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Consumer; @@ -256,7 +258,15 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } if (this.containerProperties.getAckMode() == PulsarContainerProperties.AckMode.BATCH) { try { - this.consumer.acknowledge(messages); + if (isSharedSubsriptionType()) { + this.consumer.acknowledge(messages); + } + else { + final Stream> stream = StreamSupport.stream(messages.spliterator(), + true); + Message last = stream.reduce((a, b) -> b).orElse(null); + this.consumer.acknowledgeCumulative(last); + } } catch (PulsarClientException pce) { this.consumer.negativeAcknowledge(messages); @@ -303,11 +313,23 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } } + private boolean isSharedSubsriptionType() { + return this.containerProperties.getSubscriptionType() == SubscriptionType.Shared + || this.containerProperties.getSubscriptionType() == SubscriptionType.Key_Shared; + } + private void handleAcks(Messages messages) { if (this.nackableMessages.isEmpty()) { try { if (messages.size() > 0) { - this.consumer.acknowledge(messages); + if (isSharedSubsriptionType()) { + this.consumer.acknowledge(messages); + } + else { + final Stream> stream = StreamSupport.stream(messages.spliterator(), true); + Message last = stream.reduce((a, b) -> b).orElse(null); + this.consumer.acknowledgeCumulative(last); + } } } catch (PulsarClientException pce) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java index c34cda5e..30458d3b 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java @@ -63,7 +63,7 @@ public class PulsarContainerProperties { private String subscriptionName; - private SubscriptionType subscriptionType; + private SubscriptionType subscriptionType = SubscriptionType.Exclusive; private Schema schema; diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarMessageListenerContainerTests.java index b496430e..d227e20f 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarMessageListenerContainerTests.java @@ -128,7 +128,7 @@ class PulsarMessageListenerContainerTests extends AbstractContainerBaseTests { } assertThat(latch.await(30, TimeUnit.SECONDS)).isTrue(); verify(containerConsumer, never()).acknowledge(any(Message.class)); - verify(containerConsumer, atLeastOnce()).acknowledge(any(Messages.class)); + verify(containerConsumer, atLeastOnce()).acknowledgeCumulative(any(Message.class)); container.stop(); pulsarClient.close(); } @@ -270,7 +270,7 @@ class PulsarMessageListenerContainerTests extends AbstractContainerBaseTests { } assertThat(latch.await(30, TimeUnit.SECONDS)).isTrue(); verify(pulsarBatchMessageListener, times(1)).received(any(Consumer.class), any(Messages.class)); - verify(containerConsumer, times(1)).acknowledge(any(Messages.class)); + verify(containerConsumer, times(1)).acknowledgeCumulative(any(Message.class)); container.stop(); pulsarClient.close(); }