From 1945f29d5518e3c4a9950ba82135420dfb61e808 Mon Sep 17 00:00:00 2001 From: Vladimir Tsanev Date: Fri, 25 Aug 2017 16:10:20 +0300 Subject: [PATCH] GH-401: Add API to access consumer metrics Fix spring-projects/spring-kafka#401 **Cherry-pick to 1.3.x** Polish javadoc and remove dependency to jdk 8. Return always return unmodifiable collections from metrics() * Polishing JavaDocs * Add author * Mention the new feature in the Docs --- .../ConcurrentMessageListenerContainer.java | 17 ++++++++++++++++- .../listener/KafkaMessageListenerContainer.java | 17 +++++++++++++++++ .../listener/MessageListenerContainer.java | 16 +++++++++++++++- ...ConcurrentMessageListenerContainerTests.java | 2 ++ src/reference/asciidoc/kafka.adoc | 4 ++++ 5 files changed, 54 insertions(+), 2 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java index 3a3e8df6..cfdc7514 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java @@ -1,5 +1,5 @@ /* - * Copyright 2015-2016 the original author or authors. + * Copyright 2015-2017 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -19,9 +19,14 @@ package org.springframework.kafka.listener; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; +import org.apache.kafka.common.Metric; +import org.apache.kafka.common.MetricName; + import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.listener.config.ContainerProperties; import org.springframework.kafka.support.TopicPartitionInitialOffset; @@ -42,6 +47,7 @@ import org.springframework.util.Assert; * @author Murali Reddy * @author Jerome Mirc * @author Artem Bilan + * @author Vladimir Tsanev */ public class ConcurrentMessageListenerContainer extends AbstractMessageListenerContainer { @@ -89,6 +95,15 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis return Collections.unmodifiableList(this.containers); } + @Override + public Map> metrics() { + Map> metrics = new HashMap<>(); + for (KafkaMessageListenerContainer container : this.containers) { + metrics.putAll(container.metrics()); + } + return Collections.unmodifiableMap(metrics); + } + /* * Under lifecycle lock. */ diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 31e554fd..0cf36766 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -41,6 +41,8 @@ import org.apache.kafka.clients.consumer.NoOffsetForPartitionException; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.clients.consumer.OffsetCommitCallback; import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.common.Metric; +import org.apache.kafka.common.MetricName; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.WakeupException; @@ -80,6 +82,7 @@ import org.springframework.util.concurrent.ListenableFutureCallback; * @author Martin Dam * @author Artem Bilan * @author Loic Talhouarne + * @author Vladimir Tsanev */ public class KafkaMessageListenerContainer extends AbstractMessageListenerContainer { @@ -153,6 +156,20 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } } + @Override + public Map> metrics() { + ListenerConsumer listenerConsumer = this.listenerConsumer; + if (listenerConsumer != null) { + Map metrics = listenerConsumer.consumer.metrics(); + Iterator metricIterator = metrics.keySet().iterator(); + if (metricIterator.hasNext()) { + String clientId = metricIterator.next().tags().get("client-id"); + return Collections.singletonMap(clientId, metrics); + } + } + return Collections.emptyMap(); + } + @Override protected void doStart() { if (isRunning()) { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java index a9db3d0c..fadd267f 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2017 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,6 +16,11 @@ package org.springframework.kafka.listener; +import java.util.Map; + +import org.apache.kafka.common.Metric; +import org.apache.kafka.common.MetricName; + import org.springframework.context.SmartLifecycle; /** @@ -24,6 +29,7 @@ import org.springframework.context.SmartLifecycle; * * @author Stephane Nicoll * @author Gary Russell + * @author Vladimir Tsanev */ public interface MessageListenerContainer extends SmartLifecycle { @@ -34,4 +40,12 @@ public interface MessageListenerContainer extends SmartLifecycle { */ void setupMessageListener(Object messageListener); + /** + * Return metrics kept by this container's consumer(s), grouped by {@code client-id}. + * @return the consumer(s) metrics grouped by {@code client-id} + * @since 1.3 + * @see org.apache.kafka.clients.consumer.Consumer#metrics() + */ + Map> metrics(); + } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java index 4a1e491d..42e23a5f 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java @@ -67,6 +67,7 @@ import org.springframework.kafka.test.utils.KafkaTestUtils; * @author Jerome Mirc * @author Marius Bogoevici * @author Artem Yakshin + * @author Vladimir Tsanev */ public class ConcurrentMessageListenerContainerTests { @@ -140,6 +141,7 @@ public class ConcurrentMessageListenerContainerTests { assertThat(KafkaTestUtils.getPropertyValue(containers.get(i), "listenerConsumer.acks", Collection.class) .size()).isEqualTo(0); } + assertThat(container.metrics()).isNotNull(); container.stop(); this.logger.info("Stop auto"); } diff --git a/src/reference/asciidoc/kafka.adoc b/src/reference/asciidoc/kafka.adoc index 0446b88f..85437fa4 100644 --- a/src/reference/asciidoc/kafka.adoc +++ b/src/reference/asciidoc/kafka.adoc @@ -505,6 +505,10 @@ each container will get one partition. NOTE: The `client.id` property (if set) will be appended with `-n` where `n` is the consumer instance according to the concurrency. This is required to provide unique names for MBeans when JMX is enabled. +Starting with _version 1.3_, the `MessageListenerContainer` provides an access to the metrics of the underlying `KafkaConsumer`. +In case of `ConcurrentMessageListenerContainer` the `metrics()` method returns the metrics for all the target `KafkaMessageListenerContainer` instances. +The metrics are grouped into the `Map` by the `client-id` provided for the underlying `KafkaConsumer`. + [[committing-offsets]] ====== Committing Offsets