diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java index 03dbcbd95..7b4cc82f4 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java @@ -20,6 +20,7 @@ import java.util.HashMap; import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -75,10 +76,10 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener metadataConsumer; + private Map> metadataConsumers; private int timeout = DEFAULT_TIMEOUT; - + public KafkaBinderMetrics(KafkaMessageChannelBinder binder, KafkaBinderConfigurationProperties binderConfigurationProperties, ConsumerFactory defaultConsumerFactory, @Nullable MeterRegistry meterRegistry) { @@ -87,6 +88,7 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener(); } public KafkaBinderMetrics(KafkaMessageChannelBinder binder, @@ -126,13 +128,9 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener metadataConsumer = metadataConsumers.computeIfAbsent( + group, + g -> createConsumerFactory().createConsumer(g, "monitoring")); synchronized (metadataConsumer) { List partitionInfos = metadataConsumer.partitionsFor(topic); List topicPartitions = new LinkedList<>(); @@ -171,21 +169,24 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener createConsumerFactory(String group) { + private ConsumerFactory createConsumerFactory() { if (this.defaultConsumerFactory == null) { - Map props = new HashMap<>(); - props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); - Map mergedConfig = this.binderConfigurationProperties.mergedConsumerConfiguration(); - if (!ObjectUtils.isEmpty(mergedConfig)) { - props.putAll(mergedConfig); + synchronized (this) { + if (this.defaultConsumerFactory == null) { + Map props = new HashMap<>(); + props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); + Map mergedConfig = this.binderConfigurationProperties.mergedConsumerConfiguration(); + if (!ObjectUtils.isEmpty(mergedConfig)) { + props.putAll(mergedConfig); + } + if (!props.containsKey(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)) { + props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, + this.binderConfigurationProperties.getKafkaConnectionString()); + } + this.defaultConsumerFactory = new DefaultKafkaConsumerFactory<>(props); + } } - if (!props.containsKey(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)) { - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, - this.binderConfigurationProperties.getKafkaConnectionString()); - } - props.put("group.id", group); - this.defaultConsumerFactory = new DefaultKafkaConsumerFactory<>(props); } return this.defaultConsumerFactory; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java index 956397967..6b892b35c 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java @@ -42,6 +42,7 @@ import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfi import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; /** * @author Henryk Konsek @@ -73,7 +74,7 @@ public class KafkaBinderMetricsTest { @Before public void setup() { MockitoAnnotations.initMocks(this); - org.mockito.BDDMockito.given(consumerFactory.createConsumer()).willReturn(consumer); + org.mockito.BDDMockito.given(consumerFactory.createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any())).willReturn(consumer); org.mockito.BDDMockito.given(binder.getTopicsInUse()).willReturn(topicsInUse); metrics = new KafkaBinderMetrics(binder, kafkaBinderConfigurationProperties, consumerFactory, null); org.mockito.BDDMockito.given(consumer.endOffsets(ArgumentMatchers.anyCollection())) @@ -138,12 +139,12 @@ public class KafkaBinderMetricsTest { gauge.value(); assertThat(gauge.value()).isEqualTo(1000.0); - org.mockito.Mockito.verify(this.consumerFactory).createConsumer(); + org.mockito.Mockito.verify(this.consumerFactory).createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()); } @Test public void consumerCreationFailsFirstTime() { - org.mockito.BDDMockito.given(consumerFactory.createConsumer()).willThrow(KafkaException.class) + org.mockito.BDDMockito.given(consumerFactory.createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any())).willThrow(KafkaException.class) .willReturn(consumer); final List partitions = partitions(new Node(0, null, 0)); @@ -155,7 +156,33 @@ public class KafkaBinderMetricsTest { assertThat(gauge.value()).isEqualTo(0); assertThat(gauge.value()).isEqualTo(1000.0); - org.mockito.Mockito.verify(this.consumerFactory, Mockito.times(2)).createConsumer(); + org.mockito.Mockito.verify(this.consumerFactory, Mockito.times(2)).createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()); + } + + @Test + public void createOneConsumerPerGroup() { + final List partitions1 = partitions(new Node(0, null, 0)); + final List partitions2 = partitions(new Node(0, null, 0)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group1-metrics", partitions1, false)); + topicsInUse.put("test2", new TopicInformation("group2-metrics", partitions2, false)); + + metrics.bindTo(meterRegistry); + + KafkaConsumer consumer2 = mock(KafkaConsumer.class); + org.mockito.BDDMockito.given(consumerFactory.createConsumer(ArgumentMatchers.eq("group2-metrics"), ArgumentMatchers.any())) + .willReturn(consumer2); + org.mockito.BDDMockito.given(consumer2.endOffsets(ArgumentMatchers.anyCollection())) + .willReturn(java.util.Collections.singletonMap(new TopicPartition("test2", 0), 50L)); + + Gauge gauge1 = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge(); + Gauge gauge2 = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group2-metrics").tag("topic", "test2").gauge(); + gauge1.value(); + gauge2.value(); + assertThat(gauge1.value()).isEqualTo(1000.0); + assertThat(gauge2.value()).isEqualTo(50.0); + + org.mockito.Mockito.verify(this.consumerFactory).createConsumer(ArgumentMatchers.eq("group1-metrics"), ArgumentMatchers.any()); + org.mockito.Mockito.verify(this.consumerFactory).createConsumer(ArgumentMatchers.eq("group2-metrics"), ArgumentMatchers.any()); } private List partitions(Node... nodes) {