From 4ff4507741141864fc15baa3a409d80d12aa6d84 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 28 Sep 2017 15:51:15 +0100 Subject: [PATCH] GH-206: Close Consumer/Producer in provisioning Fixes https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/206 Close the consumer and producer after retrieving the current partition count. **Cherry pick/back port to 0.11 and 1.2.x, 2.0.x** * Destroy the Producer Factory --- .../kafka/KafkaMessageChannelBinder.java | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 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 b82f4c4e5..6ebb08548 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 @@ -21,12 +21,15 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.UUID; import java.util.concurrent.Callable; +import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.serialization.ByteArrayDeserializer; @@ -154,10 +157,16 @@ public class KafkaMessageChannelBinder extends producerProperties.getPartitionCount(), false, new Callable>() { + @Override public Collection call() throws Exception { - return producerFB.createProducer().partitionsFor(destination.getName()); + Producer producer = producerFB.createProducer(); + List partitionsFor = producer.partitionsFor(destination.getName()); + producer.close(); + producerFB.destroy(); + return partitionsFor; } + }); this.topicsInUse.put(destination.getName(), new TopicInformation(null, partitions)); if (producerProperties.getPartitionCount() < partitions.size()) { @@ -238,10 +247,15 @@ public class KafkaMessageChannelBinder extends Collection allPartitions = provisioningProvider.getPartitionsForTopic(partitionCount, extendedConsumerProperties.getExtension().isAutoRebalanceEnabled(), new Callable>() { + @Override public Collection call() throws Exception { - return consumerFactory.createConsumer().partitionsFor(destination.getName()); + Consumer consumer = consumerFactory.createConsumer(); + List partitionsFor = consumer.partitionsFor(destination.getName()); + consumer.close(); + return partitionsFor; } + }); Collection listenedPartitions;