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 62b11604..08190e02 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 @@ -83,7 +83,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess private final AbstractPulsarMessageListenerContainer thisOrParentContainer; - private AtomicReference listenerConsumerThread; + private final AtomicReference listenerConsumerThread = new AtomicReference<>(); private final AtomicBoolean receiveInProgress = new AtomicBoolean(); @@ -138,7 +138,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess setRunning(false); this.logger.info("Pausing this consumer."); this.listenerConsumer.consumer.pause(); - if (this.listenerConsumerThread != null) { + if (this.listenerConsumerThread.get() != null) { // if there is a receive operation already in progress, we want to interrupt // the listener thread. if (this.receiveInProgress.get()) { @@ -309,8 +309,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess @Override public void run() { - DefaultPulsarMessageListenerContainer.this.listenerConsumerThread = new AtomicReference<>( - Thread.currentThread()); + DefaultPulsarMessageListenerContainer.this.listenerConsumerThread.set(Thread.currentThread()); publishConsumerStartingEvent(); publishConsumerStartedEvent(); AtomicBoolean inRetryMode = new AtomicBoolean(false);