From 7355ada4613ad50fe95430f1859d4ea65f004be1 Mon Sep 17 00:00:00 2001 From: Doug Saus Date: Tue, 23 May 2017 18:37:03 -0500 Subject: [PATCH] Fix ConsumerConfig.AUTO_OFFSET_RESET_CONFIG Move set of ConsumerConfig.AUTO_OFFSET_RESET_CONFIG based to before setting of custom kafka properties. This allows users to override this behavior via spring.cloud.stream.kafka.binder.configuration Updated fix to also allow the spring.cloud.stream.kafka.bindings..consumer.startOffset value to override the anonymous-consumer-based value if set Moved setting of auto.offset.reset based on binder configuration below setting of kafka properties so that it has higher preceence. Trailing Spaces --- .../binder/kafka/KafkaMessageChannelBinder.java | 11 +++++++++-- 1 file changed, 9 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 5fb45ac14..a22959e1c 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 @@ -84,6 +84,7 @@ import org.springframework.util.concurrent.ListenableFutureCallback; * @author Mark Fisher * @author Soby Chacko * @author Henryk Konsek + * @author Doug Saus */ public class KafkaMessageChannelBinder extends AbstractMessageChannelBinder, @@ -242,7 +243,7 @@ public class KafkaMessageChannelBinder extends final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets( listenedPartitions); final ContainerProperties containerProperties = - anonymous || extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() ? + anonymous || extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() ? new ContainerProperties(destination.getName()) : new ContainerProperties(topicPartitionInitialOffsets); int concurrency = Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); final ConcurrentMessageListenerContainer messageListenerContainer = @@ -322,6 +323,8 @@ public class KafkaMessageChannelBinder extends props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 100); + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, anonymous ? "latest" : "earliest"); + if (!ObjectUtils.isEmpty(configurationProperties.getConfiguration())) { props.putAll(configurationProperties.getConfiguration()); } @@ -331,8 +334,12 @@ public class KafkaMessageChannelBinder extends if (!ObjectUtils.isEmpty(consumerProperties.getExtension().getConfiguration())) { props.putAll(consumerProperties.getExtension().getConfiguration()); } + props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup); - props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, anonymous ? "latest" : "earliest"); + if (!ObjectUtils.isEmpty(consumerProperties.getExtension().getStartOffset())) { + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, consumerProperties.getExtension().getStartOffset().name()); + } + return new DefaultKafkaConsumerFactory<>(props); }