From 9e0e248ba54e91041028bdab6e97fb7c4303d519 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 3 Aug 2016 17:22:38 -0700 Subject: [PATCH] GH-161; Memory Leak with autoCommit Fixes #161 We should not add to the `acks` collection when autoCommit is true. Polishing - PR Comments --- .../kafka/listener/KafkaMessageListenerContainer.java | 7 +++++-- .../listener/ConcurrentMessageListenerContainerTests.java | 8 ++++++++ 2 files changed, 13 insertions(+), 2 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 4ffe1524..5dd81868 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 @@ -599,7 +599,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) { @@ -607,7 +607,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) { @@ -853,6 +853,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 8d31737f..b35858d3 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"); }