diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index f7d092afa..66c55d024 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -93,6 +93,7 @@ import org.springframework.util.StringUtils; * @author Yi Liu * @author Omer Celik * @author Byungjun You + * @author Roman Akentev */ public class KafkaTopicProvisioner implements // @checkstyle:off @@ -638,14 +639,13 @@ public class KafkaTopicProvisioner implements public Collection getPartitionInfoForProducer(final String topicName, final ProducerFactory producerFB, final ExtendedProducerProperties producerProperties) { - return getPartitionsForTopic( - producerProperties.getPartitionCount(), false, () -> { - Producer producer = producerFB.createProducer(); - List partitionsFor = producer - .partitionsFor(topicName); - producer.close(); - return partitionsFor; - }, topicName); + return getPartitionsForTopic(producerProperties.getPartitionCount(), false, + () -> { + try (Producer producer = producerFB + .createProducer()) { + return producer.partitionsFor(topicName); + } + }, topicName); } /**