From 1829bd879f94e0ba56cecbf6bda435fb15450b5f Mon Sep 17 00:00:00 2001 From: Ilayaperumal Gopinathan Date: Mon, 18 Aug 2014 06:33:20 -0700 Subject: [PATCH] Make ConsumerContext disposable - This will allow shutting down the consumer connector for the consumer configuration when the associated context no longer exists. --- .../kafka/support/KafkaConsumerContext.java | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java index e019da4659..3c5905567b 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java @@ -23,6 +23,7 @@ import java.util.Map; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; +import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.integration.kafka.core.KafkaConsumerDefaults; import org.springframework.integration.support.MessageBuilder; @@ -31,9 +32,10 @@ import org.springframework.util.CollectionUtils; /** * @author Soby Chacko + * @author Ilayaperumal Gopinathan * @since 0.5 */ -public class KafkaConsumerContext implements BeanFactoryAware { +public class KafkaConsumerContext implements BeanFactoryAware, DisposableBean { private Map> consumerConfigurations; private String consumerTimeout = KafkaConsumerDefaults.CONSUMER_TIMEOUT; private ZookeeperConnect zookeeperConnect; @@ -77,4 +79,11 @@ public class KafkaConsumerContext implements BeanFactoryAware { public void setZookeeperConnect(final ZookeeperConnect zookeeperConnect) { this.zookeeperConnect = zookeeperConnect; } + + @Override + public void destroy() throws Exception { + for (ConsumerConfiguration config: consumerConfigurations.values()) { + config.getConsumerConnector().shutdown(); + } + } }