From e0fdfdecc8d24063aa037857f095ee2664439be8 Mon Sep 17 00:00:00 2001 From: akenra <37288280+akenra@users.noreply.github.com> Date: Fri, 7 Mar 2025 12:39:41 +0500 Subject: [PATCH] fix(kafka-topic-provisioner): Prevent resource leak on binding producer to KafkaMessageChannelBinder - Producer that fetches partition info now initializes within a try-with-resources block - If exceptions occur on calling producer.partitionsFor(topicName), it's now properly closed and resources are released Signed-off-by: akenra <37288280+akenra@users.noreply.github.com> --- .../provisioning/KafkaTopicProvisioner.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) 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); } /**