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 da2de243..ce2e6987 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 @@ -32,6 +32,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import java.util.stream.StreamSupport; +import org.apache.commons.logging.LogFactory; import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.DeadLetterPolicy; @@ -45,6 +46,7 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; import org.springframework.context.ApplicationEventPublisher; +import org.springframework.core.log.LogAccessor; import org.springframework.core.task.AsyncTaskExecutor; import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.pulsar.core.PulsarConsumerFactory; @@ -514,44 +516,34 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } private void handleAck(Message message) { - try { - this.consumer.acknowledge(message); - } - catch (PulsarClientException pce) { - this.consumer.negativeAcknowledge(message); - } + AbstractAcknowledgement.handleAckByMessageId(this.consumer, message.getMessageId()); } } - private static final class ConsumerAcknowledgment implements Acknowledgement { + private static abstract class AbstractAcknowledgement implements Acknowledgement { - private final Consumer consumer; + private static final LogAccessor logger = new LogAccessor(LogFactory.getLog(AbstractAcknowledgement.class)); - private final Message message; + protected final Consumer consumer; - ConsumerAcknowledgment(Consumer consumer, Message message) { + AbstractAcknowledgement(Consumer consumer) { this.consumer = consumer; - this.message = message; - } - - @Override - public void acknowledge() { - try { - this.consumer.acknowledge(this.message); - } - catch (PulsarClientException e) { - this.consumer.negativeAcknowledge(this.message); - } } @Override public void acknowledge(MessageId messageId) { + handleAckByMessageId(this.consumer, messageId); + } + + private static void handleAckByMessageId(Consumer consumer, MessageId messageId) { try { - this.consumer.acknowledge(messageId); + consumer.acknowledge(messageId); } - catch (PulsarClientException e) { - this.consumer.negativeAcknowledge(messageId); + catch (PulsarClientException pce) { + AbstractAcknowledgement.logger.warn(pce, + () -> String.format("Acknowledgment failed for message: [%s]", messageId)); + consumer.negativeAcknowledge(messageId); } } @@ -562,34 +554,43 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } catch (PulsarClientException e) { for (MessageId messageId : messageIds) { - try { - this.consumer.acknowledge(messageId); - } - catch (PulsarClientException ex) { - this.consumer.negativeAcknowledge(messageId); - } + handleAckByMessageId(this.consumer, messageId); } } } + @Override + public void nack(MessageId messageId) { + this.consumer.negativeAcknowledge(messageId); + } + + } + + private static final class ConsumerAcknowledgment extends AbstractAcknowledgement { + + private final Message message; + + ConsumerAcknowledgment(Consumer consumer, Message message) { + super(consumer); + this.message = message; + } + + @Override + public void acknowledge() { + acknowledge(this.message.getMessageId()); + } + @Override public void nack() { this.consumer.negativeAcknowledge(this.message); } - @Override - public void nack(MessageId messageId) { - this.consumer.negativeAcknowledge(messageId); - } - } - private static final class ConsumerBatchAcknowledgment implements Acknowledgement { - - private final Consumer consumer; + private static final class ConsumerBatchAcknowledgment extends AbstractAcknowledgement { ConsumerBatchAcknowledgment(Consumer consumer) { - this.consumer = consumer; + super(consumer); } @Override @@ -597,43 +598,11 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess throw new UnsupportedOperationException(); } - @Override - public void acknowledge(MessageId messageId) { - try { - this.consumer.acknowledge(messageId); - } - catch (PulsarClientException e) { - this.consumer.negativeAcknowledge(messageId); - } - } - - @Override - public void acknowledge(List messageIds) { - try { - this.consumer.acknowledge(messageIds); - } - catch (PulsarClientException e) { - for (MessageId messageId : messageIds) { - try { - this.consumer.acknowledge(messageId); - } - catch (PulsarClientException ex) { - this.consumer.negativeAcknowledge(messageId); - } - } - } - } - @Override public void nack() { throw new UnsupportedOperationException(); } - @Override - public void nack(MessageId messageId) { - this.consumer.negativeAcknowledge(messageId); - } - } } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java index 5566f8a1..32529e5f 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java @@ -41,6 +41,7 @@ import java.util.concurrent.atomic.AtomicInteger; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Messages; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; @@ -85,7 +86,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { doAnswer(invocation -> { latch.countDown(); return invocation.callRealMethod(); - }).when(containerConsumer).acknowledge(any(Message.class)); + }).when(containerConsumer).acknowledge(any(MessageId.class)); Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-011"); @@ -166,7 +167,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { doAnswer(invocation -> { ackCallCount.incrementAndGet(); return invocation.callRealMethod(); - }).when(containerConsumer).acknowledge(any(Message.class)); + }).when(containerConsumer).acknowledge(any(MessageId.class)); Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-013"); @@ -187,7 +188,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { final int ackCalls = ackCallCount.get(); if (ackCalls < 5) { await().atMost(Duration.ofSeconds(10)) - .untilAsserted(() -> verify(containerConsumer, atMost(4)).acknowledge(any(Message.class))); + .untilAsserted(() -> verify(containerConsumer, atMost(4)).acknowledge(any(MessageId.class))); await().atMost(Duration.ofSeconds(10)).untilAsserted( () -> verify(containerConsumer, atLeastOnce()).acknowledgeCumulative(any(Message.class))); if (ackCalls == 0) { @@ -236,7 +237,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { doAnswer(invocation -> { latch.countDown(); return invocation.callRealMethod(); - }).when(containerConsumer).acknowledge(any(Message.class)); + }).when(containerConsumer).acknowledge(any(MessageId.class)); Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-014"); @@ -251,7 +252,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { // invocation. assertThat(acksObjects.size()).isEqualTo(10); await().atMost(Duration.ofSeconds(10)) - .untilAsserted(() -> verify(containerConsumer, times(10)).acknowledge(any(Message.class))); + .untilAsserted(() -> verify(containerConsumer, times(10)).acknowledge(any(MessageId.class))); container.stop(); pulsarClient.close();