This is an automated email from the ASF dual-hosted git repository.
SteNicholas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new 76ce74504b [CELEBORN-2380] Support JmxSink for exposing metrics as JMX
MBeans
76ce74504b is described below
commit 76ce74504ba66e022e4b7950db6af75e5e0089e2
Author: senthh <[email protected]>
AuthorDate: Fri Jul 24 10:06:42 2026 +0800
[CELEBORN-2380] Support JmxSink for exposing metrics as JMX MBeans
### What changes were proposed in this pull request?
This PR adds a new `JmxSink` metrics sink that exposes Celeborn metrics as
JMX MBeans.
- Add `JmxSink`
(`common/src/main/scala/org/apache/celeborn/common/metrics/sink/JmxSink.scala`),
backed by Dropwizard Metrics' `JmxReporter`. It follows the existing `Sink`
contract and is loaded reflectively by `MetricsSystem` via the standard
`(Properties, MetricRegistry)` constructor, so no changes to `MetricsSystem`
are required.
- Add the `io.dropwizard.metrics:metrics-jmx` dependency (which contains
`JmxReporter` in Dropwizard Metrics 4.x) to both the Maven build (`pom.xml`
dependency management + `common/pom.xml`) and the SBT build
(`project/CelebornBuild.scala`), pinned to the existing
`${codahale.metrics.version}` (4.2.25).
- Add a commented-out example to `conf/metrics.properties.template` and
`charts/celeborn/files/conf/metrics.properties` showing how to enable the sink.
- Document `JmxSink` in `docs/monitoring.md`.
- Register `metrics-jmx-4.2.25.jar` in all `dev/deps/dependencies-*`
manifests (server + client profiles) and add
`io.dropwizard.metrics:metrics-jmx` to `LICENSE-binary`, since the artifact is
now bundled in the distributions.
Enabling the sink is opt-in:
```properties
*.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink
Test Output:
<img width="1315" height="800" alt="Celeborn_jmx"
src="https://github.com/user-attachments/assets/7af729f8-7eb9-4749-af9c-e2af64d48fb8"
/>
Closes #3758 from senthh/CELEBORN-2380.
Lead-authored-by: senthh <[email protected]>
Co-authored-by: senthh <[email protected]>
Signed-off-by: Nicholas Jiang <[email protected]>
---
LICENSE-binary | 1 +
charts/celeborn/files/conf/metrics.properties | 3 +
charts/celeborn/tests/master/statefulset_test.yaml | 6 +-
charts/celeborn/tests/worker/statefulset_test.yaml | 6 +-
common/pom.xml | 4 +
.../celeborn/common/metrics/sink/JmxSink.scala | 55 +++++++++
conf/metrics.properties.template | 5 +
dev/deps/dependencies-client-flink-1.18 | 1 +
dev/deps/dependencies-client-flink-1.19 | 1 +
dev/deps/dependencies-client-flink-1.20 | 1 +
dev/deps/dependencies-client-flink-2.0 | 1 +
dev/deps/dependencies-client-flink-2.1 | 1 +
dev/deps/dependencies-client-flink-2.2 | 1 +
dev/deps/dependencies-client-flink-2.3 | 1 +
dev/deps/dependencies-client-mr | 1 +
dev/deps/dependencies-client-spark-3.0 | 1 +
dev/deps/dependencies-client-spark-3.1 | 1 +
dev/deps/dependencies-client-spark-3.2 | 1 +
dev/deps/dependencies-client-spark-3.3 | 1 +
dev/deps/dependencies-client-spark-3.4 | 1 +
dev/deps/dependencies-client-spark-3.5 | 1 +
dev/deps/dependencies-client-spark-4.0 | 1 +
dev/deps/dependencies-client-spark-4.1 | 1 +
dev/deps/dependencies-client-spark-4.2 | 1 +
dev/deps/dependencies-client-tez | 1 +
dev/deps/dependencies-server | 1 +
docs/monitoring.md | 1 +
pom.xml | 5 +
project/CelebornBuild.scala | 2 +
.../src/test/resources/metrics-jmx.properties | 5 +-
.../server/common/metrics/sink/JmxSinkSuite.scala | 132 +++++++++++++++++++++
31 files changed, 234 insertions(+), 10 deletions(-)
diff --git a/LICENSE-binary b/LICENSE-binary
index 65e4b7d04b..e1ad752dca 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 e3b521369b..7ae31b95ed 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 a95718f93a..88f6c51fc2 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 1bea9653f5..f586f13557 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 4467279159..9ff068ece1 100644
--- a/common/pom.xml
+++ b/common/pom.xml
@@ -47,6 +47,10 @@
<groupId>io.dropwizard.metrics</groupId>
<artifactId>metrics-jvm</artifactId>
</dependency>
+ <dependency>
+ <groupId>io.dropwizard.metrics</groupId>
+ <artifactId>metrics-jmx</artifactId>
+ </dependency>
<dependency>
<groupId>org.yaml</groupId>
<artifactId>snakeyaml</artifactId>
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 0000000000..d54727b20e
--- /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=<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 e3b521369b..88cd9f7d36 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=<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 d2604d91bf..637fb02a74 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 d2604d91bf..637fb02a74 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 d2604d91bf..637fb02a74 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 b06979be85..4c94188858 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 b06979be85..4c94188858 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 b06979be85..4c94188858 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 b06979be85..4c94188858 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 82919c0804..67d140abcd 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 8b12d625d7..739aa79e5b 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 5beb35ca03..a3ce2d19ed 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 3adc3ab3a3..74c4b38d2e 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 0bd8eaec4a..87e8e4e148 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 02c0c7b192..6f82bfd3ab 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 e8b51efeca..0b1d02005a 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 f5f8ddbca3..dd165d6e22 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 aa0460df43..5431f44c2c 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 ca992cea38..690866dc4f 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 9d73b59b5f..cfeb25cdc3 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 7796756d30..2ef290587d 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 5679ef548a..1aa4b3749a 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 9632fd398c..71822101ce 100644
--- a/pom.xml
+++ b/pom.xml
@@ -282,6 +282,11 @@
<artifactId>metrics-jvm</artifactId>
<version>${codahale.metrics.version}</version>
</dependency>
+ <dependency>
+ <groupId>io.dropwizard.metrics</groupId>
+ <artifactId>metrics-jmx</artifactId>
+ <version>${codahale.metrics.version}</version>
+ </dependency>
<dependency>
<groupId>org.apache.ratis</groupId>
<artifactId>ratis-common</artifactId>
diff --git a/project/CelebornBuild.scala b/project/CelebornBuild.scala
index 4aa615d969..b2688fa976 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/conf/metrics.properties.template
b/service/src/test/resources/metrics-jmx.properties
similarity index 76%
copy from conf/metrics.properties.template
copy to service/src/test/resources/metrics-jmx.properties
index e3b521369b..9ab1f1fd63 100644
--- a/conf/metrics.properties.template
+++ b/service/src/test/resources/metrics-jmx.properties
@@ -14,7 +14,4 @@
# See the License for the specific language governing permissions and
# limitations under the License.
#
-
-*.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
+*.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 0000000000..370f3e28a3
--- /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")
+ }
+}