From 0bbb23a57f0335113da984a6bce031f01785aa05 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 11 Apr 2019 10:56:59 -0400 Subject: [PATCH] Fix close producers after a rebalance With a transactional container if committing the offsets fails during a rebalance, we would leave the producers open. Move the close to a finally block. **cherry-pick to 2.0.x** --- .../KafkaMessageListenerContainer.java | 28 +++++++++++-------- 1 file changed, 16 insertions(+), 12 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 99d803b6..d9d45ffb 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -472,19 +472,23 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener @Override public void onPartitionsRevoked(Collection partitions) { - if (this.consumerAwareListener != null) { - this.consumerAwareListener.onPartitionsRevokedBeforeCommit(consumer, partitions); + try { + if (this.consumerAwareListener != null) { + this.consumerAwareListener.onPartitionsRevokedBeforeCommit(consumer, partitions); + } + else { + this.userListener.onPartitionsRevoked(partitions); + } + // Wait until now to commit, in case the user listener added acks + commitPendingAcks(); + if (this.consumerAwareListener != null) { + this.consumerAwareListener.onPartitionsRevokedAfterCommit(consumer, partitions); + } } - else { - this.userListener.onPartitionsRevoked(partitions); - } - // Wait until now to commit, in case the user listener added acks - commitPendingAcks(); - if (this.consumerAwareListener != null) { - this.consumerAwareListener.onPartitionsRevokedAfterCommit(consumer, partitions); - } - if (ListenerConsumer.this.kafkaTxManager != null) { - closeProducers(partitions); + finally { + if (ListenerConsumer.this.kafkaTxManager != null) { + closeProducers(partitions); + } } }