From bc1936eb28bea26028ca72bc86d8ae51acf04be0 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 1 Nov 2019 15:51:14 -0400 Subject: [PATCH] KafkaStreams binder metrics - duplicate entries When the same metric name is repeated, there are some registry implementations such as the micrometer Prometheus registry fail to register the duplicate entry. Fixing this issue by restricting the duplicate metric names not to be registered. Also, address an issue with multiple processors and metrics in the same application by prepending the application ID of the Kafka Streams processor in the metric name itself. Resolves #788 --- docs/src/main/asciidoc/kafka-streams.adoc | 4 ++ .../streams/KafkaStreamsBinderMetrics.java | 59 ++++++++++++------- ...StreamsBinderSupportAutoConfiguration.java | 12 ++-- .../kafka/streams/KafkaStreamsRegistry.java | 29 ++------- .../streams/StreamsBuilderFactoryManager.java | 15 ++--- ...reamsInteractiveQueryIntegrationTests.java | 10 ++-- 6 files changed, 66 insertions(+), 63 deletions(-) diff --git a/docs/src/main/asciidoc/kafka-streams.adoc b/docs/src/main/asciidoc/kafka-streams.adoc index c8f834099..e2f31b768 100644 --- a/docs/src/main/asciidoc/kafka-streams.adoc +++ b/docs/src/main/asciidoc/kafka-streams.adoc @@ -1094,6 +1094,10 @@ All dashes in the original metric information is replaced with dots. For e.g. the metric name `network-io-total` from the metric group `consumer-metrics` is available in the micrometer registry as `consumer.metrics.network.io.total`. Similarly, the metric `commit-total` from `stream-metrics` is available as `stream.metrics.commit.total`. +If you have multiple Kafka Streams processors in the same application, then the metric name will be prepended with the corresponding application ID of the Kafka Streams. +The application ID in this case will be preserved as is, i.e. no dashes will be converted to dots etc. +For example, if the application ID of the first processor is `processor-1`, then the metric name `network-io-total` from the metric group `consumer-metrics` is available in the micrometer registry as `processor-1.consumer.metrics.network.io.total`. + You can either programmatically access the Micrometer `MeterRegistry` in the application and then iterate through the available gauges or use Spring Boot actuator to access the metrics through a REST endpoint. When accessing through the Boot actuator endpoint, make sure to add `metrics` to the property `management.endpoints.web.exposure.include`. Then you can access `/acutator/metrics` to get a list of all the available metrics which then can be individually accessed through the same URL (`/actuator/metrics/`). diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java index 4e860a436..899c595e9 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderMetrics.java @@ -16,7 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.HashSet; import java.util.Map; +import java.util.Set; import java.util.function.ToDoubleFunction; import io.micrometer.core.instrument.Gauge; @@ -25,6 +27,9 @@ import io.micrometer.core.instrument.binder.MeterBinder; import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; import org.apache.kafka.streams.KafkaStreams; +import org.apache.kafka.streams.StreamsConfig; + +import org.springframework.kafka.config.StreamsBuilderFactoryBean; /** * Kafka Streams binder metrics implementation that exports the metrics available @@ -35,8 +40,6 @@ import org.apache.kafka.streams.KafkaStreams; */ public class KafkaStreamsBinderMetrics { - private KafkaStreams kafkaStreams; - private final MeterRegistry meterRegistry; private MeterBinder meterBinder; @@ -45,26 +48,41 @@ public class KafkaStreamsBinderMetrics { this.meterRegistry = meterRegistry; } - public void bindTo(MeterRegistry meterRegistry) { + public void bindTo(Set streamsBuilderFactoryBeans, MeterRegistry meterRegistry) { + if (this.meterBinder == null) { this.meterBinder = new MeterBinder() { @Override @SuppressWarnings("unchecked") public void bindTo(MeterRegistry registry) { - if (KafkaStreamsBinderMetrics.this.kafkaStreams != null) { - final Map metrics = KafkaStreamsBinderMetrics.this.kafkaStreams.metrics(); + if (streamsBuilderFactoryBeans != null) { + for (StreamsBuilderFactoryBean streamsBuilderFactoryBean : streamsBuilderFactoryBeans) { + KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); + final Map metrics = kafkaStreams.metrics(); - for (Map.Entry metric : metrics.entrySet()) { - final Gauge.Builder builder = - Gauge.builder(sanitize(metric.getKey().group() + "." + metric.getKey().name()), this, - toDoubleFunction(metric.getValue())); - final Map tags = metric.getKey().tags(); - for (Map.Entry tag : tags.entrySet()) { - builder.tag(tag.getKey(), tag.getValue()); + Set meterNames = new HashSet<>(); + + for (Map.Entry metric : metrics.entrySet()) { + final String sanitized = sanitize(metric.getKey().group() + "." + metric.getKey().name()); + final String applicationId = streamsBuilderFactoryBean.getStreamsConfiguration().getProperty(StreamsConfig.APPLICATION_ID_CONFIG); + + final String name = streamsBuilderFactoryBeans.size() > 1 ? applicationId + "." + sanitized : sanitized; + + final Gauge.Builder builder = + Gauge.builder(name, this, + toDoubleFunction(metric.getValue())); + final Map tags = metric.getKey().tags(); + for (Map.Entry tag : tags.entrySet()) { + builder.tag(tag.getKey(), tag.getValue()); + } + if (!meterNames.contains(name)) { + builder.description(metric.getKey().description()) + .register(meterRegistry); + meterNames.add(name); + } } - builder.description(metric.getKey().description()) - .register(meterRegistry); } + } } @@ -83,14 +101,13 @@ public class KafkaStreamsBinderMetrics { this.meterBinder.bindTo(this.meterRegistry); } - public void addMetrics(KafkaStreams kafkaStreams) { - synchronized (KafkaStreamsBinderMetrics.this) { - this.kafkaStreams = kafkaStreams; - this.bindTo(this.meterRegistry); - } - } - private static String sanitize(String value) { return value.replaceAll("-", "."); } + + public void addMetrics(Set streamsBuilderFactoryBeans) { + synchronized (KafkaStreamsBinderMetrics.this) { + this.bindTo(streamsBuilderFactoryBeans, this.meterRegistry); + } + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index e5b6f70de..fb2fb80ac 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -352,15 +352,16 @@ public class KafkaStreamsBinderSupportAutoConfiguration { } @Bean - public KafkaStreamsRegistry kafkaStreamsRegistry(@Nullable KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) { - return new KafkaStreamsRegistry(kafkaStreamsBinderMetrics); + public KafkaStreamsRegistry kafkaStreamsRegistry() { + return new KafkaStreamsRegistry(); } @Bean public StreamsBuilderFactoryManager streamsBuilderFactoryManager( KafkaStreamsBindingInformationCatalogue catalogue, - KafkaStreamsRegistry kafkaStreamsRegistry) { - return new StreamsBuilderFactoryManager(catalogue, kafkaStreamsRegistry); + KafkaStreamsRegistry kafkaStreamsRegistry, + @Nullable KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) { + return new StreamsBuilderFactoryManager(catalogue, kafkaStreamsRegistry, kafkaStreamsBinderMetrics); } @Bean @@ -393,8 +394,7 @@ public class KafkaStreamsBinderSupportAutoConfiguration { @Bean @ConditionalOnBean(MeterRegistry.class) @ConditionalOnMissingBean(KafkaStreamsBinderMetrics.class) - public KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics( - MeterRegistry meterRegistry) { + public KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics(MeterRegistry meterRegistry) { return new KafkaStreamsBinderMetrics(meterRegistry); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java index 79f65986d..efc5b7b17 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java @@ -33,13 +33,7 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean; */ class KafkaStreamsRegistry { - private final KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics; - - private Map streamsStreamsBuilderFactoryBeanMap = new HashMap<>(); - - KafkaStreamsRegistry(KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) { - this.kafkaStreamsBinderMetrics = kafkaStreamsBinderMetrics; - } + private Map streamsBuilderFactoryBeanMap = new HashMap<>(); private final Set kafkaStreams = new HashSet<>(); @@ -49,23 +43,12 @@ class KafkaStreamsRegistry { /** * Register the {@link KafkaStreams} object created in the application. - * @param kafkaStreams {@link KafkaStreams} object created in the application + * @param streamsBuilderFactoryBean {@link StreamsBuilderFactoryBean} */ - void registerKafkaStreams(KafkaStreams kafkaStreams) { - if (this.kafkaStreamsBinderMetrics != null) { - this.kafkaStreamsBinderMetrics.addMetrics(kafkaStreams); - } + void registerKafkaStreams(StreamsBuilderFactoryBean streamsBuilderFactoryBean) { + final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); this.kafkaStreams.add(kafkaStreams); - } - - /** - * Make an association between {@link KafkaStreams} and its corresponding {@link StreamsBuilderFactoryBean}. - * - * @param kafkaStreams {@link KafkaStreams} object - * @param streamsBuilderFactoryBean Associtated {@link StreamsBuilderFactoryBean} for the {@link KafkaStreams} - */ - void addToStreamBuilderFactoryBeanMap(KafkaStreams kafkaStreams, StreamsBuilderFactoryBean streamsBuilderFactoryBean) { - streamsStreamsBuilderFactoryBeanMap.put(kafkaStreams, streamsBuilderFactoryBean); + this.streamsBuilderFactoryBeanMap.put(kafkaStreams, streamsBuilderFactoryBean); } /** @@ -74,7 +57,7 @@ class KafkaStreamsRegistry { * @return Corresponding {@link StreamsBuilderFactoryBean}. */ StreamsBuilderFactoryBean streamBuilderFactoryBean(KafkaStreams kafkaStreams) { - return streamsStreamsBuilderFactoryBeanMap.get(kafkaStreams); + return this.streamsBuilderFactoryBeanMap.get(kafkaStreams); } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java index 92269b546..c5ca7fcfe 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java @@ -18,8 +18,6 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Set; -import org.apache.kafka.streams.KafkaStreams; - import org.springframework.context.SmartLifecycle; import org.springframework.kafka.KafkaException; import org.springframework.kafka.config.StreamsBuilderFactoryBean; @@ -42,14 +40,15 @@ class StreamsBuilderFactoryManager implements SmartLifecycle { private final KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue; private final KafkaStreamsRegistry kafkaStreamsRegistry; + private final KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics; private volatile boolean running; - StreamsBuilderFactoryManager( - KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, - KafkaStreamsRegistry kafkaStreamsRegistry) { + StreamsBuilderFactoryManager(KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, + KafkaStreamsRegistry kafkaStreamsRegistry, KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) { this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; this.kafkaStreamsRegistry = kafkaStreamsRegistry; + this.kafkaStreamsBinderMetrics = kafkaStreamsBinderMetrics; } @Override @@ -73,11 +72,9 @@ class StreamsBuilderFactoryManager implements SmartLifecycle { .getStreamsBuilderFactoryBeans(); for (StreamsBuilderFactoryBean streamsBuilderFactoryBean : streamsBuilderFactoryBeans) { streamsBuilderFactoryBean.start(); - final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); - this.kafkaStreamsRegistry.registerKafkaStreams( - kafkaStreams); - this.kafkaStreamsRegistry.addToStreamBuilderFactoryBeanMap(kafkaStreams, streamsBuilderFactoryBean); + this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); } + this.kafkaStreamsBinderMetrics.addMetrics(streamsBuilderFactoryBeans); this.running = true; } catch (Exception ex) { diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsInteractiveQueryIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsInteractiveQueryIntegrationTests.java index 1b34befe6..b1c2008b0 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsInteractiveQueryIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsInteractiveQueryIntegrationTests.java @@ -48,6 +48,7 @@ import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaSt import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; +import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -94,9 +95,10 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { @Test public void testStateStoreRetrievalRetry() { - KafkaStreams mock = Mockito.mock(KafkaStreams.class); - KafkaStreamsBinderMetrics mockMetrics = Mockito.mock(KafkaStreamsBinderMetrics.class); - KafkaStreamsRegistry kafkaStreamsRegistry = new KafkaStreamsRegistry(mockMetrics); + StreamsBuilderFactoryBean mock = Mockito.mock(StreamsBuilderFactoryBean.class); + KafkaStreams mockKafkaStreams = Mockito.mock(KafkaStreams.class); + Mockito.when(mock.getKafkaStreams()).thenReturn(mockKafkaStreams); + KafkaStreamsRegistry kafkaStreamsRegistry = new KafkaStreamsRegistry(); kafkaStreamsRegistry.registerKafkaStreams(mock); KafkaStreamsBinderConfigurationProperties binderConfigurationProperties = new KafkaStreamsBinderConfigurationProperties(new KafkaProperties()); @@ -112,7 +114,7 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { } - Mockito.verify(mock, times(3)).store("foo", storeType); + Mockito.verify(mockKafkaStreams, times(3)).store("foo", storeType); } @Test