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.
This commit is contained in:
@@ -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<T> extends AbstractPulsarMess
|
||||
}
|
||||
if (this.containerProperties.getAckMode() == PulsarContainerProperties.AckMode.BATCH) {
|
||||
try {
|
||||
this.consumer.acknowledge(messages);
|
||||
if (isSharedSubsriptionType()) {
|
||||
this.consumer.acknowledge(messages);
|
||||
}
|
||||
else {
|
||||
final Stream<Message<T>> stream = StreamSupport.stream(messages.spliterator(),
|
||||
true);
|
||||
Message<T> 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<T> extends AbstractPulsarMess
|
||||
}
|
||||
}
|
||||
|
||||
private boolean isSharedSubsriptionType() {
|
||||
return this.containerProperties.getSubscriptionType() == SubscriptionType.Shared
|
||||
|| this.containerProperties.getSubscriptionType() == SubscriptionType.Key_Shared;
|
||||
}
|
||||
|
||||
private void handleAcks(Messages<T> messages) {
|
||||
if (this.nackableMessages.isEmpty()) {
|
||||
try {
|
||||
if (messages.size() > 0) {
|
||||
this.consumer.acknowledge(messages);
|
||||
if (isSharedSubsriptionType()) {
|
||||
this.consumer.acknowledge(messages);
|
||||
}
|
||||
else {
|
||||
final Stream<Message<T>> stream = StreamSupport.stream(messages.spliterator(), true);
|
||||
Message<T> last = stream.reduce((a, b) -> b).orElse(null);
|
||||
this.consumer.acknowledgeCumulative(last);
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (PulsarClientException pce) {
|
||||
|
||||
@@ -63,7 +63,7 @@ public class PulsarContainerProperties {
|
||||
|
||||
private String subscriptionName;
|
||||
|
||||
private SubscriptionType subscriptionType;
|
||||
private SubscriptionType subscriptionType = SubscriptionType.Exclusive;
|
||||
|
||||
private Schema<?> schema;
|
||||
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user