diff --git a/checkstyle/import-control-core.xml b/checkstyle/import-control-core.xml index 4898c0b50a73b..6e5042fd35d8a 100644 --- a/checkstyle/import-control-core.xml +++ b/checkstyle/import-control-core.xml @@ -38,6 +38,11 @@ + + + + diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 281a8c66851c2..4fd3cf30b3504 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -448,7 +448,9 @@ static KafkaAdminClient createInternal(AdminClientConfig config, TimeoutProcesso .timeWindow(config.getLong(AdminClientConfig.METRICS_SAMPLE_WINDOW_MS_CONFIG), TimeUnit.MILLISECONDS) .recordLevel(Sensor.RecordingLevel.forName(config.getString(AdminClientConfig.METRICS_RECORDING_LEVEL_CONFIG))) .tags(metricTags); - reporters.add(new JmxReporter(JMX_PREFIX)); + JmxReporter jmxReporter = new JmxReporter(JMX_PREFIX); + jmxReporter.configure(config.originals()); + reporters.add(jmxReporter); metrics = new Metrics(metricConfig, reporters, time); String metricGrpPrefix = "admin-client"; channelBuilder = ClientUtils.createChannelBuilder(config, time, logContext); 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 20da83f36f367..f140d6a8806a1 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 @@ -883,7 +883,9 @@ private static Metrics buildMetrics(ConsumerConfig config, Time time, String cli .tags(metricsTags); List reporters = config.getConfiguredInstances(ConsumerConfig.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class, Collections.singletonMap(ConsumerConfig.CLIENT_ID_CONFIG, clientId)); - reporters.add(new JmxReporter(JMX_PREFIX)); + JmxReporter jmxReporter = new JmxReporter(JMX_PREFIX); + jmxReporter.configure(config.originals()); + reporters.add(jmxReporter); return new Metrics(metricConfig, reporters, time); } diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java index 36af5030ec600..ac659d328bb41 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java @@ -351,7 +351,9 @@ public KafkaProducer(Properties properties, Serializer keySerializer, Seriali List reporters = config.getConfiguredInstances(ProducerConfig.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class, Collections.singletonMap(ProducerConfig.CLIENT_ID_CONFIG, clientId)); - reporters.add(new JmxReporter(JMX_PREFIX)); + JmxReporter jmxReporter = new JmxReporter(JMX_PREFIX); + jmxReporter.configure(userProvidedConfigs); + reporters.add(jmxReporter); this.metrics = new Metrics(metricConfig, reporters, time); this.partitioner = config.getConfiguredInstance(ProducerConfig.PARTITIONER_CLASS_CONFIG, Partitioner.class); long retryBackoffMs = config.getLong(ProducerConfig.RETRY_BACKOFF_MS_CONFIG); diff --git a/clients/src/main/java/org/apache/kafka/common/metrics/JmxReporter.java b/clients/src/main/java/org/apache/kafka/common/metrics/JmxReporter.java index 3129564faed24..c7be865fcd82d 100644 --- a/clients/src/main/java/org/apache/kafka/common/metrics/JmxReporter.java +++ b/clients/src/main/java/org/apache/kafka/common/metrics/JmxReporter.java @@ -18,7 +18,9 @@ import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.MetricName; +import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.utils.Sanitizer; +import org.apache.kafka.common.utils.Utils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -36,16 +38,32 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.function.Predicate; +import java.util.regex.Pattern; +import java.util.regex.PatternSyntaxException; /** * Register metrics in JMX as dynamic mbeans based on the metric names */ public class JmxReporter implements MetricsReporter { + public static final String METRICS_CONFIG_PREFIX = "metrics.jmx."; + + public static final String BLACKLIST_CONFIG = METRICS_CONFIG_PREFIX + "blacklist"; + public static final String WHITELIST_CONFIG = METRICS_CONFIG_PREFIX + "whitelist"; + + public static final Set RECONFIGURABLE_CONFIGS = Utils.mkSet(WHITELIST_CONFIG, + BLACKLIST_CONFIG); + + public static final String DEFAULT_WHITELIST = ".*"; + public static final String DEFAULT_BLACKLIST = ""; + private static final Logger log = LoggerFactory.getLogger(JmxReporter.class); private static final Object LOCK = new Object(); private String prefix; private final Map mbeans = new HashMap<>(); + private Predicate mbeanPredicate = s -> true; public JmxReporter() { this(""); @@ -59,15 +77,46 @@ public JmxReporter(String prefix) { } @Override - public void configure(Map configs) {} + public void configure(Map configs) { + reconfigure(configs); + } + + @Override + public Set reconfigurableConfigs() { + return RECONFIGURABLE_CONFIGS; + } + + @Override + public void validateReconfiguration(Map configs) throws ConfigException { + compilePredicate(configs); + } + + @Override + public void reconfigure(Map configs) { + synchronized (LOCK) { + this.mbeanPredicate = JmxReporter.compilePredicate(configs); + + mbeans.forEach((name, mbean) -> { + if (mbeanPredicate.test(name)) { + reregister(mbean); + } else { + unregister(mbean); + } + }); + } + } @Override public void init(List metrics) { synchronized (LOCK) { for (KafkaMetric metric : metrics) addAttribute(metric); - for (KafkaMbean mbean : mbeans.values()) - reregister(mbean); + + mbeans.forEach((name, mbean) -> { + if (mbeanPredicate.test(name)) { + reregister(mbean); + } + }); } } @@ -78,8 +127,10 @@ public boolean containsMbean(String mbeanName) { @Override public void metricChange(KafkaMetric metric) { synchronized (LOCK) { - KafkaMbean mbean = addAttribute(metric); - reregister(mbean); + String mbeanName = addAttribute(metric); + if (mbeanName != null && mbeanPredicate.test(mbeanName)) { + reregister(mbeans.get(mbeanName)); + } } } @@ -93,7 +144,7 @@ public void metricRemoval(KafkaMetric metric) { if (mbean.metrics.isEmpty()) { unregister(mbean); mbeans.remove(mBeanName); - } else + } else if (mbeanPredicate.test(mBeanName)) reregister(mbean); } } @@ -107,7 +158,7 @@ private KafkaMbean removeAttribute(KafkaMetric metric, String mBeanName) { return mbean; } - private KafkaMbean addAttribute(KafkaMetric metric) { + private String addAttribute(KafkaMetric metric) { try { MetricName metricName = metric.metricName(); String mBeanName = getMBeanName(prefix, metricName); @@ -115,7 +166,7 @@ private KafkaMbean addAttribute(KafkaMetric metric) { mbeans.put(mBeanName, new KafkaMbean(mBeanName)); KafkaMbean mbean = this.mbeans.get(mBeanName); mbean.setAttribute(metricName.name(), metric); - return mbean; + return mBeanName; } catch (JMException e) { throw new KafkaException("Error creating mbean attribute for metricName :" + metric.metricName(), e); } @@ -244,4 +295,27 @@ public AttributeList setAttributes(AttributeList list) { } + public static Predicate compilePredicate(Map configs) { + String whitelist = (String) configs.get(WHITELIST_CONFIG); + String blacklist = (String) configs.get(BLACKLIST_CONFIG); + + if (whitelist == null) { + whitelist = DEFAULT_WHITELIST; + } + + if (blacklist == null) { + blacklist = DEFAULT_BLACKLIST; + } + + try { + Pattern whitelistPattern = Pattern.compile(whitelist); + Pattern blacklistPattern = Pattern.compile(blacklist); + + return s -> whitelistPattern.matcher(s).matches() + && !blacklistPattern.matcher(s).matches(); + } catch (PatternSyntaxException e) { + throw new ConfigException("JMX filter for configuration" + METRICS_CONFIG_PREFIX + + ".(whitelist/blacklist) is not a valid regular expression"); + } + } } diff --git a/clients/src/main/java/org/apache/kafka/common/metrics/MetricsReporter.java b/clients/src/main/java/org/apache/kafka/common/metrics/MetricsReporter.java index 49d0d6fa8104b..cc112d181ced5 100644 --- a/clients/src/main/java/org/apache/kafka/common/metrics/MetricsReporter.java +++ b/clients/src/main/java/org/apache/kafka/common/metrics/MetricsReporter.java @@ -16,16 +16,20 @@ */ package org.apache.kafka.common.metrics; +import java.util.Collections; import java.util.List; +import java.util.Map; +import java.util.Set; -import org.apache.kafka.common.Configurable; +import org.apache.kafka.common.Reconfigurable; +import org.apache.kafka.common.config.ConfigException; /** * A plugin interface to allow things to listen as new metrics are created so they can be reported. *

* Implement {@link org.apache.kafka.common.ClusterResourceListener} to receive cluster metadata once it's available. Please see the class documentation for ClusterResourceListener for more information. */ -public interface MetricsReporter extends Configurable, AutoCloseable { +public interface MetricsReporter extends Reconfigurable, AutoCloseable { /** * This is called when the reporter is first registered to initially register all existing metrics @@ -50,4 +54,15 @@ public interface MetricsReporter extends Configurable, AutoCloseable { */ void close(); + // default methods for backwards compatibility with reporters that only implement Configurable + default Set reconfigurableConfigs() { + return Collections.emptySet(); + } + + default void validateReconfiguration(Map configs) throws ConfigException { + } + + default void reconfigure(Map configs) { + } + } diff --git a/clients/src/test/java/org/apache/kafka/common/metrics/JmxReporterTest.java b/clients/src/test/java/org/apache/kafka/common/metrics/JmxReporterTest.java index 37d1182262667..ea9597bad7eef 100644 --- a/clients/src/test/java/org/apache/kafka/common/metrics/JmxReporterTest.java +++ b/clients/src/test/java/org/apache/kafka/common/metrics/JmxReporterTest.java @@ -24,6 +24,8 @@ import javax.management.MBeanServer; import javax.management.ObjectName; import java.lang.management.ManagementFactory; +import java.util.HashMap; +import java.util.Map; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -111,4 +113,46 @@ public void testJmxRegistrationSanitization() throws Exception { metrics.close(); } } + + @Test + public void testPredicateAndDynamicReload() throws Exception { + Metrics metrics = new Metrics(); + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + + Map configs = new HashMap<>(); + + configs.put(JmxReporter.BLACKLIST_CONFIG, + JmxReporter.getMBeanName("", metrics.metricName("pack.bean2.total", "grp2"))); + + try { + JmxReporter reporter = new JmxReporter(); + reporter.configure(configs); + metrics.addReporter(reporter); + + Sensor sensor = metrics.sensor("kafka.requests"); + sensor.add(metrics.metricName("pack.bean2.avg", "grp1"), new Avg()); + sensor.add(metrics.metricName("pack.bean2.total", "grp2"), new CumulativeSum()); + sensor.record(); + + assertTrue(server.isRegistered(new ObjectName(":type=grp1"))); + assertEquals(1.0, server.getAttribute(new ObjectName(":type=grp1"), "pack.bean2.avg")); + assertFalse(server.isRegistered(new ObjectName(":type=grp2"))); + + sensor.record(); + + configs.put(JmxReporter.BLACKLIST_CONFIG, + JmxReporter.getMBeanName("", metrics.metricName("pack.bean2.avg", "grp1"))); + + reporter.reconfigure(configs); + + assertFalse(server.isRegistered(new ObjectName(":type=grp1"))); + assertTrue(server.isRegistered(new ObjectName(":type=grp2"))); + assertEquals(2.0, server.getAttribute(new ObjectName(":type=grp2"), "pack.bean2.total")); + + metrics.removeMetric(metrics.metricName("pack.bean2.total", "grp2")); + assertFalse(server.isRegistered(new ObjectName(":type=grp2"))); + } finally { + metrics.close(); + } + } } diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorConnectorConfig.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorConnectorConfig.java index d922eade5f372..527c1eb4bc2a7 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorConnectorConfig.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorConnectorConfig.java @@ -270,7 +270,9 @@ Map sourceAdminConfig() { List metricsReporters() { List reporters = getConfiguredInstances( CommonClientConfigs.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class); - reporters.add(new JmxReporter("kafka.connect.mirror")); + JmxReporter jmxReporter = new JmxReporter("kafka.connect.mirror"); + jmxReporter.configure(this.originals()); + reporters.add(jmxReporter); return reporters; } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java index 3f819a50e41de..57b5595563275 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java @@ -64,21 +64,20 @@ public class ConnectMetrics { * @param time the time; may not be null */ public ConnectMetrics(String workerId, WorkerConfig config, Time time) { - this(workerId, time, config.getInt(CommonClientConfigs.METRICS_NUM_SAMPLES_CONFIG), - config.getLong(CommonClientConfigs.METRICS_SAMPLE_WINDOW_MS_CONFIG), - config.getString(CommonClientConfigs.METRICS_RECORDING_LEVEL_CONFIG), - config.getConfiguredInstances(CommonClientConfigs.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class)); - } - - public ConnectMetrics(String workerId, Time time, int numSamples, long sampleWindowMs, String metricsRecordingLevel, - List reporters) { this.workerId = workerId; this.time = time; + int numSamples = config.getInt(CommonClientConfigs.METRICS_NUM_SAMPLES_CONFIG); + long sampleWindowMs = config.getLong(CommonClientConfigs.METRICS_SAMPLE_WINDOW_MS_CONFIG); + String metricsRecordingLevel = config.getString(CommonClientConfigs.METRICS_RECORDING_LEVEL_CONFIG); + List reporters = config.getConfiguredInstances(CommonClientConfigs.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class); + MetricConfig metricConfig = new MetricConfig().samples(numSamples) .timeWindow(sampleWindowMs, TimeUnit.MILLISECONDS).recordLevel( Sensor.RecordingLevel.forName(metricsRecordingLevel)); - reporters.add(new JmxReporter(JMX_PREFIX)); + JmxReporter jmxReporter = new JmxReporter(JMX_PREFIX); + jmxReporter.configure(config.originals()); + reporters.add(jmxReporter); this.metrics = new Metrics(metricConfig, reporters, time); LOG.debug("Registering Connect metrics with JMX for worker '{}'", workerId); AppInfoParser.registerAppInfo(JMX_PREFIX, workerId, metrics, time.milliseconds()); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java index c9be10f4d7dc6..a0a00e190e850 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java @@ -88,7 +88,9 @@ public WorkerGroupMember(DistributedConfig config, List reporters = config.getConfiguredInstances(CommonClientConfigs.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class, Collections.singletonMap(CommonClientConfigs.CLIENT_ID_CONFIG, clientId)); - reporters.add(new JmxReporter(JMX_PREFIX)); + JmxReporter jmxReporter = new JmxReporter(JMX_PREFIX); + jmxReporter.configure(config.originals()); + reporters.add(jmxReporter); this.metrics = new Metrics(metricConfig, reporters, time); this.retryBackoffMs = config.getLong(CommonClientConfigs.RETRY_BACKOFF_MS_CONFIG); this.metadata = new Metadata(retryBackoffMs, config.getLong(CommonClientConfigs.METADATA_MAX_AGE_CONFIG), diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java index 0deecd129a06e..419bea97f1200 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java @@ -24,8 +24,6 @@ import org.apache.kafka.connect.runtime.ConnectMetricsRegistry; import org.apache.kafka.connect.util.ConnectorTaskId; -import java.util.ArrayList; - /** * Contains various sensors used for monitoring errors. */ @@ -45,13 +43,6 @@ public class ErrorHandlingMetrics { private final Sensor dlqProduceFailures; private long lastErrorTime = 0; - // for testing only - public ErrorHandlingMetrics() { - this(new ConnectorTaskId("noop-connector", -1), - new ConnectMetrics("noop-worker", new SystemTime(), 2, 3000, Sensor.RecordingLevel.INFO.toString(), - new ArrayList<>())); - } - public ErrorHandlingMetrics(ConnectorTaskId id, ConnectMetrics connectMetrics) { ConnectMetricsRegistry registry = connectMetrics.registry(); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java index b6c3c8d688143..68d58a7d2fc34 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java @@ -85,8 +85,7 @@ import static org.junit.Assert.assertTrue; @PowerMockIgnore({"javax.management.*", - "org.apache.log4j.*", - "org.apache.kafka.connect.runtime.isolation.*"}) + "org.apache.log4j.*"}) @RunWith(PowerMockRunner.class) public class WorkerSourceTaskTest extends ThreadedTest { private static final String TOPIC = "topic"; diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/errors/RetryWithToleranceOperatorTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/errors/RetryWithToleranceOperatorTest.java index 2d340ac32e589..be723143e6630 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/errors/RetryWithToleranceOperatorTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/errors/RetryWithToleranceOperatorTest.java @@ -16,12 +16,20 @@ */ package org.apache.kafka.connect.runtime.errors; +import org.apache.kafka.clients.CommonClientConfigs; +import org.apache.kafka.common.metrics.Sensor; import org.apache.kafka.common.utils.MockTime; +import org.apache.kafka.common.utils.SystemTime; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.errors.RetriableException; +import org.apache.kafka.connect.runtime.ConnectMetrics; import org.apache.kafka.connect.runtime.ConnectorConfig; +import org.apache.kafka.connect.runtime.WorkerConfig; import org.apache.kafka.connect.runtime.isolation.Plugins; +import org.apache.kafka.connect.runtime.isolation.PluginsTest.TestConverter; +import org.apache.kafka.connect.runtime.isolation.PluginsTest.TestableWorkerConfig; import org.apache.kafka.connect.sink.SinkTask; +import org.apache.kafka.connect.util.ConnectorTaskId; import org.easymock.EasyMock; import org.easymock.Mock; import org.junit.Test; @@ -33,6 +41,7 @@ import java.util.HashMap; import java.util.Map; +import java.util.Objects; import static java.util.Collections.emptyMap; import static java.util.Collections.singletonMap; @@ -59,7 +68,19 @@ public class RetryWithToleranceOperatorTest { public static final RetryWithToleranceOperator NOOP_OPERATOR = new RetryWithToleranceOperator( ERRORS_RETRY_TIMEOUT_DEFAULT, ERRORS_RETRY_MAX_DELAY_DEFAULT, NONE, SYSTEM); static { - NOOP_OPERATOR.metrics(new ErrorHandlingMetrics()); + Map properties = new HashMap<>(); + properties.put(CommonClientConfigs.METRICS_NUM_SAMPLES_CONFIG, Objects.toString(2)); + properties.put(CommonClientConfigs.METRICS_SAMPLE_WINDOW_MS_CONFIG, Objects.toString(3000)); + properties.put(CommonClientConfigs.METRICS_RECORDING_LEVEL_CONFIG, Sensor.RecordingLevel.INFO.toString()); + + // define required properties + properties.put(WorkerConfig.KEY_CONVERTER_CLASS_CONFIG, TestConverter.class.getName()); + properties.put(WorkerConfig.VALUE_CONVERTER_CLASS_CONFIG, TestConverter.class.getName()); + + NOOP_OPERATOR.metrics(new ErrorHandlingMetrics( + new ConnectorTaskId("noop-connector", -1), + new ConnectMetrics("noop-worker", new TestableWorkerConfig(properties), new SystemTime())) + ); } @SuppressWarnings("unused") diff --git a/core/src/main/java/kafka/metrics/FilteringJmxReporter.java b/core/src/main/java/kafka/metrics/FilteringJmxReporter.java new file mode 100644 index 0000000000000..be2ba14b3df24 --- /dev/null +++ b/core/src/main/java/kafka/metrics/FilteringJmxReporter.java @@ -0,0 +1,62 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package kafka.metrics; + +import com.yammer.metrics.core.Metric; +import com.yammer.metrics.core.MetricName; +import com.yammer.metrics.core.MetricsRegistry; +import com.yammer.metrics.reporting.JmxReporter; + +import java.util.function.Predicate; + +public class FilteringJmxReporter extends JmxReporter { + + private volatile Predicate metricPredicate; + + public FilteringJmxReporter(MetricsRegistry registry, Predicate metricPredicate) { + super(registry); + this.metricPredicate = metricPredicate; + } + + @Override + public void onMetricAdded(MetricName name, Metric metric) { + if (metricPredicate.test(name)) { + super.onMetricAdded(name, metric); + } + } + + @Override + public void onMetricRemoved(MetricName name) { + super.onMetricRemoved(name); + } + + public void updatePredicate(Predicate predicate) { + this.metricPredicate = predicate; + // re-register metrics on update + getMetricsRegistry() + .allMetrics() + .forEach((name, metric) -> { + if (metricPredicate.test(name)) { + super.onMetricAdded(name, metric); + } else { + super.onMetricRemoved(name); + } + } + ); + } +} diff --git a/core/src/main/java/kafka/metrics/KafkaYammerMetrics.java b/core/src/main/java/kafka/metrics/KafkaYammerMetrics.java new file mode 100644 index 0000000000000..dd650fdd0f79e --- /dev/null +++ b/core/src/main/java/kafka/metrics/KafkaYammerMetrics.java @@ -0,0 +1,76 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package kafka.metrics; + +import com.yammer.metrics.core.MetricsRegistry; + +import org.apache.kafka.common.Reconfigurable; +import org.apache.kafka.common.config.ConfigException; +import org.apache.kafka.common.metrics.JmxReporter; + +import java.util.Map; +import java.util.Set; +import java.util.function.Predicate; + +/** + * This class encapsulates the default yammer metrics registry for Kafka server, + * and configures the set of exported JMX metrics for Yammer metrics. + * + * KafkaYammerMetrics.defaultRegistry() should always be used instead of Metrics.defaultRegistry() + */ +public class KafkaYammerMetrics implements Reconfigurable { + + public static final KafkaYammerMetrics INSTANCE = new KafkaYammerMetrics(); + + /** + * convenience method to replace {@link com.yammer.metrics.Metrics#defaultRegistry()} + */ + public static MetricsRegistry defaultRegistry() { + return INSTANCE.metricsRegistry; + } + + private final MetricsRegistry metricsRegistry = new MetricsRegistry(); + private final FilteringJmxReporter jmxReporter = new FilteringJmxReporter(metricsRegistry, + metricName -> true); + + private KafkaYammerMetrics() { + jmxReporter.start(); + Runtime.getRuntime().addShutdownHook(new Thread(jmxReporter::shutdown)); + } + + @Override + public void configure(Map configs) { + reconfigure(configs); + } + + @Override + public Set reconfigurableConfigs() { + return JmxReporter.RECONFIGURABLE_CONFIGS; + } + + @Override + public void validateReconfiguration(Map configs) throws ConfigException { + JmxReporter.compilePredicate(configs); + } + + @Override + public void reconfigure(Map configs) { + Predicate mBeanPredicate = JmxReporter.compilePredicate(configs); + jmxReporter.updatePredicate(metricName -> mBeanPredicate.test(metricName.getMBeanName())); + } +} diff --git a/core/src/main/scala/kafka/metrics/KafkaCSVMetricsReporter.scala b/core/src/main/scala/kafka/metrics/KafkaCSVMetricsReporter.scala index 82317277814c5..0d8354728ac8d 100755 --- a/core/src/main/scala/kafka/metrics/KafkaCSVMetricsReporter.scala +++ b/core/src/main/scala/kafka/metrics/KafkaCSVMetricsReporter.scala @@ -20,7 +20,6 @@ package kafka.metrics -import com.yammer.metrics.Metrics import java.io.File import java.nio.file.Files @@ -52,7 +51,7 @@ private class KafkaCSVMetricsReporter extends KafkaMetricsReporter csvDir = new File(props.getString("kafka.csv.metrics.dir", "kafka_metrics")) Utils.delete(csvDir) Files.createDirectories(csvDir.toPath()) - underlying = new CsvReporter(Metrics.defaultRegistry(), csvDir) + underlying = new CsvReporter(KafkaYammerMetrics.defaultRegistry(), csvDir) if (props.getBoolean("kafka.csv.metrics.reporter.enabled", default = false)) { initialized = true startReporter(metricsConfig.pollingIntervalSecs) @@ -79,7 +78,7 @@ private class KafkaCSVMetricsReporter extends KafkaMetricsReporter underlying.shutdown() running = false info("Stopped Kafka CSV metrics reporter") - underlying = new CsvReporter(Metrics.defaultRegistry(), csvDir) + underlying = new CsvReporter(KafkaYammerMetrics.defaultRegistry(), csvDir) } } } diff --git a/core/src/main/scala/kafka/metrics/KafkaMetricsGroup.scala b/core/src/main/scala/kafka/metrics/KafkaMetricsGroup.scala index 53088b347d9bf..a63be1fe2693c 100644 --- a/core/src/main/scala/kafka/metrics/KafkaMetricsGroup.scala +++ b/core/src/main/scala/kafka/metrics/KafkaMetricsGroup.scala @@ -19,7 +19,6 @@ package kafka.metrics import java.util.concurrent.TimeUnit -import com.yammer.metrics.Metrics import com.yammer.metrics.core.{Gauge, MetricName, Meter, Histogram, Timer} import kafka.utils.Logging import org.apache.kafka.common.utils.Sanitizer @@ -66,19 +65,19 @@ trait KafkaMetricsGroup extends Logging { } def newGauge[T](name: String, metric: Gauge[T], tags: scala.collection.Map[String, String] = Map.empty): Gauge[T] = - Metrics.defaultRegistry().newGauge(metricName(name, tags), metric) + KafkaYammerMetrics.defaultRegistry().newGauge(metricName(name, tags), metric) def newMeter(name: String, eventType: String, timeUnit: TimeUnit, tags: scala.collection.Map[String, String] = Map.empty): Meter = - Metrics.defaultRegistry().newMeter(metricName(name, tags), eventType, timeUnit) + KafkaYammerMetrics.defaultRegistry().newMeter(metricName(name, tags), eventType, timeUnit) def newHistogram(name: String, biased: Boolean = true, tags: scala.collection.Map[String, String] = Map.empty): Histogram = - Metrics.defaultRegistry().newHistogram(metricName(name, tags), biased) + KafkaYammerMetrics.defaultRegistry().newHistogram(metricName(name, tags), biased) def newTimer(name: String, durationUnit: TimeUnit, rateUnit: TimeUnit, tags: scala.collection.Map[String, String] = Map.empty): Timer = - Metrics.defaultRegistry().newTimer(metricName(name, tags), durationUnit, rateUnit) + KafkaYammerMetrics.defaultRegistry().newTimer(metricName(name, tags), durationUnit, rateUnit) def removeMetric(name: String, tags: scala.collection.Map[String, String] = Map.empty): Unit = - Metrics.defaultRegistry().removeMetric(metricName(name, tags)) + KafkaYammerMetrics.defaultRegistry().removeMetric(metricName(name, tags)) private def toMBeanName(tags: collection.Map[String, String]): Option[String] = { val filteredTags = tags.filter { case (_, tagValue) => tagValue != "" } diff --git a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala index 75081847e3e98..a0b8e562c63ed 100755 --- a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala +++ b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala @@ -237,6 +237,7 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging case Some(authz: Reconfigurable) => addReconfigurable(authz) case _ => } + addReconfigurable(kafkaServer.kafkaYammerMetrics) addReconfigurable(new DynamicMetricsReporters(kafkaConfig.brokerId, kafkaServer)) addReconfigurable(new DynamicClientQuotaCallback(kafkaConfig.brokerId, kafkaServer)) diff --git a/core/src/main/scala/kafka/server/KafkaServer.scala b/core/src/main/scala/kafka/server/KafkaServer.scala index 8614c0c9be1da..4713f33615fe1 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -30,7 +30,7 @@ import kafka.controller.KafkaController import kafka.coordinator.group.GroupCoordinator import kafka.coordinator.transaction.TransactionCoordinator import kafka.log.{LogConfig, LogManager} -import kafka.metrics.{KafkaMetricsGroup, KafkaMetricsReporter} +import kafka.metrics.{KafkaMetricsGroup, KafkaMetricsReporter, KafkaYammerMetrics} import kafka.network.SocketServer import kafka.security.CredentialProvider import kafka.utils._ @@ -133,6 +133,7 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP private var logContext: LogContext = null + var kafkaYammerMetrics: KafkaYammerMetrics = null var metrics: Metrics = null val brokerState: BrokerState = new BrokerState @@ -186,7 +187,7 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP newGauge("BrokerState", () => brokerState.currentState) newGauge("ClusterId", () => clusterId) - newGauge("yammer-metrics-count", () => com.yammer.metrics.Metrics.defaultRegistry.allMetrics.size) + newGauge("yammer-metrics-count", () => KafkaYammerMetrics.defaultRegistry.allMetrics.size) /** * Start up API for bringing up a single instance of the Kafka server. @@ -236,8 +237,14 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP kafkaScheduler.startup() /* create and configure metrics */ + kafkaYammerMetrics = KafkaYammerMetrics.INSTANCE + kafkaYammerMetrics.configure(config.originals) + + val jmxReporter = new JmxReporter(jmxPrefix) + jmxReporter.configure(config.originals) + val reporters = new util.ArrayList[MetricsReporter] - reporters.add(new JmxReporter(jmxPrefix)) + reporters.add(jmxReporter) val metricConfig = KafkaServer.metricConfig(config) metrics = new Metrics(metricConfig, reporters, time, true) diff --git a/core/src/test/scala/integration/kafka/api/EndToEndAuthorizationTest.scala b/core/src/test/scala/integration/kafka/api/EndToEndAuthorizationTest.scala index 90ae39a70c02e..e6d05f6b9b07f 100644 --- a/core/src/test/scala/integration/kafka/api/EndToEndAuthorizationTest.scala +++ b/core/src/test/scala/integration/kafka/api/EndToEndAuthorizationTest.scala @@ -17,13 +17,13 @@ package kafka.api -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Gauge import java.io.File import java.util.Collections import java.util.concurrent.ExecutionException import kafka.admin.AclCommand +import kafka.metrics.KafkaYammerMetrics import kafka.security.authorizer.AclAuthorizer import kafka.security.authorizer.AclEntry.WildcardHost import kafka.server._ @@ -232,7 +232,7 @@ abstract class EndToEndAuthorizationTest extends IntegrationTestHarness with Sas } private def getGauge(metricName: String) = { - Metrics.defaultRegistry.allMetrics.asScala + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala .filterKeys(k => k.getName == metricName) .headOption .getOrElse { fail( "Unable to find metric " + metricName ) } diff --git a/core/src/test/scala/integration/kafka/api/MetricsTest.scala b/core/src/test/scala/integration/kafka/api/MetricsTest.scala index 83c75dfb77f2b..a851b4ee02de4 100644 --- a/core/src/test/scala/integration/kafka/api/MetricsTest.scala +++ b/core/src/test/scala/integration/kafka/api/MetricsTest.scala @@ -17,8 +17,8 @@ import java.util.{Locale, Properties} import kafka.log.LogConfig import kafka.server.{KafkaConfig, KafkaServer} import kafka.utils.{JaasTestUtils, TestUtils} -import com.yammer.metrics.Metrics import com.yammer.metrics.core.{Gauge, Histogram, Meter} +import kafka.metrics.KafkaYammerMetrics import org.apache.kafka.clients.consumer.KafkaConsumer import org.apache.kafka.clients.producer.{KafkaProducer, ProducerConfig, ProducerRecord} import org.apache.kafka.common.{Metric, MetricName, TopicPartition} @@ -222,7 +222,7 @@ class MetricsTest extends IntegrationTestHarness with SaslSetup { private def verifyBrokerErrorMetrics(server: KafkaServer): Unit = { - def errorMetricCount = Metrics.defaultRegistry.allMetrics.keySet.asScala.filter(_.getName == "ErrorsPerSec").size + def errorMetricCount = KafkaYammerMetrics.defaultRegistry.allMetrics.keySet.asScala.filter(_.getName == "ErrorsPerSec").size val startErrorMetricCount = errorMetricCount val errorMetricPrefix = "kafka.network:type=RequestMetrics,name=ErrorsPerSec" @@ -271,7 +271,7 @@ class MetricsTest extends IntegrationTestHarness with SaslSetup { } private def yammerMetricValue(name: String): Any = { - val allMetrics = Metrics.defaultRegistry.allMetrics.asScala + val allMetrics = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala val (_, metric) = allMetrics.find { case (n, _) => n.getMBeanName.endsWith(name) } .getOrElse(fail(s"Unable to find broker metric $name: allMetrics: ${allMetrics.keySet.map(_.getMBeanName)}")) metric match { @@ -283,7 +283,7 @@ class MetricsTest extends IntegrationTestHarness with SaslSetup { } private def yammerHistogram(name: String): Histogram = { - val allMetrics = Metrics.defaultRegistry.allMetrics.asScala + val allMetrics = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala val (_, metric) = allMetrics.find { case (n, _) => n.getMBeanName.endsWith(name) } .getOrElse(fail(s"Unable to find broker metric $name: allMetrics: ${allMetrics.keySet.map(_.getMBeanName)}")) metric match { @@ -299,7 +299,7 @@ class MetricsTest extends IntegrationTestHarness with SaslSetup { } private def verifyNoRequestMetrics(errorMessage: String): Unit = { - val metrics = Metrics.defaultRegistry.allMetrics.asScala.filter { case (n, _) => + val metrics = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.filter { case (n, _) => n.getMBeanName.startsWith("kafka.network:type=RequestMetrics") } assertTrue(s"$errorMessage: ${metrics.keys}", metrics.isEmpty) diff --git a/core/src/test/scala/integration/kafka/api/SslAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/SslAdminIntegrationTest.scala index 142dbca5942ec..9b1a645a4c929 100644 --- a/core/src/test/scala/integration/kafka/api/SslAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/SslAdminIntegrationTest.scala @@ -17,8 +17,8 @@ import java.util import java.util.Collections import java.util.concurrent._ -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Gauge +import kafka.metrics.KafkaYammerMetrics import kafka.security.authorizer.AclAuthorizer import kafka.security.authorizer.AclEntry.{WildcardHost, WildcardPrincipalString} import kafka.server.KafkaConfig @@ -259,7 +259,7 @@ class SslAdminIntegrationTest extends SaslSslAdminIntegrationTest { } private def purgatoryMetric(name: String): Int = { - val allMetrics = Metrics.defaultRegistry.allMetrics.asScala + val allMetrics = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala val metrics = allMetrics.filter { case (metricName, _) => metricName.getMBeanName.contains("delayedOperation=AlterAcls") && metricName.getMBeanName.contains(s"name=$name") }.values.toList diff --git a/core/src/test/scala/integration/kafka/server/DynamicBrokerReconfigurationTest.scala b/core/src/test/scala/integration/kafka/server/DynamicBrokerReconfigurationTest.scala index 00df09af782c1..4705b0e47a19f 100644 --- a/core/src/test/scala/integration/kafka/server/DynamicBrokerReconfigurationTest.scala +++ b/core/src/test/scala/integration/kafka/server/DynamicBrokerReconfigurationTest.scala @@ -26,15 +26,15 @@ import java.time.Duration import java.util import java.util.{Collections, Properties} import java.util.concurrent._ -import javax.management.ObjectName -import com.yammer.metrics.Metrics +import javax.management.ObjectName import com.yammer.metrics.core.MetricName import kafka.admin.ConfigCommand import kafka.api.{KafkaSasl, SaslSetup} import kafka.controller.{ControllerBrokerStateInfo, ControllerChannelManager} import kafka.log.LogConfig import kafka.message.ProducerCompressionCodec +import kafka.metrics.KafkaYammerMetrics import kafka.network.{Processor, RequestChannel} import kafka.utils._ import kafka.utils.Implicits._ @@ -780,8 +780,8 @@ class DynamicBrokerReconfigurationTest extends ZooKeeperTestHarness with SaslSet } private def clearLeftOverProcessorMetrics(): Unit = { - val metricsFromOldTests = Metrics.defaultRegistry.allMetrics.keySet.asScala.filter(isProcessorMetric) - metricsFromOldTests.foreach(Metrics.defaultRegistry.removeMetric) + val metricsFromOldTests = KafkaYammerMetrics.defaultRegistry.allMetrics.keySet.asScala.filter(isProcessorMetric) + metricsFromOldTests.foreach(KafkaYammerMetrics.defaultRegistry.removeMetric) } // Verify that metrics from processors that were removed have been deleted. @@ -795,7 +795,7 @@ class DynamicBrokerReconfigurationTest extends ZooKeeperTestHarness with SaslSet .groupBy(_.tags.get(Processor.NetworkProcessorMetricTag)) assertEquals(numProcessors, kafkaMetrics.size) - Metrics.defaultRegistry.allMetrics.keySet.asScala + KafkaYammerMetrics.defaultRegistry.allMetrics.keySet.asScala .filter(isProcessorMetric) .groupBy(_.getName) .foreach { case (name, set) => assertEquals(s"Metrics not deleted $name", numProcessors, set.size) } diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index 13dd2f1eb4872..05b26dfc84842 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -18,14 +18,14 @@ package kafka.cluster import java.nio.ByteBuffer import java.util.{Optional, Properties} -import java.util.concurrent.{CountDownLatch, Executors, TimeoutException, TimeUnit} +import java.util.concurrent.{CountDownLatch, Executors, TimeUnit, TimeoutException} import java.util.concurrent.atomic.AtomicBoolean -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Metric import kafka.api.{ApiVersion, LeaderAndIsr} import kafka.common.UnexpectedAppendOffsetException import kafka.log.{Defaults => _, _} +import kafka.metrics.KafkaYammerMetrics import kafka.server._ import kafka.utils._ import org.apache.kafka.common.{IsolationLevel, TopicPartition} @@ -1485,7 +1485,7 @@ class PartitionTest extends AbstractPartitionTest { "AtMinIsr") def getMetric(metric: String): Option[Metric] = { - Metrics.defaultRegistry().allMetrics().asScala.filterKeys { metricName => + KafkaYammerMetrics.defaultRegistry().allMetrics().asScala.filterKeys { metricName => metricName.getName == metric && metricName.getType == "Partition" }.headOption.map(_._2) } @@ -1494,7 +1494,7 @@ class PartitionTest extends AbstractPartitionTest { Partition.removeMetrics(topicPartition) - assertEquals(Set(), Metrics.defaultRegistry().allMetrics().asScala.keySet.filter(_.getType == "Partition")) + assertEquals(Set(), KafkaYammerMetrics.defaultRegistry().allMetrics().asScala.keySet.filter(_.getType == "Partition")) } @Test diff --git a/core/src/test/scala/unit/kafka/controller/ControllerEventManagerTest.scala b/core/src/test/scala/unit/kafka/controller/ControllerEventManagerTest.scala index 7c28ff85c794a..4aa02ccf97a4b 100644 --- a/core/src/test/scala/unit/kafka/controller/ControllerEventManagerTest.scala +++ b/core/src/test/scala/unit/kafka/controller/ControllerEventManagerTest.scala @@ -20,9 +20,9 @@ package kafka.controller import java.util.concurrent.CountDownLatch import java.util.concurrent.atomic.AtomicInteger -import com.yammer.metrics.Metrics import com.yammer.metrics.core.{Histogram, MetricName, Timer} import kafka.controller +import kafka.metrics.KafkaYammerMetrics import kafka.utils.TestUtils import org.apache.kafka.common.message.UpdateMetadataResponseData import org.apache.kafka.common.protocol.Errors @@ -54,7 +54,7 @@ class ControllerEventManagerTest { } def allEventManagerMetrics: Set[MetricName] = { - Metrics.defaultRegistry.allMetrics.asScala.keySet + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.keySet .filter(_.getMBeanName.startsWith("kafka.controller:type=ControllerEventManager")) .toSet } @@ -111,7 +111,7 @@ class ControllerEventManagerTest { } // The metric should not already exist - assertTrue(Metrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.isEmpty) + assertTrue(KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.isEmpty) controllerEventManager = new ControllerEventManager(0, eventProcessor, time, controllerStats.rateAndTimeMetrics) @@ -124,7 +124,7 @@ class ControllerEventManagerTest { TestUtils.waitUntilTrue(() => processedEvents.get() == 2, "Timed out waiting for processing of all events") - val queueTimeHistogram = Metrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.headOption + val queueTimeHistogram = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.headOption .getOrElse(fail(s"Unable to find metric $metricName")).asInstanceOf[Histogram] assertEquals(2, queueTimeHistogram.count) @@ -179,7 +179,7 @@ class ControllerEventManagerTest { } private def timer(metricName: String): Timer = { - Metrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.headOption + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.headOption .getOrElse(fail(s"Unable to find metric $metricName")).asInstanceOf[Timer] } diff --git a/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala b/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala index 018a0bb8cb014..771c81c162f68 100644 --- a/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala +++ b/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala @@ -20,9 +20,9 @@ package kafka.controller import java.util.Properties import java.util.concurrent.{CountDownLatch, LinkedBlockingQueue} -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Timer import kafka.api.LeaderAndIsr +import kafka.metrics.KafkaYammerMetrics import kafka.server.{KafkaConfig, KafkaServer} import kafka.utils.TestUtils import kafka.zk._ @@ -698,7 +698,7 @@ class ControllerIntegrationTest extends ZooKeeperTestHarness { } private def timer(metricName: String): Timer = { - Metrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.headOption + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getMBeanName == metricName).values.headOption .getOrElse(fail(s"Unable to find metric $metricName")).asInstanceOf[Timer] } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 3e3fb66e95a58..783e94450d702 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -22,13 +22,13 @@ import java.nio.ByteBuffer import java.util.concurrent.locks.ReentrantLock import java.util.{Collections, Optional} -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Gauge import javax.management.ObjectName import kafka.api._ import kafka.cluster.Partition import kafka.common.OffsetAndMetadata import kafka.log.{AppendOrigin, Log, LogAppendInfo} +import kafka.metrics.KafkaYammerMetrics import kafka.server.{FetchDataInfo, FetchLogEnd, HostedPartition, KafkaConfig, LogOffsetMetadata, ReplicaManager} import kafka.utils.{KafkaScheduler, MockTime, TestUtils} import kafka.zk.KafkaZkClient @@ -2327,7 +2327,7 @@ class GroupMetadataManagerTest { } private def getGauge(manager: GroupMetadataManager, name: String): Gauge[Int] = { - Metrics.defaultRegistry().allMetrics().get(manager.metricName(name, Map.empty)).asInstanceOf[Gauge[Int]] + KafkaYammerMetrics.defaultRegistry().allMetrics().get(manager.metricName(name, Map.empty)).asInstanceOf[Gauge[Int]] } private def expectMetrics(manager: GroupMetadataManager, diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerChannelManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerChannelManagerTest.scala index 2b7b805c71ce9..2374275a52606 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerChannelManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionMarkerChannelManagerTest.scala @@ -29,8 +29,8 @@ import org.apache.kafka.common.{Node, TopicPartition} import org.easymock.{Capture, EasyMock} import org.junit.Assert._ import org.junit.Test -import com.yammer.metrics.Metrics import kafka.common.RequestAndCompletionHandler +import kafka.metrics.KafkaYammerMetrics import org.apache.kafka.common.protocol.{ApiKeys, Errors} import org.apache.kafka.common.record.RecordBatch @@ -423,7 +423,7 @@ class TransactionMarkerChannelManagerTest { @Test def shouldCreateMetricsOnStarting(): Unit = { - val metrics = Metrics.defaultRegistry.allMetrics.asScala + val metrics = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala assertEquals(1, metrics .filterKeys(_.getMBeanName == "kafka.coordinator.transaction:type=TransactionMarkerChannelManager,name=UnknownDestinationQueueSize") diff --git a/core/src/test/scala/unit/kafka/integration/MetricsDuringTopicCreationDeletionTest.scala b/core/src/test/scala/unit/kafka/integration/MetricsDuringTopicCreationDeletionTest.scala index 48f9d79d6b6a0..c4fea92d43daf 100644 --- a/core/src/test/scala/unit/kafka/integration/MetricsDuringTopicCreationDeletionTest.scala +++ b/core/src/test/scala/unit/kafka/integration/MetricsDuringTopicCreationDeletionTest.scala @@ -18,14 +18,15 @@ package kafka.integration import java.util.Properties + import kafka.server.KafkaConfig import kafka.utils.{Logging, TestUtils} + import scala.collection.JavaConverters.mapAsScalaMapConverter import org.scalatest.Assertions.fail - import org.junit.{Before, Test} -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Gauge +import kafka.metrics.KafkaYammerMetrics class MetricsDuringTopicCreationDeletionTest extends KafkaServerTestHarness with Logging { @@ -57,8 +58,8 @@ class MetricsDuringTopicCreationDeletionTest extends KafkaServerTestHarness with // This is a test workaround to the issue that prior harness runs may have left a populated registry. // see https://issues.apache.org/jira/browse/KAFKA-4605 for (m <- testedMetrics) { - val metricName = Metrics.defaultRegistry.allMetrics.asScala.keys.find(_.getName.endsWith(m)) - metricName.foreach(Metrics.defaultRegistry.removeMetric) + val metricName = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.keys.find(_.getName.endsWith(m)) + metricName.foreach(KafkaYammerMetrics.defaultRegistry.removeMetric) } super.setUp @@ -122,11 +123,11 @@ class MetricsDuringTopicCreationDeletionTest extends KafkaServerTestHarness with } private def getGauge(metricName: String) = { - Metrics.defaultRegistry.allMetrics.asScala - .filterKeys(k => k.getName.endsWith(metricName)) - .headOption - .getOrElse { fail( "Unable to find metric " + metricName ) } - ._2.asInstanceOf[Gauge[Int]] + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala + .filterKeys(k => k.getName.endsWith(metricName)) + .headOption + .getOrElse { fail( "Unable to find metric " + metricName ) } + ._2.asInstanceOf[Gauge[Int]] } private def createDeleteTopics(): Unit = { diff --git a/core/src/test/scala/unit/kafka/log/LogCleanerIntegrationTest.scala b/core/src/test/scala/unit/kafka/log/LogCleanerIntegrationTest.scala index 7aea6c1fcdd80..52c4d7e53f5da 100644 --- a/core/src/test/scala/unit/kafka/log/LogCleanerIntegrationTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogCleanerIntegrationTest.scala @@ -19,9 +19,8 @@ package kafka.log import java.io.PrintWriter -import com.yammer.metrics.Metrics import com.yammer.metrics.core.{Gauge, MetricName} -import kafka.metrics.KafkaMetricsGroup +import kafka.metrics.{KafkaMetricsGroup, KafkaYammerMetrics} import kafka.utils.{MockTime, TestUtils} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.record.{CompressionType, RecordBatch} @@ -90,7 +89,7 @@ class LogCleanerIntegrationTest extends AbstractLogCleanerIntegrationTest with K } private def getGauge[T](filter: MetricName => Boolean): Gauge[T] = { - Metrics.defaultRegistry.allMetrics.asScala + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala .filterKeys(filter(_)) .headOption .getOrElse { fail(s"Unable to find metric") } diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index 6bbedbbba1dab..6d147b4acabb5 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -23,10 +23,10 @@ import java.nio.file.{Files, Paths} import java.util.regex.Pattern import java.util.{Collections, Optional, Properties} -import com.yammer.metrics.Metrics import kafka.api.{ApiVersion, KAFKA_0_11_0_IV0} import kafka.common.{OffsetsOutOfOrderException, RecordValidationException, UnexpectedAppendOffsetException} import kafka.log.Log.DeleteDirSuffix +import kafka.metrics.KafkaYammerMetrics import kafka.server.checkpoints.LeaderEpochCheckpointFile import kafka.server.epoch.{EpochEntry, LeaderEpochFileCache} import kafka.server.{BrokerTopicStats, FetchDataInfo, FetchHighWatermark, FetchIsolation, FetchLogEnd, FetchTxnCommitted, KafkaConfig, LogDirFailureChannel, LogOffsetMetadata} @@ -56,7 +56,7 @@ class LogTest { val tmpDir = TestUtils.tempDir() val logDir = TestUtils.randomPartitionLogDir(tmpDir) val mockTime = new MockTime() - def metricsKeySet = Metrics.defaultRegistry.allMetrics.keySet.asScala + def metricsKeySet = KafkaYammerMetrics.defaultRegistry.allMetrics.keySet.asScala @Before def setUp(): Unit = { diff --git a/core/src/test/scala/unit/kafka/log/LogValidatorTest.scala b/core/src/test/scala/unit/kafka/log/LogValidatorTest.scala index 3dc43740e93c3..5de188f8214f8 100644 --- a/core/src/test/scala/unit/kafka/log/LogValidatorTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogValidatorTest.scala @@ -19,10 +19,10 @@ package kafka.log import java.nio.ByteBuffer import java.util.concurrent.TimeUnit -import com.yammer.metrics.Metrics import kafka.api.{ApiVersion, KAFKA_2_0_IV1, KAFKA_2_3_IV1} import kafka.common.{LongRef, RecordValidationException} import kafka.message._ +import kafka.metrics.KafkaYammerMetrics import kafka.server.BrokerTopicStats import kafka.utils.TestUtils.meterCount import org.apache.kafka.common.InvalidRecordException @@ -42,7 +42,7 @@ class LogValidatorTest { val time = Time.SYSTEM val topicPartition = new TopicPartition("topic", 0) val brokerTopicStats = new BrokerTopicStats - val metricsKeySet = Metrics.defaultRegistry.allMetrics.keySet.asScala + val metricsKeySet = KafkaYammerMetrics.defaultRegistry.allMetrics.keySet.asScala @Test def testOnlyOneBatch(): Unit = { diff --git a/core/src/test/scala/unit/kafka/metrics/MetricsTest.scala b/core/src/test/scala/unit/kafka/metrics/MetricsTest.scala index 62040d5e262f9..80bcbc71034dd 100644 --- a/core/src/test/scala/unit/kafka/metrics/MetricsTest.scala +++ b/core/src/test/scala/unit/kafka/metrics/MetricsTest.scala @@ -17,10 +17,10 @@ package kafka.metrics +import java.lang.management.ManagementFactory import java.util.Properties import javax.management.ObjectName -import com.yammer.metrics.Metrics import com.yammer.metrics.core.MetricPredicate import org.junit.Test import org.junit.Assert._ @@ -32,6 +32,7 @@ import scala.collection._ import scala.collection.JavaConverters._ import kafka.log.LogConfig import org.apache.kafka.common.TopicPartition +import org.apache.kafka.common.metrics.JmxReporter class MetricsTest extends KafkaServerTestHarness with Logging { val numNodes = 2 @@ -39,6 +40,7 @@ class MetricsTest extends KafkaServerTestHarness with Logging { val overridingProps = new Properties overridingProps.put(KafkaConfig.NumPartitionsProp, numParts.toString) + overridingProps.put(JmxReporter.BLACKLIST_CONFIG, "kafka.server:type=KafkaServer,name=ClusterId") def generateConfigs = TestUtils.createBrokerConfigs(numNodes, zkConnect).map(KafkaConfig.fromProps(_, overridingProps)) @@ -71,10 +73,31 @@ class MetricsTest extends KafkaServerTestHarness with Logging { @Test def testClusterIdMetric(): Unit = { // Check if clusterId metric exists. - val metrics = Metrics.defaultRegistry.allMetrics + val metrics = KafkaYammerMetrics.defaultRegistry.allMetrics assertEquals(metrics.keySet.asScala.count(_.getMBeanName == "kafka.server:type=KafkaServer,name=ClusterId"), 1) } + @Test + def testJMXFilter(): Unit = { + // Check if cluster id metrics is not exposed in JMX + assertTrue(ManagementFactory.getPlatformMBeanServer + .isRegistered(new ObjectName("kafka.controller:type=KafkaController,name=ActiveControllerCount"))) + assertFalse(ManagementFactory.getPlatformMBeanServer + .isRegistered(new ObjectName("kafka.server:type=KafkaServer,name=ClusterId"))) + } + + @Test + def testUpdateJMXFilter(): Unit = { + // verify previously exposed metrics are removed and existing matching metrics are added + servers.foreach(server => server.kafkaYammerMetrics.reconfigure( + Map(JmxReporter.BLACKLIST_CONFIG -> "kafka.controller:type=KafkaController,name=ActiveControllerCount").asJava + )) + assertFalse(ManagementFactory.getPlatformMBeanServer + .isRegistered(new ObjectName("kafka.controller:type=KafkaController,name=ActiveControllerCount"))) + assertTrue(ManagementFactory.getPlatformMBeanServer + .isRegistered(new ObjectName("kafka.server:type=KafkaServer,name=ClusterId"))) + } + @Test def testGeneralBrokerTopicMetricsAreGreedilyRegistered(): Unit = { val topic = "test-broker-topic-metric" @@ -148,7 +171,7 @@ class MetricsTest extends KafkaServerTestHarness with Logging { @Test def testControllerMetrics(): Unit = { - val metrics = Metrics.defaultRegistry.allMetrics + val metrics = KafkaYammerMetrics.defaultRegistry.allMetrics assertEquals(metrics.keySet.asScala.count(_.getMBeanName == "kafka.controller:type=KafkaController,name=ActiveControllerCount"), 1) assertEquals(metrics.keySet.asScala.count(_.getMBeanName == "kafka.controller:type=KafkaController,name=OfflinePartitionsCount"), 1) @@ -167,7 +190,7 @@ class MetricsTest extends KafkaServerTestHarness with Logging { */ @Test def testSessionExpireListenerMetrics(): Unit = { - val metrics = Metrics.defaultRegistry.allMetrics + val metrics = KafkaYammerMetrics.defaultRegistry.allMetrics assertEquals(metrics.keySet.asScala.count(_.getMBeanName == "kafka.server:type=SessionExpireListener,name=SessionState"), 1) assertEquals(metrics.keySet.asScala.count(_.getMBeanName == "kafka.server:type=SessionExpireListener,name=ZooKeeperExpiresPerSec"), 1) @@ -175,12 +198,12 @@ class MetricsTest extends KafkaServerTestHarness with Logging { } private def topicMetrics(topic: Option[String]): Set[String] = { - val metricNames = Metrics.defaultRegistry.allMetrics().keySet.asScala.map(_.getMBeanName) + val metricNames = KafkaYammerMetrics.defaultRegistry.allMetrics().keySet.asScala.map(_.getMBeanName) filterByTopicMetricRegex(metricNames, topic) } private def topicMetricGroups(topic: String): Set[String] = { - val metricGroups = Metrics.defaultRegistry.groupedMetrics(MetricPredicate.ALL).keySet.asScala + val metricGroups = KafkaYammerMetrics.defaultRegistry.groupedMetrics(MetricPredicate.ALL).keySet.asScala filterByTopicMetricRegex(metricGroups, Some(topic)) } diff --git a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala index 91fce5b23c3db..3b84d6f0e94e0 100644 --- a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala +++ b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala @@ -27,9 +27,8 @@ import java.util.concurrent.{CompletableFuture, ConcurrentLinkedQueue, Executors import java.util.{HashMap, Properties, Random} import com.yammer.metrics.core.{Gauge, Meter} -import com.yammer.metrics.{Metrics => YammerMetrics} import javax.net.ssl._ - +import kafka.metrics.KafkaYammerMetrics import kafka.security.CredentialProvider import kafka.server.{KafkaConfig, ThrottledChannel} import kafka.utils.Implicits._ @@ -1100,7 +1099,7 @@ class SocketServerTest { s"kafka.network:type=RequestMetrics,name=RequestsPerSec,request=Produce,version=$version2" -> 1, "kafka.network:type=RequestMetrics,name=ErrorsPerSec,request=Produce,error=NONE" -> 1) - def requestMetricMeters = YammerMetrics + def requestMetricMeters = KafkaYammerMetrics .defaultRegistry .allMetrics.asScala .collect { case (k, metric: Meter) if k.getType == "RequestMetrics" => (k.toString, metric.count) } @@ -1114,7 +1113,7 @@ class SocketServerTest { def testMetricCollectionAfterShutdown(): Unit = { server.shutdown() - val nonZeroMetricNamesAndValues = YammerMetrics + val nonZeroMetricNamesAndValues = KafkaYammerMetrics .defaultRegistry .allMetrics.asScala .filter { case (k, _) => k.getName.endsWith("IdlePercent") || k.getName.endsWith("NetworkProcessorAvgIdlePercent") } @@ -1135,7 +1134,7 @@ class SocketServerTest { } // legacy metrics not tagged - val yammerMetricsNames = YammerMetrics.defaultRegistry.allMetrics.asScala + val yammerMetricsNames = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala .filterKeys(_.getType.equals("Processor")) .collect { case (k, _: Gauge[_]) => k } assertFalse(yammerMetricsNames.isEmpty) @@ -1684,7 +1683,7 @@ class SocketServerTest { private def verifyAcceptorBlockedPercent(listenerName: String, expectBlocked: Boolean): Unit = { val blockedPercentMetricMBeanName = "kafka.network:type=Acceptor,name=AcceptorBlockedPercent,listener=PLAINTEXT" - val blockedPercentMetrics = YammerMetrics.defaultRegistry.allMetrics.asScala + val blockedPercentMetrics = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala .filterKeys(_.getMBeanName == blockedPercentMetricMBeanName).values assertEquals(1, blockedPercentMetrics.size) val blockedPercentMetric = blockedPercentMetrics.head.asInstanceOf[Meter] diff --git a/core/src/test/scala/unit/kafka/server/AbstractFetcherManagerTest.scala b/core/src/test/scala/unit/kafka/server/AbstractFetcherManagerTest.scala index ecd92bfba2961..211bf56d6e509 100644 --- a/core/src/test/scala/unit/kafka/server/AbstractFetcherManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/AbstractFetcherManagerTest.scala @@ -16,9 +16,9 @@ */ package kafka.server -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Gauge import kafka.cluster.BrokerEndPoint +import kafka.metrics.KafkaYammerMetrics import kafka.utils.TestUtils import org.apache.kafka.common.TopicPartition import org.easymock.EasyMock @@ -35,7 +35,7 @@ class AbstractFetcherManagerTest { } private def getMetricValue(name: String): Any = { - Metrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getName == name).values.headOption.get. + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.filterKeys(_.getName == name).values.headOption.get. asInstanceOf[Gauge[Int]].value() } diff --git a/core/src/test/scala/unit/kafka/server/AbstractFetcherThreadTest.scala b/core/src/test/scala/unit/kafka/server/AbstractFetcherThreadTest.scala index 6ff62ecee9b98..4fca8a444c100 100644 --- a/core/src/test/scala/unit/kafka/server/AbstractFetcherThreadTest.scala +++ b/core/src/test/scala/unit/kafka/server/AbstractFetcherThreadTest.scala @@ -21,10 +21,10 @@ import java.nio.ByteBuffer import java.util.Optional import java.util.concurrent.atomic.AtomicInteger -import com.yammer.metrics.Metrics import kafka.cluster.BrokerEndPoint import kafka.log.LogAppendInfo import kafka.message.NoCompressionCodec +import kafka.metrics.KafkaYammerMetrics import kafka.server.AbstractFetcherThread.ReplicaFetch import kafka.server.AbstractFetcherThread.ResultWithPartitions import kafka.utils.TestUtils @@ -56,7 +56,7 @@ class AbstractFetcherThreadTest { TestUtils.clearYammerMetrics() } - private def allMetricsNames: Set[String] = Metrics.defaultRegistry().allMetrics().asScala.keySet.map(_.getName) + private def allMetricsNames: Set[String] = KafkaYammerMetrics.defaultRegistry().allMetrics().asScala.keySet.map(_.getName) private def mkBatch(baseOffset: Long, leaderEpoch: Int, records: SimpleRecord*): RecordBatch = { MemoryRecords.withRecords(baseOffset, CompressionType.NONE, leaderEpoch, records: _*) @@ -89,7 +89,7 @@ class AbstractFetcherThreadTest { fetcher.shutdown() // verify that all the fetcher metrics are removed and only brokerTopicStats left - val metricNames = Metrics.defaultRegistry().allMetrics().asScala.keySet.map(_.getName).toSet + val metricNames = KafkaYammerMetrics.defaultRegistry().allMetrics().asScala.keySet.map(_.getName).toSet assertTrue(metricNames.intersect(fetcherMetrics).isEmpty) assertEquals(brokerTopicStatsMetrics, metricNames.intersect(brokerTopicStatsMetrics)) } diff --git a/core/src/test/scala/unit/kafka/server/ProduceRequestTest.scala b/core/src/test/scala/unit/kafka/server/ProduceRequestTest.scala index 1b43d02067ef0..da3f419f76c83 100644 --- a/core/src/test/scala/unit/kafka/server/ProduceRequestTest.scala +++ b/core/src/test/scala/unit/kafka/server/ProduceRequestTest.scala @@ -20,9 +20,9 @@ package kafka.server import java.nio.ByteBuffer import java.util.Properties -import com.yammer.metrics.Metrics import kafka.log.LogConfig import kafka.message.ZStdCompressionCodec +import kafka.metrics.KafkaYammerMetrics import kafka.utils.TestUtils import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.protocol.Errors @@ -40,7 +40,7 @@ import scala.collection.JavaConverters._ */ class ProduceRequestTest extends BaseRequestTest { - val metricsKeySet = Metrics.defaultRegistry.allMetrics.keySet.asScala + val metricsKeySet = KafkaYammerMetrics.defaultRegistry.allMetrics.keySet.asScala @Test def testSimpleProduceRequest(): Unit = { diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index 6edff6ea73726..d1017a705f392 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -37,9 +37,9 @@ import kafka.security.auth.{Acl, Authorizer => LegacyAuthorizer, Resource} import kafka.server._ import kafka.server.checkpoints.OffsetCheckpointFile import Implicits._ -import com.yammer.metrics.Metrics import com.yammer.metrics.core.Meter import kafka.controller.LeaderIsrAndControllerEpoch +import kafka.metrics.KafkaYammerMetrics import kafka.zk._ import org.apache.kafka.clients.CommonClientConfigs import org.apache.kafka.clients.admin.AlterConfigOp.OpType @@ -1626,7 +1626,7 @@ object TestUtils extends Logging { } def meterCount(metricName: String): Long = { - Metrics.defaultRegistry.allMetrics.asScala + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala .filterKeys(_.getMBeanName.endsWith(metricName)) .values .headOption @@ -1636,8 +1636,8 @@ object TestUtils extends Logging { } def clearYammerMetrics(): Unit = { - for (metricName <- Metrics.defaultRegistry.allMetrics.keySet.asScala) - Metrics.defaultRegistry.removeMetric(metricName) + for (metricName <- KafkaYammerMetrics.defaultRegistry.allMetrics.keySet.asScala) + KafkaYammerMetrics.defaultRegistry.removeMetric(metricName) } def stringifyTopicPartitions(partitions: Set[TopicPartition]): String = { diff --git a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala index eb9a88992fc0a..8dc9a361817d6 100644 --- a/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala +++ b/core/src/test/scala/unit/kafka/zookeeper/ZooKeeperClientTest.scala @@ -22,9 +22,9 @@ import java.util.concurrent.atomic.{AtomicBoolean, AtomicInteger} import java.util.concurrent.{ArrayBlockingQueue, ConcurrentLinkedQueue, CountDownLatch, Executors, Semaphore, TimeUnit} import scala.collection.Seq -import com.yammer.metrics.Metrics import com.yammer.metrics.core.{Gauge, Meter, MetricName} import kafka.server.KafkaConfig +import kafka.metrics.KafkaYammerMetrics import kafka.zk.ZooKeeperTestHarness import org.apache.kafka.common.security.JaasUtils import org.apache.kafka.common.utils.Time @@ -642,7 +642,7 @@ class ZooKeeperClientTest extends ZooKeeperTestHarness { @Test def testZooKeeperStateChangeRateMetrics(): Unit = { def checkMeterCount(name: String, expected: Long): Unit = { - val meter = Metrics.defaultRegistry.allMetrics.asScala.collectFirst { + val meter = KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.collectFirst { case (metricName, meter: Meter) if isExpectedMetricName(metricName, name) => meter }.getOrElse(sys.error(s"Unable to find meter with name $name")) assertEquals(s"Unexpected meter count for $name", expected, meter.count) @@ -665,7 +665,7 @@ class ZooKeeperClientTest extends ZooKeeperTestHarness { @Test def testZooKeeperSessionStateMetric(): Unit = { def gaugeValue(name: String): Option[String] = { - Metrics.defaultRegistry.allMetrics.asScala.collectFirst { + KafkaYammerMetrics.defaultRegistry.allMetrics.asScala.collectFirst { case (metricName, gauge: Gauge[_]) if isExpectedMetricName(metricName, name) => gauge.value.asInstanceOf[String] } } @@ -680,7 +680,7 @@ class ZooKeeperClientTest extends ZooKeeperTestHarness { } private def cleanMetricsRegistry(): Unit = { - val metrics = Metrics.defaultRegistry + val metrics = KafkaYammerMetrics.defaultRegistry metrics.allMetrics.keySet.asScala.foreach(metrics.removeMetric) } diff --git a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java index 58333944bdd03..63b6e5c397714 100644 --- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java +++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java @@ -681,7 +681,9 @@ private KafkaStreams(final InternalTopologyBuilder internalTopologyBuilder, final List reporters = config.getConfiguredInstances(StreamsConfig.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class, Collections.singletonMap(StreamsConfig.CLIENT_ID_CONFIG, clientId)); - reporters.add(new JmxReporter(JMX_PREFIX)); + final JmxReporter jmxReporter = new JmxReporter(JMX_PREFIX); + jmxReporter.configure(config.originals()); + reporters.add(jmxReporter); metrics = new Metrics(metricConfig, reporters, time); streamsMetrics = new StreamsMetricsImpl(metrics, clientId, config.getString(StreamsConfig.BUILT_IN_METRICS_VERSION_CONFIG));