From 73460aa12e62fa57b1b4411f72be090cf00a3fa7 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Tue, 23 Aug 2016 16:12:24 -0400 Subject: [PATCH] Close ProducerFactory on unbind --- .../binder/kafka/KafkaMessageChannelBinder.java | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 05407533f..777956068 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -63,7 +63,6 @@ import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; -import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.listener.config.ContainerProperties; @@ -206,12 +205,12 @@ public class KafkaMessageChannelBinder extends this.topicsInUse.put(destination, partitions); - ProducerFactory producerFB = getProducerFactory(producerProperties); + DefaultKafkaProducerFactory producerFB = getProducerFactory(producerProperties); KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFB); if (this.producerListener != null) { kafkaTemplate.setProducerListener(this.producerListener); } - return new ProducerConfigurationMessageHandler(kafkaTemplate, destination, producerProperties); + return new ProducerConfigurationMessageHandler(kafkaTemplate, destination, producerProperties, producerFB); } @Override @@ -233,7 +232,7 @@ public class KafkaMessageChannelBinder extends return name; } - private ProducerFactory getProducerFactory( + private DefaultKafkaProducerFactory getProducerFactory( ExtendedProducerProperties producerProperties) { Map props = new HashMap<>(); if (!ObjectUtils.isEmpty(configurationProperties.getConfiguration())) { @@ -550,8 +549,11 @@ public class KafkaMessageChannelBinder extends private boolean running = true; + private final DefaultKafkaProducerFactory producerFactory; + private ProducerConfigurationMessageHandler(KafkaTemplate kafkaTemplate, String topic, - ExtendedProducerProperties producerProperties) { + ExtendedProducerProperties producerProperties, + DefaultKafkaProducerFactory producerFactory) { super(kafkaTemplate); setTopicExpression(new LiteralExpression(topic)); setBeanFactory(KafkaMessageChannelBinder.this.getBeanFactory()); @@ -562,6 +564,7 @@ public class KafkaMessageChannelBinder extends if (producerProperties.getExtension().isSync()) { setSync(true); } + this.producerFactory = producerFactory; } @Override @@ -577,6 +580,7 @@ public class KafkaMessageChannelBinder extends @Override public void stop() { + producerFactory.stop(); this.running = false; }