From 546fe6934683844668bb84bfce80aac4d27bf309 Mon Sep 17 00:00:00 2001 From: "Nadelson, Mark" Date: Wed, 21 Sep 2016 17:37:52 -0400 Subject: [PATCH] set Kafka ackMode to MANUAL if autoCommitOffset = false polishing --- .../cloud/stream/binder/kafka/KafkaMessageChannelBinder.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 9894872e8..f3098034b 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -61,6 +61,7 @@ import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.listener.config.ContainerProperties; @@ -335,6 +336,9 @@ public class KafkaMessageChannelBinder extends }; messageListenerContainer.setConcurrency(concurrency); messageListenerContainer.getContainerProperties().setAckOnError(isAutoCommitOnError(properties)); + if (!properties.getExtension().isAutoCommitOffset()) { + messageListenerContainer.getContainerProperties().setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL); + } if (this.logger.isDebugEnabled()) { this.logger.debug(