From b2b1438ad77bb7d28b01c98c1dbf27691d11d9fe Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 19 Jul 2019 18:47:41 -0400 Subject: [PATCH] Fix resetOffsets for manual partition assignment --- .../stream/binder/kafka/KafkaMessageChannelBinder.java | 9 +++++++-- 1 file changed, 7 insertions(+), 2 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 1664baab5..64b70d8f0 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 @@ -733,12 +733,17 @@ public class KafkaMessageChannelBinder extends } private Object checkReset(boolean resetOffsets, final Object resetTo) { - if (resetOffsets && !"earliest".equals(resetTo) && !"latest".equals(resetTo)) { + if (!resetOffsets) { + return null; + } + else if (!"earliest".equals(resetTo) && !"latest".equals(resetTo)) { logger.warn("no (or unknown) " + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG + " property cannot reset"); return null; } - return resetTo; + else { + return resetTo; + } } @Override