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 8f14c077..62b11604 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 @@ -24,10 +24,11 @@ import java.util.List; import java.util.Map; import java.util.Properties; import java.util.Set; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; -import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; import java.util.stream.Stream; import java.util.stream.StreamSupport; @@ -74,7 +75,7 @@ import io.micrometer.observation.ObservationRegistry; */ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMessageListenerContainer { - private volatile Future listenerConsumerFuture; + private volatile CompletableFuture listenerConsumerFuture; private volatile Listener listenerConsumer; @@ -82,6 +83,10 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess private final AbstractPulsarMessageListenerContainer thisOrParentContainer; + private AtomicReference listenerConsumerThread; + + private final AtomicBoolean receiveInProgress = new AtomicBoolean(); + public DefaultPulsarMessageListenerContainer(PulsarConsumerFactory pulsarConsumerFactory, PulsarContainerProperties pulsarContainerProperties) { this(pulsarConsumerFactory, pulsarContainerProperties, null); @@ -113,7 +118,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.getObservationRegistry()); setRunning(true); this.startLatch = new CountDownLatch(1); - this.listenerConsumerFuture = consumerExecutor.submit(this.listenerConsumer); + this.listenerConsumerFuture = consumerExecutor.submitCompletable(this.listenerConsumer); try { if (!this.startLatch.await(containerProperties.getConsumerStartTimeout().toMillis(), @@ -133,6 +138,24 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess setRunning(false); this.logger.info("Pausing this consumer."); this.listenerConsumer.consumer.pause(); + if (this.listenerConsumerThread != null) { + // if there is a receive operation already in progress, we want to interrupt + // the listener thread. + if (this.receiveInProgress.get()) { + // All the records received so far in the current batch receive will be + // re-delivered. + this.listenerConsumerThread.get().interrupt(); + } + // if there is something other than receive operations are in progress, + // such as ack operations, wait for the listener thread to complete them. + try { + this.listenerConsumerThread.get().join(); + } + catch (InterruptedException e) { + this.logger.error(e, () -> "Interrupting the main thread"); + Thread.currentThread().interrupt(); + } + } try { this.logger.info("Closing this consumer."); this.listenerConsumer.consumer.close(); @@ -286,6 +309,8 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess @Override public void run() { + DefaultPulsarMessageListenerContainer.this.listenerConsumerThread = new AtomicReference<>( + Thread.currentThread()); publishConsumerStartingEvent(); publishConsumerStartedEvent(); AtomicBoolean inRetryMode = new AtomicBoolean(false); @@ -296,13 +321,28 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess // Always receive messages in batch mode. try { if (!inRetryMode.get() && !messagesPendingInBatch.get()) { + DefaultPulsarMessageListenerContainer.this.receiveInProgress.set(true); messages = this.consumer.batchReceive(); } } catch (PulsarClientException e) { - DefaultPulsarMessageListenerContainer.this.logger.error(e, () -> "Error receiving messages."); + if (e.getCause() instanceof InterruptedException) { + DefaultPulsarMessageListenerContainer.this.logger.debug(e, + () -> "Error receiving messages due to a thread interrupt call from upstream."); + } + else { + DefaultPulsarMessageListenerContainer.this.logger.error(e, () -> "Error receiving messages."); + } + messages = null; } - Assert.isTrue(messages != null, "Messages cannot be null."); + finally { + DefaultPulsarMessageListenerContainer.this.receiveInProgress.set(false); + } + + if (messages == null) { + continue; + } + if (this.isBatchListener) { if (!inRetryMode.get() && !messagesPendingInBatch.get()) { messageList = new ArrayList<>(); @@ -345,7 +385,8 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess messageList, e); } else { - // the whole batch is negatively acknowledged in the event of + // the whole batch is negatively acknowledged in the event + // of // an exception from the handler method. this.consumer.negativeAcknowledge(messages); } 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 c57adeca..f6c15582 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 @@ -361,4 +361,61 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarClient.close(); } + @Test + void messagesAreProperlyAckdOnContainerStopBeforeExitingListenerThread() throws Exception { + Map config = new HashMap<>(); + config.put("topicNames", Set.of("duplicate-message-test")); + config.put("subscriptionName", "duplicate-sub-1"); + final PulsarClient pulsarClient = PulsarClient.builder() + .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); + final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( + pulsarClient, config); + + PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); + final AtomicInteger counter1 = new AtomicInteger(0); + pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> { + counter1.getAndIncrement(); + }); + pulsarContainerProperties.setSchema(Schema.STRING); + DefaultPulsarMessageListenerContainer container1 = new DefaultPulsarMessageListenerContainer<>( + pulsarConsumerFactory, pulsarContainerProperties); + container1.start(); + + Map prodConfig = Collections.singletonMap("topicName", "duplicate-message-test"); + final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( + pulsarClient, prodConfig); + final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + pulsarTemplate.send("hello john doe"); + + while (counter1.get() == 0) { + // busy wait until counter1 is > 0 + } + // When we stop, if any acks are in progress, that should all be + // taken care of before exiting the listener thread, so that the + // next consumer under the same subscription will not receive the + // unacked message. + container1.stop(); + + final AtomicInteger counter2 = new AtomicInteger(0); + pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> { + counter2.getAndIncrement(); + }); + pulsarContainerProperties.setSchema(Schema.STRING); + DefaultPulsarMessageListenerContainer container2 = new DefaultPulsarMessageListenerContainer<>( + pulsarConsumerFactory, pulsarContainerProperties); + container2.start(); + + pulsarTemplate.send("hello john doe"); + + while (counter2.get() == 0) { + // busy wait until counter2 > 0 + } + // Asserting that both consumers are only receiving the expected data. + assertThat(counter1.get()).isEqualTo(1); + assertThat(counter2.get()).isEqualTo(1); + + container2.stop(); + pulsarClient.close(); + } + } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java index a3c2427d..05683425 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java @@ -262,6 +262,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS // Normal consumer should receive 5 msg + 1 re-delivery assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); container.stop(); + dlqContainer.stop(); pulsarClient.close(); } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/observation/ObservationIntegrationTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/observation/ObservationIntegrationTests.java index 593bab7e..10a462dd 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/observation/ObservationIntegrationTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/observation/ObservationIntegrationTests.java @@ -93,14 +93,6 @@ public class ObservationIntegrationTests extends SampleTestRunner implements Pul assertThat(listen2Completed).withFailMessage( "Message %s not received in listen2 (latchesByMessageListen1 = %s and latchesByMessageListen2 = %s)", msg, listeners.latchesByMessageListen1, listeners.latchesByMessageListen2).isTrue(); - - // Without this sleep, the 2nd tracingSetup run sometimes fails due to - // messages from 1st run being - // delivered during the 2nd run. The test runs share the same listener - // config, including the - // same subscription names. Seems like the listener in run2 is getting - // message from run1. - Thread.sleep(5000); } List finishedSpans = bb.getFinishedSpans();