From 74689bf007f69d272fccdd89fe81e83071b7dda0 Mon Sep 17 00:00:00 2001 From: Thomas Cheyney Date: Tue, 24 Apr 2018 18:30:32 +0100 Subject: [PATCH] Reuse Kafka consumer Metrics Polishing --- .../binder/kafka/KafkaBinderMetrics.java | 8 ++++- .../binder/kafka/KafkaBinderMetricsTest.java | 35 +++++++++++++++++++ 2 files changed, 42 insertions(+), 1 deletion(-) 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 057c1d3ba..6df5f5011 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 @@ -50,6 +50,7 @@ import org.springframework.util.ObjectUtils; * @author Artem Bilan * @author Oleg Zhurakousky * @author Jon Schneider + * @author Thomas Cheyney */ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener { @@ -65,6 +66,8 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener metadataConsumer; + public KafkaBinderMetrics(KafkaMessageChannelBinder binder, KafkaBinderConfigurationProperties binderConfigurationProperties, ConsumerFactory defaultConsumerFactory, @Nullable MeterRegistry meterRegistry) { @@ -104,7 +107,10 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener metadataConsumer = createConsumerFactory(group).createConsumer()) { + try { + if (metadataConsumer == null) { + metadataConsumer = createConsumerFactory(group).createConsumer(); + } List partitionInfos = metadataConsumer.partitionsFor(topic); List topicPartitions = new LinkedList<>(); for (PartitionInfo partitionInfo : partitionInfos) { 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 2b0ab1d27..394c2f75b 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 @@ -23,9 +23,11 @@ import java.util.Map; import java.util.concurrent.TimeUnit; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.TimeGauge; import io.micrometer.core.instrument.simple.SimpleMeterRegistry; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; +import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.Node; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; @@ -33,6 +35,7 @@ import org.junit.Before; import org.junit.Test; import org.mockito.ArgumentMatchers; import org.mockito.Mock; +import org.mockito.Mockito; import org.mockito.MockitoAnnotations; import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.TopicInformation; @@ -43,6 +46,7 @@ import static org.assertj.core.api.Assertions.assertThat; /** * @author Henryk Konsek + * @author Thomas Cheyney */ public class KafkaBinderMetricsTest { @@ -123,6 +127,37 @@ public class KafkaBinderMetricsTest { assertThat(meterRegistry.getMeters()).isEmpty(); } + @Test + public void createsConsumerOnceWhenInvokedMultipleTimes() { + final List partitions = partitions(new Node(0, null, 0)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); + + metrics.bindTo(meterRegistry); + + TimeGauge gauge = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge(); + gauge.value(TimeUnit.MILLISECONDS); + assertThat(gauge.value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); + + org.mockito.Mockito.verify(this.consumerFactory).createConsumer(); + } + + @Test + public void consumerCreationFailsFirstTime() { + org.mockito.BDDMockito.given(consumerFactory.createConsumer()).willThrow(KafkaException.class) + .willReturn(consumer); + + final List partitions = partitions(new Node(0, null, 0)); + topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); + + metrics.bindTo(meterRegistry); + + TimeGauge gauge = meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge(); + assertThat(gauge.value(TimeUnit.MILLISECONDS)).isEqualTo(0); + assertThat(gauge.value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); + + org.mockito.Mockito.verify(this.consumerFactory, Mockito.times(2)).createConsumer(); + } + private List partitions(Node... nodes) { List partitions = new ArrayList<>(); for (int i = 0; i < nodes.length; i++) {