diff --git a/LICENSE-binary b/LICENSE-binary index 65e4b7d04bf..e1ad752dcac 100644 --- a/LICENSE-binary +++ b/LICENSE-binary @@ -226,6 +226,7 @@ com.zaxxer:HikariCP info.picocli:picocli io.dropwizard.metrics:metrics-core io.dropwizard.metrics:metrics-graphite +io.dropwizard.metrics:metrics-jmx io.dropwizard.metrics:metrics-jvm io.netty:netty-all io.netty:netty-buffer diff --git a/charts/celeborn/files/conf/metrics.properties b/charts/celeborn/files/conf/metrics.properties index e3b521369b9..7ae31b95ed7 100644 --- a/charts/celeborn/files/conf/metrics.properties +++ b/charts/celeborn/files/conf/metrics.properties @@ -18,3 +18,6 @@ *.sink.prometheusServlet.class=org.apache.celeborn.common.metrics.sink.PrometheusServlet *.sink.jsonServlet.class=org.apache.celeborn.common.metrics.sink.JsonServlet *.sink.loggerSink.class=org.apache.celeborn.common.metrics.sink.LoggerSink + +# Expose metrics as JMX MBeans. +*.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink diff --git a/charts/celeborn/tests/master/statefulset_test.yaml b/charts/celeborn/tests/master/statefulset_test.yaml index a95718f93ac..88f6c51fc23 100644 --- a/charts/celeborn/tests/master/statefulset_test.yaml +++ b/charts/celeborn/tests/master/statefulset_test.yaml @@ -40,7 +40,7 @@ tests: asserts: - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 7e9a27719ab1f2c1cea53e4879a807782abd48d87ec2a45774084990af6125b3 + value: 035c23a53d85eb407eecf3fe864be5e623afae38e8ecdd4c10579e76d85f5c6b - it: Should change checksum annotation when celeborn config changes template: master/statefulset.yaml @@ -50,10 +50,10 @@ tests: asserts: - notEqual: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 7e9a27719ab1f2c1cea53e4879a807782abd48d87ec2a45774084990af6125b3 + value: 035c23a53d85eb407eecf3fe864be5e623afae38e8ecdd4c10579e76d85f5c6b - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 118d5c045d52fbbd4e8f05cf0522044373165960aa57893f645a6f9e8844dd68 + value: d15c180987eca7e3b13dc37e3277c057a4060799ddd45cf20490ebb91da21ee1 - it: Should add extra pod annotations if `master.annotations` is specified template: master/statefulset.yaml diff --git a/charts/celeborn/tests/worker/statefulset_test.yaml b/charts/celeborn/tests/worker/statefulset_test.yaml index 1bea9653f51..f586f135572 100644 --- a/charts/celeborn/tests/worker/statefulset_test.yaml +++ b/charts/celeborn/tests/worker/statefulset_test.yaml @@ -40,7 +40,7 @@ tests: asserts: - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 7e9a27719ab1f2c1cea53e4879a807782abd48d87ec2a45774084990af6125b3 + value: 035c23a53d85eb407eecf3fe864be5e623afae38e8ecdd4c10579e76d85f5c6b - it: Should change checksum annotation when celeborn config changes template: worker/statefulset.yaml @@ -50,10 +50,10 @@ tests: asserts: - notEqual: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 7e9a27719ab1f2c1cea53e4879a807782abd48d87ec2a45774084990af6125b3 + value: 035c23a53d85eb407eecf3fe864be5e623afae38e8ecdd4c10579e76d85f5c6b - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 118d5c045d52fbbd4e8f05cf0522044373165960aa57893f645a6f9e8844dd68 + value: d15c180987eca7e3b13dc37e3277c057a4060799ddd45cf20490ebb91da21ee1 - it: Should add extra pod annotations if `worker.annotations` is specified template: worker/statefulset.yaml diff --git a/common/pom.xml b/common/pom.xml index 4467279159a..9ff068ece19 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -47,6 +47,10 @@ io.dropwizard.metrics metrics-jvm + + io.dropwizard.metrics + metrics-jmx + org.yaml snakeyaml diff --git a/common/src/main/scala/org/apache/celeborn/common/metrics/sink/JmxSink.scala b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/JmxSink.scala new file mode 100644 index 00000000000..d54727b20e7 --- /dev/null +++ b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/JmxSink.scala @@ -0,0 +1,55 @@ +/* + * 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 org.apache.celeborn.common.metrics.sink + +import java.util.Properties + +import com.codahale.metrics.MetricRegistry +import com.codahale.metrics.jmx.JmxReporter + +class JmxSink(val property: Properties, val registry: MetricRegistry) extends Sink { + + // Publish MBeans under a configurable, Celeborn-specific JMX domain (defaulting to + // `celeborn`) rather than JmxReporter's global default domain `metrics`. This avoids + // MBean name collisions when other components using Dropwizard metrics run in the + // same JVM. The domain can be overridden via `*.sink.jmx.domain=`. + val domain: String = + Option(property.getProperty(JmxSink.JMX_DOMAIN_KEY)) + .map(_.trim) + .filter(_.nonEmpty) + .getOrElse(JmxSink.JMX_DEFAULT_DOMAIN) + + val reporter: JmxReporter = JmxReporter.forRegistry(registry) + .inDomain(domain) + .build() + + override def start(): Unit = { + reporter.start() + } + + override def stop(): Unit = { + reporter.stop() + } + + override def report(): Unit = {} +} + +object JmxSink { + val JMX_DOMAIN_KEY = "domain" + val JMX_DEFAULT_DOMAIN = "celeborn" +} diff --git a/conf/metrics.properties.template b/conf/metrics.properties.template index e3b521369b9..88cd9f7d365 100644 --- a/conf/metrics.properties.template +++ b/conf/metrics.properties.template @@ -18,3 +18,8 @@ *.sink.prometheusServlet.class=org.apache.celeborn.common.metrics.sink.PrometheusServlet *.sink.jsonServlet.class=org.apache.celeborn.common.metrics.sink.JsonServlet *.sink.loggerSink.class=org.apache.celeborn.common.metrics.sink.LoggerSink + +# Expose metrics as JMX MBeans. Disabled by default; uncomment to enable. +# MBeans are published under the `celeborn` JMX domain by default; override with +# `*.sink.jmx.domain=` if needed. +# *.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink diff --git a/dev/deps/dependencies-client-flink-1.18 b/dev/deps/dependencies-client-flink-1.18 index d2604d91bfc..637fb02a747 100644 --- a/dev/deps/dependencies-client-flink-1.18 +++ b/dev/deps/dependencies-client-flink-1.18 @@ -36,6 +36,7 @@ lz4-java/1.10.4//lz4-java-1.10.4.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-flink-1.19 b/dev/deps/dependencies-client-flink-1.19 index d2604d91bfc..637fb02a747 100644 --- a/dev/deps/dependencies-client-flink-1.19 +++ b/dev/deps/dependencies-client-flink-1.19 @@ -36,6 +36,7 @@ lz4-java/1.10.4//lz4-java-1.10.4.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-flink-1.20 b/dev/deps/dependencies-client-flink-1.20 index d2604d91bfc..637fb02a747 100644 --- a/dev/deps/dependencies-client-flink-1.20 +++ b/dev/deps/dependencies-client-flink-1.20 @@ -36,6 +36,7 @@ lz4-java/1.10.4//lz4-java-1.10.4.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-flink-2.0 b/dev/deps/dependencies-client-flink-2.0 index b06979be85e..4c94188858c 100644 --- a/dev/deps/dependencies-client-flink-2.0 +++ b/dev/deps/dependencies-client-flink-2.0 @@ -35,6 +35,7 @@ leveldbjni-all/1.8//leveldbjni-all-1.8.jar lz4-java/1.10.4//lz4-java-1.10.4.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-flink-2.1 b/dev/deps/dependencies-client-flink-2.1 index b06979be85e..4c94188858c 100644 --- a/dev/deps/dependencies-client-flink-2.1 +++ b/dev/deps/dependencies-client-flink-2.1 @@ -35,6 +35,7 @@ leveldbjni-all/1.8//leveldbjni-all-1.8.jar lz4-java/1.10.4//lz4-java-1.10.4.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-flink-2.2 b/dev/deps/dependencies-client-flink-2.2 index b06979be85e..4c94188858c 100644 --- a/dev/deps/dependencies-client-flink-2.2 +++ b/dev/deps/dependencies-client-flink-2.2 @@ -35,6 +35,7 @@ leveldbjni-all/1.8//leveldbjni-all-1.8.jar lz4-java/1.10.4//lz4-java-1.10.4.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-flink-2.3 b/dev/deps/dependencies-client-flink-2.3 index b06979be85e..4c94188858c 100644 --- a/dev/deps/dependencies-client-flink-2.3 +++ b/dev/deps/dependencies-client-flink-2.3 @@ -35,6 +35,7 @@ leveldbjni-all/1.8//leveldbjni-all-1.8.jar lz4-java/1.10.4//lz4-java-1.10.4.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-mr b/dev/deps/dependencies-client-mr index 82919c08041..67d140abcd2 100644 --- a/dev/deps/dependencies-client-mr +++ b/dev/deps/dependencies-client-mr @@ -138,6 +138,7 @@ lz4-java/1.10.4//lz4-java-1.10.4.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar mssql-jdbc/6.2.1.jre7//mssql-jdbc-6.2.1.jre7.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-3.0 b/dev/deps/dependencies-client-spark-3.0 index 8b12d625d79..739aa79e5b4 100644 --- a/dev/deps/dependencies-client-spark-3.0 +++ b/dev/deps/dependencies-client-spark-3.0 @@ -36,6 +36,7 @@ lz4-java/1.7.1//lz4-java-1.7.1.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-3.1 b/dev/deps/dependencies-client-spark-3.1 index 5beb35ca037..a3ce2d19ed0 100644 --- a/dev/deps/dependencies-client-spark-3.1 +++ b/dev/deps/dependencies-client-spark-3.1 @@ -36,6 +36,7 @@ lz4-java/1.7.1//lz4-java-1.7.1.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-3.2 b/dev/deps/dependencies-client-spark-3.2 index 3adc3ab3a30..74c4b38d2e7 100644 --- a/dev/deps/dependencies-client-spark-3.2 +++ b/dev/deps/dependencies-client-spark-3.2 @@ -36,6 +36,7 @@ lz4-java/1.7.1//lz4-java-1.7.1.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-3.3 b/dev/deps/dependencies-client-spark-3.3 index 0bd8eaec4a6..87e8e4e1489 100644 --- a/dev/deps/dependencies-client-spark-3.3 +++ b/dev/deps/dependencies-client-spark-3.3 @@ -36,6 +36,7 @@ lz4-java/1.8.0//lz4-java-1.8.0.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-3.4 b/dev/deps/dependencies-client-spark-3.4 index 02c0c7b1928..6f82bfd3abf 100644 --- a/dev/deps/dependencies-client-spark-3.4 +++ b/dev/deps/dependencies-client-spark-3.4 @@ -36,6 +36,7 @@ lz4-java/1.8.0//lz4-java-1.8.0.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-3.5 b/dev/deps/dependencies-client-spark-3.5 index e8b51efecac..0b1d02005ad 100644 --- a/dev/deps/dependencies-client-spark-3.5 +++ b/dev/deps/dependencies-client-spark-3.5 @@ -36,6 +36,7 @@ lz4-java/1.8.0//lz4-java-1.8.0.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-4.0 b/dev/deps/dependencies-client-spark-4.0 index f5f8ddbca37..dd165d6e22e 100644 --- a/dev/deps/dependencies-client-spark-4.0 +++ b/dev/deps/dependencies-client-spark-4.0 @@ -35,6 +35,7 @@ leveldbjni-all/1.8//leveldbjni-all-1.8.jar lz4-java/1.8.0//lz4-java-1.8.0.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-4.1 b/dev/deps/dependencies-client-spark-4.1 index aa0460df433..5431f44c2c8 100644 --- a/dev/deps/dependencies-client-spark-4.1 +++ b/dev/deps/dependencies-client-spark-4.1 @@ -35,6 +35,7 @@ leveldbjni-all/1.8//leveldbjni-all-1.8.jar lz4-java/1.8.0//lz4-java-1.8.0.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-spark-4.2 b/dev/deps/dependencies-client-spark-4.2 index ca992cea38d..690866dc4f0 100644 --- a/dev/deps/dependencies-client-spark-4.2 +++ b/dev/deps/dependencies-client-spark-4.2 @@ -35,6 +35,7 @@ leveldbjni-all/1.8//leveldbjni-all-1.8.jar lz4-java/1.11.0//lz4-java-1.11.0.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-client-tez b/dev/deps/dependencies-client-tez index 9d73b59b5fd..cfeb25cdc3e 100644 --- a/dev/deps/dependencies-client-tez +++ b/dev/deps/dependencies-client-tez @@ -111,6 +111,7 @@ lz4-java/1.10.4//lz4-java-1.10.4.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar netty-all/4.2.10.Final//netty-all-4.2.10.Final.jar netty-buffer/4.2.10.Final//netty-buffer-4.2.10.Final.jar diff --git a/dev/deps/dependencies-server b/dev/deps/dependencies-server index 7796756d305..2ef290587d6 100644 --- a/dev/deps/dependencies-server +++ b/dev/deps/dependencies-server @@ -83,6 +83,7 @@ lz4-java/1.10.4//lz4-java-1.10.4.jar maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar metrics-core/4.2.25//metrics-core-4.2.25.jar metrics-graphite/4.2.25//metrics-graphite-4.2.25.jar +metrics-jmx/4.2.25//metrics-jmx-4.2.25.jar metrics-jvm/4.2.25//metrics-jvm-4.2.25.jar mimepull/1.9.15//mimepull-1.9.15.jar mybatis/3.5.15//mybatis-3.5.15.jar diff --git a/docs/monitoring.md b/docs/monitoring.md index 5679ef548a5..1aa4b3749a8 100644 --- a/docs/monitoring.md +++ b/docs/monitoring.md @@ -47,6 +47,7 @@ Each instance can report to zero or more _sinks_. Sinks are contained in the * `PrometheusServlet`: Adds a servlet within the existing Celeborn REST API to serve metrics data in Prometheus format. * `JsonServlet`: Adds a servlet within the existing Celeborn REST API to serve metrics data in JSON format. * `GraphiteSink`: Sends metrics to a Graphite node. +* `JmxSink`: Registers metrics for viewing in a JMX console. * `LoggerSink`: Scrape metrics periodically and output them to the logger files if you have enabled `celeborn.metrics.loggerSink.output.enabled`. This is used as safety valve to make sure the metrics data won't exist in the memory for a long time. If you don't have a metrics collector to diff --git a/pom.xml b/pom.xml index 9632fd398ca..71822101cee 100644 --- a/pom.xml +++ b/pom.xml @@ -282,6 +282,11 @@ metrics-jvm ${codahale.metrics.version} + + io.dropwizard.metrics + metrics-jmx + ${codahale.metrics.version} + org.apache.ratis ratis-common diff --git a/project/CelebornBuild.scala b/project/CelebornBuild.scala index 4aa615d969a..b2688fa9763 100644 --- a/project/CelebornBuild.scala +++ b/project/CelebornBuild.scala @@ -142,6 +142,7 @@ object Dependencies { val ioDropwizardMetricsGraphite = "io.dropwizard.metrics" % "metrics-graphite" % metricsVersion excludeAll ( ExclusionRule("com.rabbitmq", "amqp-client")) val ioDropwizardMetricsJvm = "io.dropwizard.metrics" % "metrics-jvm" % metricsVersion + val ioDropwizardMetricsJmx = "io.dropwizard.metrics" % "metrics-jmx" % metricsVersion val ioNetty = "io.netty" % "netty-all" % nettyVersion excludeAll( ExclusionRule("io.netty", "netty-codec-haproxy"), ExclusionRule("io.netty", "netty-codec-memcache"), @@ -679,6 +680,7 @@ object CelebornCommon { Dependencies.ioDropwizardMetricsCore, Dependencies.ioDropwizardMetricsGraphite, Dependencies.ioDropwizardMetricsJvm, + Dependencies.ioDropwizardMetricsJmx, Dependencies.ioNetty, Dependencies.ioNettyEpollLinuxX8664, Dependencies.ioNettyEpollLinuxAarch64, diff --git a/service/src/test/resources/metrics-jmx.properties b/service/src/test/resources/metrics-jmx.properties new file mode 100644 index 00000000000..9ab1f1fd638 --- /dev/null +++ b/service/src/test/resources/metrics-jmx.properties @@ -0,0 +1,17 @@ +# +# 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. +# +*.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink diff --git a/service/src/test/scala/org/apache/celeborn/server/common/metrics/sink/JmxSinkSuite.scala b/service/src/test/scala/org/apache/celeborn/server/common/metrics/sink/JmxSinkSuite.scala new file mode 100644 index 00000000000..370f3e28a3b --- /dev/null +++ b/service/src/test/scala/org/apache/celeborn/server/common/metrics/sink/JmxSinkSuite.scala @@ -0,0 +1,132 @@ +/* + * 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 org.apache.celeborn.server.common.metrics.sink + +import java.lang.management.ManagementFactory +import java.util.Properties +import javax.management.ObjectName + +import scala.collection.JavaConverters._ + +import com.codahale.metrics.MetricRegistry + +import org.apache.celeborn.CelebornFunSuite +import org.apache.celeborn.common.CelebornConf +import org.apache.celeborn.common.metrics.MetricsSystem +import org.apache.celeborn.common.metrics.sink.JmxSink +import org.apache.celeborn.common.metrics.source.JVMSource +import org.apache.celeborn.common.network.TestHelper + +class JmxSinkSuite extends CelebornFunSuite { + + private def newMetricsSystem(): MetricsSystem = { + val celebornConf = new CelebornConf() + celebornConf + .set(CelebornConf.METRICS_ENABLED.key, "true") + .set( + CelebornConf.METRICS_CONF.key, + TestHelper.getResourceAsAbsolutePath("/metrics-jmx.properties")) + val metricsSystem = MetricsSystem.createMetricsSystem("test", celebornConf) + metricsSystem.registerSource(new JVMSource(celebornConf, "test")) + metricsSystem + } + + test("test load jmx sink case") { + val metricsSystem = newMetricsSystem() + metricsSystem.start(true) + + try { + // JmxSink has no dedicated branch in MetricsSystem.registerSinks, so it must be + // instantiated reflectively via its (Properties, MetricRegistry) constructor. + assert(metricsSystem.sinks.exists(_.isInstanceOf[JmxSink])) + } finally { + metricsSystem.stop() + } + } + + test("test jmx sink registers and unregisters MBeans lifecycle case") { + val metricsSystem = newMetricsSystem() + val mBeanServer = ManagementFactory.getPlatformMBeanServer + // JmxSink publishes metrics under the Celeborn-specific default domain. + val jmxDomainPattern = new ObjectName(s"${JmxSink.JMX_DEFAULT_DOMAIN}:*") + + val beforeStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala + var registeredByJmxSink = Set.empty[ObjectName] + + try { + // start() runs registerSinks() (reflection-based loading) followed by Sink.start(), + // which makes the JmxSink's JmxReporter register the registry metrics as MBeans. + metricsSystem.start(true) + // report() is a no-op for JmxSink (JmxReporter reports on registration) but must be safe. + metricsSystem.report() + + val afterStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala + registeredByJmxSink = (afterStart -- beforeStart).toSet + + assert(metricsSystem.sinks.exists(_.isInstanceOf[JmxSink])) + assert( + registeredByJmxSink.nonEmpty, + s"JmxSink.start() should register metric MBeans under the " + + s"'${JmxSink.JMX_DEFAULT_DOMAIN}' domain") + } finally { + metricsSystem.stop() + } + + val afterStop = mBeanServer.queryNames(jmxDomainPattern, null).asScala + assert( + registeredByJmxSink.forall(name => !afterStop.contains(name)), + "JmxSink.stop() should unregister the MBeans it registered") + } + + test("test jmx sink honors configured domain case") { + // Use a unique domain per run so the assertions are robust against MBeans that may + // already exist in (or were leaked by a prior failed run into) the JVM-global + // MBeanServer, and assert on the before/after delta rather than absolute presence. + val domain = s"celeborn-jmx-test-${java.util.UUID.randomUUID()}" + val properties = new Properties() + properties.setProperty(JmxSink.JMX_DOMAIN_KEY, domain) + val registry = new MetricRegistry() + registry.counter("test-counter").inc() + + val sink = new JmxSink(properties, registry) + assert(sink.domain == domain) + + val mBeanServer = ManagementFactory.getPlatformMBeanServer + val jmxDomainPattern = new ObjectName(s"$domain:*") + + val beforeStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala + var registeredByJmxSink = Set.empty[ObjectName] + + try { + sink.start() + sink.report() + val afterStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala + registeredByJmxSink = (afterStart -- beforeStart).toSet + assert( + registeredByJmxSink.nonEmpty, + s"JmxSink should register MBeans under the configured '$domain' domain") + } finally { + sink.stop() + } + + val afterStop = mBeanServer.queryNames(jmxDomainPattern, null).asScala + assert( + registeredByJmxSink.forall(name => !afterStop.contains(name)), + "JmxSink.stop() should unregister the MBeans under the configured domain") + } +}