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 67de6d30d..007a75dcf 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 @@ -504,7 +504,7 @@ public class KafkaMessageChannelBinder extends boolean resetOffsets = extendedConsumerProperties.getExtension().isResetOffsets(); final Object resetTo = consumerFactory.getConfigurationProperties().get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG); final AtomicBoolean initialAssignment = new AtomicBoolean(true); - if (!"earliest".equals(resetTo) && "!latest".equals(resetTo)) { + if (!"earliest".equals(resetTo) && !"latest".equals(resetTo)) { logger.warn("no (or unknown) " + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG + " property cannot reset"); resetOffsets = false;