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 7bf5f877..88e30d8e 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 @@ -600,7 +600,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener else { this.listener.onMessage(record); } - if (!this.isAnyManualAck) { + if (!this.isAnyManualAck && !this.autoCommit) { this.acks.add(record); } if (this.isRecordAck) { @@ -608,7 +608,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } } catch (Exception e) { - if (this.containerProperties.isAckOnError()) { + if (this.containerProperties.isAckOnError() && !this.autoCommit) { this.acks.add(record); } if (this.containerProperties.getErrorHandler() != null) { @@ -848,6 +848,9 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener @Override public void acknowledge() { try { + if (ListenerConsumer.this.autoCommit) { + throw new IllegalStateException("Manual acks are not allowed when auto commit is used"); + } ListenerConsumer.this.acks.put(this.record); } catch (InterruptedException e) { diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java index dad8029f..13117337 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java @@ -130,6 +130,14 @@ public class ConcurrentMessageListenerContainerTests { template.flush(); assertThat(latch.await(60, TimeUnit.SECONDS)).isTrue(); assertThat(listenerThreadNames).allMatch(threadName -> threadName.contains("-consumer-")); + @SuppressWarnings("unchecked") + List> containers = KafkaTestUtils.getPropertyValue(container, + "containers", List.class); + assertThat(containers.size()).isEqualTo(2); + for (int i = 0; i < 2; i++) { + assertThat(KafkaTestUtils.getPropertyValue(containers.get(i), "listenerConsumer.acks", Collection.class) + .size()).isEqualTo(0); + } container.stop(); this.logger.info("Stop auto"); }