From aa14aff0f15efb55886401cb9f2c95dd6e36e562 Mon Sep 17 00:00:00 2001 From: "Colin P. Mccabe" Date: Tue, 17 Dec 2019 16:42:48 -0800 Subject: [PATCH 1/2] KAFKA-9306: Clean up KafkaConsumerMetrics on consumer close --- .../kafka/clients/consumer/KafkaConsumer.java | 1 + .../internals/KafkaConsumerMetrics.java | 23 ++++++++++++------- 2 files changed, 16 insertions(+), 8 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java index b5a1047f0efd3..97fc44ffa21e3 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java @@ -2301,6 +2301,7 @@ private void close(long timeoutMs, boolean swallowException) { } Utils.closeQuietly(fetcher, "fetcher", firstException); Utils.closeQuietly(interceptors, "consumer interceptors", firstException); + Utils.closeQuietly(kafkaConsumerMetrics, "kafka consumer metrics", firstException); Utils.closeQuietly(metrics, "consumer metrics", firstException); Utils.closeQuietly(client, "consumer network client", firstException); Utils.closeQuietly(keyDeserializer, "consumer key deserializer", firstException); diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetrics.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetrics.java index ae61ff1b76cba..aab313320d4bd 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetrics.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/KafkaConsumerMetrics.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients.consumer.internals; +import org.apache.kafka.common.MetricName; import org.apache.kafka.common.metrics.Measurable; import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.metrics.Sensor; @@ -24,11 +25,11 @@ import java.util.concurrent.TimeUnit; -public class KafkaConsumerMetrics { +public class KafkaConsumerMetrics implements AutoCloseable { private final Metrics metrics; - - private Sensor timeBetweenPollSensor; - private Sensor pollIdleSensor; + private final MetricName lastPollMetricName; + private final Sensor timeBetweenPollSensor; + private final Sensor pollIdleSensor; private long lastPollMs; private long pollStartMs; private long timeSinceLastPollMs; @@ -44,10 +45,9 @@ public KafkaConsumerMetrics(Metrics metrics, String metricGrpPrefix) { else return TimeUnit.SECONDS.convert(now - lastPollMs, TimeUnit.MILLISECONDS); }; - metrics.addMetric(metrics.metricName("last-poll-seconds-ago", - metricGroupName, - "The number of seconds since the last poll() invocation."), - lastPoll); + this.lastPollMetricName = metrics.metricName("last-poll-seconds-ago", + metricGroupName, "The number of seconds since the last poll() invocation."); + metrics.addMetric(lastPollMetricName, lastPoll); this.timeBetweenPollSensor = metrics.sensor("time-between-poll"); this.timeBetweenPollSensor.add(metrics.metricName("time-between-poll-avg", @@ -78,4 +78,11 @@ public void recordPollEnd(long pollEndMs) { double pollIdleRatio = pollTimeMs * 1.0 / (pollTimeMs + timeSinceLastPollMs); this.pollIdleSensor.record(pollIdleRatio); } + + @Override + public void close() { + metrics.removeMetric(lastPollMetricName); + metrics.removeSensor(timeBetweenPollSensor.name()); + metrics.removeSensor(pollIdleSensor.name()); + } } From c4b236d66195cdc83038a9f5bde3539fcdbf750b Mon Sep 17 00:00:00 2001 From: "Colin P. Mccabe" Date: Wed, 18 Dec 2019 14:44:16 -0800 Subject: [PATCH 2/2] Add unit test --- .../clients/consumer/KafkaConsumerTest.java | 24 +++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java index 7455ee1281259..ccfe2a0e57483 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java @@ -2258,4 +2258,28 @@ public void testPollIdleRatio() { // Avg of three data points assertEquals((1.0d + 0.0d + 0.5d) / 3, consumer.metrics().get(pollIdleRatio).metricValue()); } + + private static boolean consumerMetricPresent(KafkaConsumer consumer, String name) { + MetricName metricName = new MetricName(name, "consumer-metrics", "", Collections.emptyMap()); + return consumer.metrics.metrics().containsKey(metricName); + } + + @Test + public void testClosingConsumerUnregistersConsumerMetrics() { + Time time = new MockTime(); + SubscriptionState subscription = new SubscriptionState(new LogContext(), OffsetResetStrategy.EARLIEST); + ConsumerMetadata metadata = createMetadata(subscription); + MockClient client = new MockClient(time, metadata); + initMetadata(client, Collections.singletonMap(topic, 1)); + KafkaConsumer consumer = newConsumer(time, client, subscription, metadata, + new RoundRobinAssignor(), true, groupInstanceId); + consumer.subscribe(singletonList(topic)); + assertTrue(consumerMetricPresent(consumer, "last-poll-seconds-ago")); + assertTrue(consumerMetricPresent(consumer, "time-between-poll-avg")); + assertTrue(consumerMetricPresent(consumer, "time-between-poll-max")); + consumer.close(); + assertFalse(consumerMetricPresent(consumer, "last-poll-seconds-ago")); + assertFalse(consumerMetricPresent(consumer, "time-between-poll-avg")); + assertFalse(consumerMetricPresent(consumer, "time-between-poll-max")); + } }