From 8c65d62a247a0f247ccade8c8a060113bc010937 Mon Sep 17 00:00:00 2001 From: Yang Qiju <362991493@qq.com> Date: Fri, 24 Nov 2017 01:32:10 +0800 Subject: [PATCH] Destroy internal TaskScheduler in container **Cherry-pick to 2.0.x & 1.3.x** --- .../kafka/listener/KafkaMessageListenerContainer.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 5b414829..41926667 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -87,6 +87,7 @@ import org.springframework.util.concurrent.ListenableFutureCallback; * @author Artem Bilan * @author Loic Talhouarne * @author Vladimir Tsanev + * @author Yang Qiju */ public class KafkaMessageListenerContainer extends AbstractMessageListenerContainer { @@ -343,6 +344,8 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private boolean fatalError; + private boolean taskSchedulerExplicitlySet; + @SuppressWarnings("unchecked") ListenerConsumer(GenericMessageListener listener, ListenerType listenerType) { Assert.state(!this.isAnyManualAck || !this.autoCommit, @@ -406,13 +409,14 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener if (this.transactionManager != null) { this.transactionTemplate = new TransactionTemplate(this.transactionManager); Assert.state(!(this.errorHandler instanceof RemainingRecordsErrorHandler), - "You cannot use a 'RemainingRecordsErrorHandler' with transactions"); + "You cannot use a 'RemainingRecordsErrorHandler' with transactions"); } else { this.transactionTemplate = null; } if (this.containerProperties.getScheduler() != null) { this.taskScheduler = this.containerProperties.getScheduler(); + this.taskSchedulerExplicitlySet = true; } else { ThreadPoolTaskScheduler threadPoolTaskScheduler = new ThreadPoolTaskScheduler(); @@ -663,6 +667,9 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener KafkaMessageListenerContainer.this.stop(); } this.monitorTask.cancel(true); + if (!this.taskSchedulerExplicitlySet) { + ((ThreadPoolTaskScheduler) this.taskScheduler).destroy(); + } this.consumer.close(); if (this.logger.isInfoEnabled()) { this.logger.info("Consumer stopped");