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

Reply via email to