From 0d5d092f84a8c8bccb9f4685b567994c1b0ba0c1 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 19 Oct 2018 15:42:35 -0400 Subject: [PATCH] Fix merge issues in 40c1a8b8 --- .../cloud/stream/binder/kafka/KafkaMessageChannelBinder.java | 4 ++-- 1 file changed, 2 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 eb6b8a0e6..67de6d30d 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 @@ -567,13 +567,13 @@ public class KafkaMessageChannelBinder extends // not just the ones this binding is listening to; doesn't seem right for a health check. Collection partitionInfos = getPartitionInfo(destination.getName(), consumerProperties, consumerFactory, -1); - this.topicsInUse.put(destination.getName(), new TopicInformation(group, partitionInfos, false)); + this.topicsInUse.put(destination.getName(), new TopicInformation(consumerGroup, partitionInfos, false)); } else { for (int i = 0; i < topics.length; i++) { Collection partitionInfos = getPartitionInfo(topics[i], consumerProperties, consumerFactory, -1); - this.topicsInUse.put(topics[i], new TopicInformation(group, partitionInfos, false)); + this.topicsInUse.put(topics[i], new TopicInformation(consumerGroup, partitionInfos, false)); } }