From 33603c62f056c03affd5a6994ddafe42ba17ba39 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 8 Apr 2019 14:18:43 -0400 Subject: [PATCH] Transactional binder producer factory With a transactional binder, the producer factory should not be destroyed. Resolves #626 --- .../cloud/stream/binder/kafka/KafkaMessageChannelBinder.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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 c77a870f5..add27a2c1 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 @@ -314,7 +314,9 @@ public class KafkaMessageChannelBinder extends List partitionsFor = producer .partitionsFor(destination.getName()); producer.close(); - ((DisposableBean) producerFB).destroy(); + if (this.transactionManager == null) { + ((DisposableBean) producerFB).destroy(); + } return partitionsFor; }, destination.getName()); this.topicsInUse.put(destination.getName(),