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