From 2360a9d2d1e1b5d16e526e7bfda14cc5474b4dce Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 12 Oct 2022 18:38:31 -0400 Subject: [PATCH] Consumer restart causes duplicate messages When we stop the consumer, there is a race condition that forces the messages to not get acknowledged. This results in duplicate messages on the restart of the consumer. This PR addresses the issue by waiting for the consumer thread to complete before stopping the consumer. Resolves https://github.com/spring-projects-experimental/spring-pulsar/issues/161 --- ...DefaultPulsarMessageListenerContainer.java | 53 +++++++++++++++-- .../core/ConsumerAcknowledgmentTests.java | 57 +++++++++++++++++++ ...ltPulsarMessageListenerContainerTests.java | 1 + .../ObservationIntegrationTests.java | 8 --- 4 files changed, 105 insertions(+), 14 deletions(-) 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();