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 0675730e..d2db82ff 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 @@ -21,7 +21,6 @@ import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; -import java.util.Optional; import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -167,7 +166,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess private Consumer consumer; - private final Set> nackableMessages = new HashSet<>(); + private final Set nackableMessages = new HashSet<>(); private final PulsarContainerProperties containerProperties = getPulsarContainerProperties(); @@ -291,7 +290,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.consumer.negativeAcknowledge(message); } else if (this.containerProperties.getAckMode() == PulsarContainerProperties.AckMode.BATCH) { - this.nackableMessages.add(message); + this.nackableMessages.add(message.getMessageId()); } } } @@ -316,13 +315,9 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } else { for (Message message : messages) { - final Optional> foundMessage = this.nackableMessages.stream() - .filter(m -> m.getMessageId().compareTo(message.getMessageId()) == 0) - .findFirst(); - if (foundMessage.isPresent()) { - final Message msg = foundMessage.get(); - this.consumer.negativeAcknowledge(msg); - this.nackableMessages.remove(msg); + if (this.nackableMessages.contains(message.getMessageId())) { + this.consumer.negativeAcknowledge(message); + this.nackableMessages.remove(message.getMessageId()); } else { handleAck(message); 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 70745ae5..ad299ca7 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 @@ -172,7 +172,7 @@ class PulsarMessageListenerContainerTests extends AbstractContainerBaseTests { pulsarTemplate.sendAsync("hello john doe"); } assertThat(latch.await(30, TimeUnit.SECONDS)).isTrue(); - Thread.sleep(5_000); //Rework this + Thread.sleep(1_000); // Half of the message get acknowledged, and the other half gets negatively acknowledged. verify(containerConsumer, times(5)).acknowledge(any(Message.class)); verify(containerConsumer, times(5)).negativeAcknowledge(any(Message.class));