From 1f52ec5020da96bcb1dd1fe9d0ff21995d4c523a Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 19 Feb 2021 12:58:59 -0500 Subject: [PATCH] Fix minor issue processing ackMode --- .../stream/binder/kafka/KafkaMessageChannelBinder.java | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) 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 7858f5969..beae1040a 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 @@ -707,14 +707,10 @@ public class KafkaMessageChannelBinder extends else { messageListenerContainer.getContainerProperties() .setAckOnError(isAutoCommitOnError(extendedConsumerProperties)); - if (extendedConsumerProperties.getExtension().isAckEachRecord()) { - messageListenerContainer.getContainerProperties() - .setAckMode(ContainerProperties.AckMode.RECORD); - } } } } - else { + if (ackMode != null) { if ((extendedConsumerProperties.isBatchMode() && ackMode != ContainerProperties.AckMode.RECORD) || !extendedConsumerProperties.isBatchMode()) { messageListenerContainer.getContainerProperties() @@ -722,8 +718,6 @@ public class KafkaMessageChannelBinder extends } } - - if (this.logger.isDebugEnabled()) { this.logger.debug("Listened partitions: " + StringUtils.collectionToCommaDelimitedString(listenedPartitions));