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")
+ }
+}