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); + } } }