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 4ebff2a4..a110c609 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 @@ -2233,7 +2233,7 @@ public class KafkaMessageListenerContainer // NOSONAR line count try { this.batchFailed = true; invokeBatchErrorHandler(records, recordList, e); - commitOffsetsIfNeeded(records); + commitOffsetsIfNeededAfterHandlingError(records); } catch (KafkaException ke) { ke.selfLog(ERROR_HANDLER_THREW_AN_EXCEPTION, this.logger); @@ -2254,8 +2254,8 @@ public class KafkaMessageListenerContainer // NOSONAR line count return null; } - private void commitOffsetsIfNeeded(final ConsumerRecords records) { - if ((!this.autoCommit && this.commonErrorHandler.isAckAfterHandle()) + private void commitOffsetsIfNeededAfterHandlingError(final ConsumerRecords records) { + if ((!this.autoCommit && this.commonErrorHandler.isAckAfterHandle() && this.consumerGroupId != null) || this.producer != null) { if (this.remainingRecords != null) { ConsumerRecord firstUncommitted = this.remainingRecords.iterator().next(); @@ -2720,7 +2720,7 @@ public class KafkaMessageListenerContainer // NOSONAR line count } try { invokeErrorHandler(record, iterator, e); - commitOffsetsIfNeeded(record); + commitOffsetsIfNeededAfterHandlingError(record); } catch (KafkaException ke) { ke.selfLog(ERROR_HANDLER_THREW_AN_EXCEPTION, this.logger); @@ -2738,8 +2738,8 @@ public class KafkaMessageListenerContainer // NOSONAR line count return null; } - private void commitOffsetsIfNeeded(final ConsumerRecord cRecord) { - if ((!this.autoCommit && this.commonErrorHandler.isAckAfterHandle()) + private void commitOffsetsIfNeededAfterHandlingError(final ConsumerRecord cRecord) { + if ((!this.autoCommit && this.commonErrorHandler.isAckAfterHandle() && this.consumerGroupId != null) || this.producer != null) { if (this.isManualAck) { this.commitRecovered = true;