diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java index c7db111b..2abc5ebe 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java @@ -425,14 +425,24 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta } @Override - protected void doStop() { - super.doStop(); + public void stop(Runnable callback) { + super.stop(callback); + cleanUpTaskScheduler(); + } + + private void cleanUpTaskScheduler() { if (!this.taskSchedulerSet && this.taskScheduler != null) { ((ThreadPoolTaskScheduler) this.taskScheduler).shutdown(); this.taskScheduler = null; } } + @Override + protected void doStop() { + super.doStop(); + cleanUpTaskScheduler(); + } + protected void actualStart() { this.aborted = false; this.hasStopped = false; @@ -975,6 +985,12 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta // default empty } + @Override + public void destroy() { + super.destroy(); + cleanUpTaskScheduler(); + } + /** * The consumer object. */ diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ContainerShutDownTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ContainerShutDownTests.java index aec9f33c..2927b563 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ContainerShutDownTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ContainerShutDownTests.java @@ -18,6 +18,7 @@ package org.springframework.amqp.rabbit.listener; import java.util.Map; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import com.rabbitmq.client.AMQP.BasicProperties; @@ -139,4 +140,26 @@ public class ContainerShutDownTests { } } + @Test + void directMessageListenerContainerShutdownsItsSchedulerOnStopWithCallback() { + DirectMessageListenerContainer container = new DirectMessageListenerContainer(); + CachingConnectionFactory cf = new CachingConnectionFactory("localhost"); + container.setConnectionFactory(cf); + container.setQueueNames("test.shutdown"); + container.setMessageListener(m -> { + }); + + container.start(); + + ScheduledExecutorService scheduledExecutorService = + TestUtils.getPropertyValue(container, "taskScheduler.scheduledExecutor", ScheduledExecutorService.class); + + container.stop(() -> { + }); + + cf.destroy(); + + assertThat(scheduledExecutorService.isShutdown()).isTrue(); + } + }