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 11c2db0b9..f5308b63b 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,10 +20,13 @@ import java.util.HashMap; import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; +import io.micrometer.core.instrument.Gauge; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Tags; +import io.micrometer.core.instrument.TimeGauge; import io.micrometer.core.instrument.binder.MeterBinder; - import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.clients.consumer.Consumer; @@ -53,7 +56,7 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener calculateConsumerLagOnTopic(topic, group)); + TimeGauge.builder(METRIC_NAME, this, TimeUnit.MILLISECONDS, + o -> calculateConsumerLagOnTopic(topic, group)) + .tag("group", group) + .tag("topic", topic) + .description("Consumer lag for a particular group and topic") + .register(registry); } } @@ -104,6 +111,7 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener endOffsets = metadataConsumer.endOffsets(topicPartitions); for (Map.Entry endOffset : endOffsets.entrySet()) { 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 8de9b9eaa..2b0ab1d27 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 @@ -20,9 +20,9 @@ import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import io.micrometer.core.instrument.MeterRegistry; -import io.micrometer.core.instrument.search.Search; import io.micrometer.core.instrument.simple.SimpleMeterRegistry; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; @@ -31,6 +31,7 @@ import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.junit.Before; import org.junit.Test; +import org.mockito.ArgumentMatchers; import org.mockito.Mock; import org.mockito.MockitoAnnotations; @@ -71,20 +72,20 @@ public class KafkaBinderMetricsTest { org.mockito.BDDMockito.given(consumerFactory.createConsumer()).willReturn(consumer); org.mockito.BDDMockito.given(binder.getTopicsInUse()).willReturn(topicsInUse); metrics = new KafkaBinderMetrics(binder, kafkaBinderConfigurationProperties, consumerFactory, null); - org.mockito.BDDMockito.given(consumer.endOffsets(org.mockito.Matchers.anyCollectionOf(TopicPartition.class))) + org.mockito.BDDMockito.given(consumer.endOffsets(ArgumentMatchers.anyCollection())) .willReturn(java.util.Collections.singletonMap(new TopicPartition(TEST_TOPIC, 0), 1000L)); } @Test public void shouldIndicateLag() { - org.mockito.BDDMockito.given(consumer.committed(org.mockito.Matchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); + org.mockito.BDDMockito.given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); List partitions = partitions(new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - Search group = meterRegistry.find(String.format("%s.%s.%s.lag", KafkaBinderMetrics.METRIC_PREFIX, "group", TEST_TOPIC)); - assertThat(group.gauge().value()).isEqualTo(500.0); + assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge() + .value(TimeUnit.MILLISECONDS)).isEqualTo(500.0); } @Test @@ -92,15 +93,15 @@ public class KafkaBinderMetricsTest { Map endOffsets = new HashMap<>(); endOffsets.put(new TopicPartition(TEST_TOPIC, 0), 1000L); endOffsets.put(new TopicPartition(TEST_TOPIC, 1), 1000L); - org.mockito.BDDMockito.given(consumer.endOffsets(org.mockito.Matchers.anyCollectionOf(TopicPartition.class))).willReturn(endOffsets); - org.mockito.BDDMockito.given(consumer.committed(org.mockito.Matchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); + org.mockito.BDDMockito.given(consumer.endOffsets(ArgumentMatchers.anyCollection())).willReturn(endOffsets); + org.mockito.BDDMockito.given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); List partitions = partitions(new Node(0, null, 0), new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - Search group = meterRegistry.find(String.format("%s.%s.%s.lag", KafkaBinderMetrics.METRIC_PREFIX, "group", TEST_TOPIC)); - assertThat(group.gauge().value()).isEqualTo(1000.0); + assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge() + .value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); } @Test @@ -110,8 +111,8 @@ public class KafkaBinderMetricsTest { org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - Search group = meterRegistry.find(String.format("%s.%s.%s.lag", KafkaBinderMetrics.METRIC_PREFIX, "group", TEST_TOPIC)); - assertThat(group.gauge().value()).isEqualTo(1000.0); + assertThat(meterRegistry.get(KafkaBinderMetrics.METRIC_NAME).tag("group", "group").tag("topic", TEST_TOPIC).timeGauge() + .value(TimeUnit.MILLISECONDS)).isEqualTo(1000.0); } @Test diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java index 80ef873bb..4b4362cb3 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java @@ -74,15 +74,13 @@ public class KafkaBinderActuatorTests { @Test public void testKafkaBinderMetricsExposed() { - Search search = this.meterRegistry.find( - String.format("%s.%s.%s.lag", "spring.cloud.stream.binder.kafka", TEST_CONSUMER_GROUP, Sink.INPUT)); - - assertThat(search.gauge()).isNotNull(); - this.kafkaTemplate.send(Sink.INPUT, null, "foo".getBytes()); this.kafkaTemplate.flush(); - assertThat(search.gauge().value()).isGreaterThan(0); + assertThat(this.meterRegistry.get("spring.cloud.stream.binder.kafka.offset") + .tag("group", TEST_CONSUMER_GROUP) + .tag("topic", Sink.INPUT) + .timeGauge().value()).isGreaterThan(0); } @EnableBinding(Sink.class)