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;