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
This commit is contained in:
@@ -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/<metric-name>`).
|
||||
|
||||
@@ -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<StreamsBuilderFactoryBean> 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<MetricName, ? extends Metric> metrics = KafkaStreamsBinderMetrics.this.kafkaStreams.metrics();
|
||||
if (streamsBuilderFactoryBeans != null) {
|
||||
for (StreamsBuilderFactoryBean streamsBuilderFactoryBean : streamsBuilderFactoryBeans) {
|
||||
KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams();
|
||||
final Map<MetricName, ? extends Metric> metrics = kafkaStreams.metrics();
|
||||
|
||||
for (Map.Entry<MetricName, ? extends Metric> metric : metrics.entrySet()) {
|
||||
final Gauge.Builder<KafkaStreamsBinderMetrics> builder =
|
||||
Gauge.builder(sanitize(metric.getKey().group() + "." + metric.getKey().name()), this,
|
||||
toDoubleFunction(metric.getValue()));
|
||||
final Map<String, String> tags = metric.getKey().tags();
|
||||
for (Map.Entry<String, String> tag : tags.entrySet()) {
|
||||
builder.tag(tag.getKey(), tag.getValue());
|
||||
Set<String> meterNames = new HashSet<>();
|
||||
|
||||
for (Map.Entry<MetricName, ? extends Metric> 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<KafkaStreamsBinderMetrics> builder =
|
||||
Gauge.builder(name, this,
|
||||
toDoubleFunction(metric.getValue()));
|
||||
final Map<String, String> tags = metric.getKey().tags();
|
||||
for (Map.Entry<String, String> 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<StreamsBuilderFactoryBean> streamsBuilderFactoryBeans) {
|
||||
synchronized (KafkaStreamsBinderMetrics.this) {
|
||||
this.bindTo(streamsBuilderFactoryBeans, this.meterRegistry);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -33,13 +33,7 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean;
|
||||
*/
|
||||
class KafkaStreamsRegistry {
|
||||
|
||||
private final KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics;
|
||||
|
||||
private Map<KafkaStreams, StreamsBuilderFactoryBean> streamsStreamsBuilderFactoryBeanMap = new HashMap<>();
|
||||
|
||||
KafkaStreamsRegistry(KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) {
|
||||
this.kafkaStreamsBinderMetrics = kafkaStreamsBinderMetrics;
|
||||
}
|
||||
private Map<KafkaStreams, StreamsBuilderFactoryBean> streamsBuilderFactoryBeanMap = new HashMap<>();
|
||||
|
||||
private final Set<KafkaStreams> 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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user