From 3faddca84a4dd1be1600999b095521c0d880bf4f Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 22 Jul 2022 20:11:59 -0400 Subject: [PATCH] Rework batch acknowledgment --- ...ulsarBatchMessagingMessageListenerAdapter.java | 15 +++++++++------ .../core/PulsarMessageListenerContainerTests.java | 2 +- 2 files changed, 10 insertions(+), 7 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java index 03a7bd4c..14e90fc4 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java @@ -23,9 +23,11 @@ import java.util.List; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Messages; +import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; -import org.springframework.pulsar.listener.PulsarBatchMessageListener; +import org.springframework.pulsar.listener.Acknowledgement; +import org.springframework.pulsar.listener.PulsarBatchAcknowledgingMessageListener; import org.springframework.pulsar.support.converter.PulsarBatchMessageConverter; import org.springframework.pulsar.support.converter.PulsarBatchMessagingMessageConverter; import org.springframework.pulsar.support.converter.PulsarRecordMessageConverter; @@ -42,7 +44,7 @@ import org.springframework.util.Assert; */ @SuppressWarnings("serial") public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagingMessageListenerAdapter - implements PulsarBatchMessageListener { + implements PulsarBatchAcknowledgingMessageListener { private PulsarBatchMessageConverter batchMessageConverter = new PulsarBatchMessagingMessageConverter(); @@ -63,7 +65,8 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi return this.batchMessageConverter; } - public void received(Consumer consumer, Messages msg) { + @Override + public void received(Consumer consumer, Messages msg, @Nullable Acknowledgement acknowledgement) { Message message; if (!isConsumerRecordList()) { if (isMessageList()) { @@ -81,15 +84,15 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi message = null; // optimization since we won't need any conversion to invoke } logger.debug(() -> "Processing [" + message + "]"); - invoke(msg, consumer, message); + invoke(msg, consumer, message, acknowledgement); } protected void invoke(Object records, Consumer consumer, - final Message messageArg) { + final Message messageArg, Acknowledgement acknowledgement) { Message message = messageArg; try { - Object result = invokeHandler(records, message, consumer, null); + Object result = invokeHandler(records, message, consumer, acknowledgement); // if (result != null) { // handleResult(result, records, 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 ad299ca7..70745ae5 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(1_000); + Thread.sleep(5_000); //Rework this // 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));