GH-2485: Improvements in Kafka Binder Metrics
Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2485 These improvements are threading related Initial implementation Readd missing 'cache' for recent offset lags and general refactoring Add a test case for the schedule-only variant and externalize scheduling interval Consistency of configuration property naming and minor cleanup of existing code Review feedback: documentation, author addition and copyright adjustments Doc wording changes Fix checkstyle issues Adjust reference doc and improve new property names according to review Move documentation to the correct file Use JDK 8 compatible APIs and make the new test case public to satisfy junit
This commit is contained in:
@@ -1,5 +1,5 @@
|
|||||||
/*
|
/*
|
||||||
* Copyright 2015-2021 the original author or authors.
|
* Copyright 2015-2022 the original author or authors.
|
||||||
*
|
*
|
||||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
* you may not use this file except in compliance with the License.
|
* you may not use this file except in compliance with the License.
|
||||||
@@ -62,6 +62,7 @@ import org.springframework.util.StringUtils;
|
|||||||
* @author Aldo Sinanaj
|
* @author Aldo Sinanaj
|
||||||
* @author Lukasz Kaminski
|
* @author Lukasz Kaminski
|
||||||
* @author Chukwubuikem Ume-Ugwa
|
* @author Chukwubuikem Ume-Ugwa
|
||||||
|
* @author Nico Heller
|
||||||
*/
|
*/
|
||||||
@ConfigurationProperties(prefix = "spring.cloud.stream.kafka.binder")
|
@ConfigurationProperties(prefix = "spring.cloud.stream.kafka.binder")
|
||||||
public class KafkaBinderConfigurationProperties {
|
public class KafkaBinderConfigurationProperties {
|
||||||
@@ -72,6 +73,8 @@ public class KafkaBinderConfigurationProperties {
|
|||||||
|
|
||||||
private final Transaction transaction = new Transaction();
|
private final Transaction transaction = new Transaction();
|
||||||
|
|
||||||
|
private final Metrics metrics = new Metrics();
|
||||||
|
|
||||||
private final KafkaProperties kafkaProperties;
|
private final KafkaProperties kafkaProperties;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -150,6 +153,10 @@ public class KafkaBinderConfigurationProperties {
|
|||||||
return this.transaction;
|
return this.transaction;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public Metrics getMetrics() {
|
||||||
|
return metrics;
|
||||||
|
}
|
||||||
|
|
||||||
public String getKafkaConnectionString() {
|
public String getKafkaConnectionString() {
|
||||||
// We need to do a check on certificate file locations to see if they are given as classpath resources.
|
// We need to do a check on certificate file locations to see if they are given as classpath resources.
|
||||||
// If that is the case, then we will move them to a file system location and use those as the certificate locations.
|
// If that is the case, then we will move them to a file system location and use those as the certificate locations.
|
||||||
@@ -666,4 +673,34 @@ public class KafkaBinderConfigurationProperties {
|
|||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static class Metrics {
|
||||||
|
/**
|
||||||
|
* When set to true, the offset lag metric of each consumer topic is computed whenever the metric is queried (default: true).
|
||||||
|
* When set to false only the periodically calculated offset lag is used.
|
||||||
|
* See {@link #offsetLagMetricsInterval} for more information.
|
||||||
|
*/
|
||||||
|
private boolean defaultOffsetLagMetricsEnabled = true;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The interval in which the offset lag for each consumer topic is computed (default: 60 seconds).
|
||||||
|
* This value is used whenever {@link #defaultOffsetLagMetricsEnabled} is disabled or its computation is taking too long.
|
||||||
|
*/
|
||||||
|
private Duration offsetLagMetricsInterval = Duration.ofSeconds(60);
|
||||||
|
|
||||||
|
public boolean isDefaultOffsetLagMetricsEnabled() {
|
||||||
|
return defaultOffsetLagMetricsEnabled;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setDefaultOffsetLagMetricsEnabled(boolean defaultOffsetLagMetricsEnabled) {
|
||||||
|
this.defaultOffsetLagMetricsEnabled = defaultOffsetLagMetricsEnabled;
|
||||||
|
}
|
||||||
|
|
||||||
|
public Duration getOffsetLagMetricsInterval() {
|
||||||
|
return offsetLagMetricsInterval;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setOffsetLagMetricsInterval(Duration offsetLagMetricsInterval) {
|
||||||
|
this.offsetLagMetricsInterval = offsetLagMetricsInterval;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ import java.util.concurrent.ScheduledExecutorService;
|
|||||||
import java.util.concurrent.ScheduledThreadPoolExecutor;
|
import java.util.concurrent.ScheduledThreadPoolExecutor;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
import java.util.concurrent.TimeoutException;
|
import java.util.concurrent.TimeoutException;
|
||||||
|
import java.util.function.ToDoubleFunction;
|
||||||
|
|
||||||
import io.micrometer.core.instrument.Gauge;
|
import io.micrometer.core.instrument.Gauge;
|
||||||
import io.micrometer.core.instrument.MeterRegistry;
|
import io.micrometer.core.instrument.MeterRegistry;
|
||||||
@@ -64,14 +65,13 @@ import org.springframework.util.ObjectUtils;
|
|||||||
* @author Gary Russell
|
* @author Gary Russell
|
||||||
* @author Lars Bilger
|
* @author Lars Bilger
|
||||||
* @author Tomek Szmytka
|
* @author Tomek Szmytka
|
||||||
|
* @author Nico Heller
|
||||||
*/
|
*/
|
||||||
public class KafkaBinderMetrics
|
public class KafkaBinderMetrics
|
||||||
implements MeterBinder, ApplicationListener<BindingCreatedEvent>, AutoCloseable {
|
implements MeterBinder, ApplicationListener<BindingCreatedEvent>, AutoCloseable {
|
||||||
|
|
||||||
private static final int DEFAULT_TIMEOUT = 5;
|
private static final int DEFAULT_TIMEOUT = 5;
|
||||||
|
|
||||||
private static final int DELAY_BETWEEN_TASK_EXECUTION = 60;
|
|
||||||
|
|
||||||
private static final Log LOG = LogFactory.getLog(KafkaBinderMetrics.class);
|
private static final Log LOG = LogFactory.getLog(KafkaBinderMetrics.class);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -91,9 +91,9 @@ public class KafkaBinderMetrics
|
|||||||
|
|
||||||
private int timeout = DEFAULT_TIMEOUT;
|
private int timeout = DEFAULT_TIMEOUT;
|
||||||
|
|
||||||
ScheduledExecutorService scheduler;
|
private final Map<String, Long> lastUnconsumedMessagesValues = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
Map<String, Long> unconsumedMessages = new ConcurrentHashMap<>();
|
ScheduledExecutorService scheduler;
|
||||||
|
|
||||||
public KafkaBinderMetrics(KafkaMessageChannelBinder binder,
|
public KafkaBinderMetrics(KafkaMessageChannelBinder binder,
|
||||||
KafkaBinderConfigurationProperties binderConfigurationProperties,
|
KafkaBinderConfigurationProperties binderConfigurationProperties,
|
||||||
@@ -141,58 +141,58 @@ public class KafkaBinderMetrics
|
|||||||
String topic = topicInfo.getKey();
|
String topic = topicInfo.getKey();
|
||||||
String group = topicInfo.getValue().getConsumerGroup();
|
String group = topicInfo.getValue().getConsumerGroup();
|
||||||
|
|
||||||
final Gauge register = Gauge.builder(OFFSET_LAG_METRIC_NAME, this,
|
ToDoubleFunction<KafkaBinderMetrics> offsetComputation = computeOffsetComputationFunction(topic, group);
|
||||||
(o) -> computeAndGetUnconsumedMessages(topic, group)).tag("group", group)
|
final Gauge register = Gauge.builder(OFFSET_LAG_METRIC_NAME, this, offsetComputation)
|
||||||
.tag("topic", topic)
|
.tag("group", group)
|
||||||
.description("Unconsumed messages for a particular group and topic")
|
.tag("topic", topic)
|
||||||
.register(registry);
|
.description("Unconsumed messages for a particular group and topic")
|
||||||
|
.register(registry);
|
||||||
|
|
||||||
if (!(register instanceof NoopGauge)) {
|
if (!(register instanceof NoopGauge)) {
|
||||||
//Schedule a task to compute the unconsumed messages for this group/topic every minute.
|
lastUnconsumedMessagesValues.put(topic + "-" + group, 0L);
|
||||||
this.scheduler.scheduleWithFixedDelay(computeUnconsumedMessagesRunnable(topic, group, this.metadataConsumers),
|
this.scheduler.scheduleWithFixedDelay(
|
||||||
10, DELAY_BETWEEN_TASK_EXECUTION, TimeUnit.SECONDS);
|
() -> computeUnconsumedMessages(topic, group),
|
||||||
|
1,
|
||||||
|
binderConfigurationProperties.getMetrics().getOffsetLagMetricsInterval().getSeconds(),
|
||||||
|
TimeUnit.SECONDS
|
||||||
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private Runnable computeUnconsumedMessagesRunnable(String topic, String group, Map<String, Consumer<?, ?>> metadataConsumers) {
|
private ToDoubleFunction<KafkaBinderMetrics> computeOffsetComputationFunction(String topic, String group) {
|
||||||
return () -> {
|
if (this.binderConfigurationProperties.getMetrics().isDefaultOffsetLagMetricsEnabled()) {
|
||||||
try {
|
return (o) -> computeAndGetUnconsumedMessagesWithTimeout(topic, group);
|
||||||
long lag = findTotalTopicGroupLag(topic, group, this.metadataConsumers);
|
}
|
||||||
this.unconsumedMessages.put(topic + "-" + group, lag);
|
else {
|
||||||
}
|
return (o) -> lastUnconsumedMessagesValues.get(topic + "-" + group);
|
||||||
catch (Exception ex) {
|
}
|
||||||
LOG.debug("Cannot generate metric for topic: " + topic, ex);
|
|
||||||
}
|
|
||||||
};
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private long computeAndGetUnconsumedMessages(String topic, String group) {
|
private long computeAndGetUnconsumedMessagesWithTimeout(String topic, String group) {
|
||||||
ExecutorService exec = Executors.newCachedThreadPool();
|
Future<Long> future = scheduler.submit(() -> computeUnconsumedMessages(topic, group));
|
||||||
Future<Long> future = exec.submit(() -> {
|
|
||||||
|
|
||||||
long lag = 0;
|
|
||||||
try {
|
|
||||||
lag = findTotalTopicGroupLag(topic, group, this.metadataConsumers);
|
|
||||||
}
|
|
||||||
catch (Exception ex) {
|
|
||||||
LOG.debug("Cannot generate metric for topic: " + topic, ex);
|
|
||||||
}
|
|
||||||
return lag;
|
|
||||||
});
|
|
||||||
try {
|
try {
|
||||||
return future.get(this.timeout, TimeUnit.SECONDS);
|
return future.get(this.timeout, TimeUnit.SECONDS);
|
||||||
}
|
}
|
||||||
catch (InterruptedException ex) {
|
catch (InterruptedException ex) {
|
||||||
Thread.currentThread().interrupt();
|
Thread.currentThread().interrupt();
|
||||||
return this.unconsumedMessages.getOrDefault(topic + "-" + group, 0L);
|
return lastUnconsumedMessagesValues.get(topic + "-" + group);
|
||||||
}
|
}
|
||||||
catch (ExecutionException | TimeoutException ex) {
|
catch (ExecutionException | TimeoutException ex) {
|
||||||
return this.unconsumedMessages.getOrDefault(topic + "-" + group, 0L);
|
return lastUnconsumedMessagesValues.get(topic + "-" + group);
|
||||||
}
|
}
|
||||||
finally {
|
}
|
||||||
exec.shutdownNow();
|
|
||||||
|
private long computeUnconsumedMessages(String topic, String group) {
|
||||||
|
long lag = 0;
|
||||||
|
try {
|
||||||
|
lag = findTotalTopicGroupLag(topic, group, this.metadataConsumers);
|
||||||
|
this.lastUnconsumedMessagesValues.put(topic + "-" + group, lag);
|
||||||
}
|
}
|
||||||
|
catch (Exception ex) {
|
||||||
|
LOG.debug("Cannot generate metric for topic: " + topic, ex);
|
||||||
|
}
|
||||||
|
return lag;
|
||||||
}
|
}
|
||||||
|
|
||||||
private long findTotalTopicGroupLag(String topic, String group, Map<String, Consumer<?, ?>> metadataConsumers) {
|
private long findTotalTopicGroupLag(String topic, String group, Map<String, Consumer<?, ?>> metadataConsumers) {
|
||||||
|
|||||||
@@ -16,6 +16,7 @@
|
|||||||
|
|
||||||
package org.springframework.cloud.stream.binder.kafka;
|
package org.springframework.cloud.stream.binder.kafka;
|
||||||
|
|
||||||
|
import java.time.Duration;
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.HashMap;
|
import java.util.HashMap;
|
||||||
@@ -32,8 +33,10 @@ import org.apache.kafka.common.KafkaException;
|
|||||||
import org.apache.kafka.common.Node;
|
import org.apache.kafka.common.Node;
|
||||||
import org.apache.kafka.common.PartitionInfo;
|
import org.apache.kafka.common.PartitionInfo;
|
||||||
import org.apache.kafka.common.TopicPartition;
|
import org.apache.kafka.common.TopicPartition;
|
||||||
|
import org.awaitility.Awaitility;
|
||||||
import org.junit.Before;
|
import org.junit.Before;
|
||||||
import org.junit.Test;
|
import org.junit.Test;
|
||||||
|
import org.mockito.Answers;
|
||||||
import org.mockito.ArgumentMatchers;
|
import org.mockito.ArgumentMatchers;
|
||||||
import org.mockito.Mock;
|
import org.mockito.Mock;
|
||||||
import org.mockito.Mockito;
|
import org.mockito.Mockito;
|
||||||
@@ -52,6 +55,7 @@ import static org.mockito.Mockito.mock;
|
|||||||
* @author Soby Chacko
|
* @author Soby Chacko
|
||||||
* @author Lars Bilger
|
* @author Lars Bilger
|
||||||
* @author Tomek Szmytka
|
* @author Tomek Szmytka
|
||||||
|
* @author Nico Heller
|
||||||
*/
|
*/
|
||||||
public class KafkaBinderMetricsTest {
|
public class KafkaBinderMetricsTest {
|
||||||
|
|
||||||
@@ -72,7 +76,7 @@ public class KafkaBinderMetricsTest {
|
|||||||
|
|
||||||
private Map<String, TopicInformation> topicsInUse = new HashMap<>();
|
private Map<String, TopicInformation> topicsInUse = new HashMap<>();
|
||||||
|
|
||||||
@Mock
|
@Mock(answer = Answers.RETURNS_DEEP_STUBS)
|
||||||
private KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties;
|
private KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties;
|
||||||
|
|
||||||
@Before
|
@Before
|
||||||
@@ -82,6 +86,10 @@ public class KafkaBinderMetricsTest {
|
|||||||
.createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()))
|
.createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()))
|
||||||
.willReturn(consumer);
|
.willReturn(consumer);
|
||||||
org.mockito.BDDMockito.given(binder.getTopicsInUse()).willReturn(topicsInUse);
|
org.mockito.BDDMockito.given(binder.getTopicsInUse()).willReturn(topicsInUse);
|
||||||
|
org.mockito.BDDMockito.given(kafkaBinderConfigurationProperties.getMetrics().isDefaultOffsetLagMetricsEnabled())
|
||||||
|
.willReturn(true);
|
||||||
|
org.mockito.BDDMockito.given(kafkaBinderConfigurationProperties.getMetrics().getOffsetLagMetricsInterval())
|
||||||
|
.willReturn(Duration.ofSeconds(60));
|
||||||
metrics = new KafkaBinderMetrics(binder, kafkaBinderConfigurationProperties,
|
metrics = new KafkaBinderMetrics(binder, kafkaBinderConfigurationProperties,
|
||||||
consumerFactory, null
|
consumerFactory, null
|
||||||
);
|
);
|
||||||
@@ -113,6 +121,39 @@ public class KafkaBinderMetricsTest {
|
|||||||
.isEqualTo(500.0);
|
.isEqualTo(500.0);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void shouldFallbackToScheduledOffsetLagComputationWhenRealtimeOffsetLagIsDisabled() {
|
||||||
|
final Map<TopicPartition, OffsetAndMetadata> committed = new HashMap<>();
|
||||||
|
TopicPartition topicPartition = new TopicPartition(TEST_TOPIC, 0);
|
||||||
|
committed.put(topicPartition, new OffsetAndMetadata(500));
|
||||||
|
org.mockito.BDDMockito
|
||||||
|
.given(consumer.committed(ArgumentMatchers.anySet()))
|
||||||
|
.willReturn(committed);
|
||||||
|
List<PartitionInfo> partitions = partitions(new Node(0, null, 0));
|
||||||
|
topicsInUse.put(
|
||||||
|
TEST_TOPIC,
|
||||||
|
new TopicInformation("group1-metrics", partitions, false)
|
||||||
|
);
|
||||||
|
org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC))
|
||||||
|
.willReturn(partitions);
|
||||||
|
org.mockito.BDDMockito.given(kafkaBinderConfigurationProperties.getMetrics().isDefaultOffsetLagMetricsEnabled())
|
||||||
|
.willReturn(false);
|
||||||
|
metrics.bindTo(meterRegistry);
|
||||||
|
assertThat(meterRegistry.getMeters()).hasSize(1);
|
||||||
|
assertThat(meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME)
|
||||||
|
.tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge().value())
|
||||||
|
.isEqualTo(0);
|
||||||
|
|
||||||
|
org.mockito.BDDMockito.given(kafkaBinderConfigurationProperties.getMetrics().getOffsetLagMetricsInterval())
|
||||||
|
.willReturn(Duration.ofSeconds(1));
|
||||||
|
metrics.bindTo(meterRegistry);
|
||||||
|
Awaitility.waitAtMost(Duration.ofSeconds(5)).untilAsserted(() -> {
|
||||||
|
assertThat(meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME)
|
||||||
|
.tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge().value())
|
||||||
|
.isEqualTo(500);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void shouldNotContainAnyMetricsWhenUsingNoopGauge() {
|
public void shouldNotContainAnyMetricsWhenUsingNoopGauge() {
|
||||||
// Adding NoopGauge for the offset metric.
|
// Adding NoopGauge for the offset metric.
|
||||||
|
|||||||
@@ -51,6 +51,7 @@ This section contains the configuration options used by the Apache Kafka binder.
|
|||||||
|
|
||||||
For common configuration options and properties pertaining to the binder, see the https://cloud.spring.io/spring-cloud-static/spring-cloud-stream/current/reference/html/spring-cloud-stream.html#binding-properties[binding properties] in core documentation.
|
For common configuration options and properties pertaining to the binder, see the https://cloud.spring.io/spring-cloud-static/spring-cloud-stream/current/reference/html/spring-cloud-stream.html#binding-properties[binding properties] in core documentation.
|
||||||
|
|
||||||
|
[[kafka-binder-properties]]
|
||||||
==== Kafka Binder Properties
|
==== Kafka Binder Properties
|
||||||
|
|
||||||
spring.cloud.stream.kafka.binder.brokers::
|
spring.cloud.stream.kafka.binder.brokers::
|
||||||
@@ -160,6 +161,19 @@ This is also true, if this value is present, but the directory cannot be found o
|
|||||||
+
|
+
|
||||||
Default: none.
|
Default: none.
|
||||||
|
|
||||||
|
spring.cloud.stream.kafka.binder.metrics.defaultOffsetLagMetricsEnabled::
|
||||||
|
When set to true, the offset lag metric of each consumer topic is computed whenever the metric is accessed.
|
||||||
|
When set to false only the periodically calculated offset lag is used.
|
||||||
|
+
|
||||||
|
Default: true
|
||||||
|
|
||||||
|
spring.cloud.stream.kafka.binder.metrics.offsetLagMetricsInterval::
|
||||||
|
The interval in which the offset lag for each consumer topic is computed.
|
||||||
|
This value is used whenever `metrics.defaultOffsetLagMetricsEnabled` is disabled or its
|
||||||
|
computation is taking too long.
|
||||||
|
+
|
||||||
|
Default: 60 seconds
|
||||||
|
|
||||||
[[kafka-consumer-properties]]
|
[[kafka-consumer-properties]]
|
||||||
==== Kafka Consumer Properties
|
==== Kafka Consumer Properties
|
||||||
|
|
||||||
@@ -744,6 +758,9 @@ The binder creates the `KafkaBinderMetrics` bean if Micrometer is on the classpa
|
|||||||
The metric contains the consumer group information, topic and the actual lag in committed offset from the latest offset on the topic.
|
The metric contains the consumer group information, topic and the actual lag in committed offset from the latest offset on the topic.
|
||||||
This metric is particularly useful for providing auto-scaling feedback to a PaaS platform.
|
This metric is particularly useful for providing auto-scaling feedback to a PaaS platform.
|
||||||
|
|
||||||
|
The metric collection behaviour can be configured by setting properties in the `spring.cloud.stream.kafka.binder.metrics` namespace,
|
||||||
|
refer to the <<kafka-binder-properties, kafka binder properties section>> for more information.
|
||||||
|
|
||||||
You can exclude `KafkaBinderMetrics` from creating the necessary infrastructure like consumers and then reporting the metrics by providing the following component in the application.
|
You can exclude `KafkaBinderMetrics` from creating the necessary infrastructure like consumers and then reporting the metrics by providing the following component in the application.
|
||||||
|
|
||||||
```
|
```
|
||||||
|
|||||||
Reference in New Issue
Block a user