Fix minor issue processing ackMode
This commit is contained in:
committed by
Gary Russell
parent
394b8a6685
commit
1f52ec5020
@@ -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));
|
||||
|
||||
Reference in New Issue
Block a user