From 20e9a18b7b56b4eec42acb7777c0b563bb30da77 Mon Sep 17 00:00:00 2001 From: senthh Date: Tue, 14 Jul 2026 10:56:12 +0530 Subject: [PATCH 1/8] [CELEBORN-2380] Support JmxSink for exposing metrics as JMX MBeans --- LICENSE-binary | 1 + charts/celeborn/files/conf/metrics.properties | 3 ++ common/pom.xml | 4 ++ .../common/metrics/sink/JmxSink.scala | 38 +++++++++++++++++++ conf/metrics.properties.template | 3 ++ 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-tez | 1 + dev/deps/dependencies-server | 1 + docs/monitoring.md | 1 + pom.xml | 5 +++ project/CelebornBuild.scala | 2 + 26 files changed, 75 insertions(+) create mode 100644 common/src/main/scala/org/apache/celeborn/common/metrics/sink/JmxSink.scala 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..a83a77ec6cb 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. Disabled by default; uncomment to enable. +#*.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink 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..37963c59c5b --- /dev/null +++ b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/JmxSink.scala @@ -0,0 +1,38 @@ +/* + * 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 + +private class JmxSink(val property: Properties, val registry: MetricRegistry) extends Sink { + + val reporter: JmxReporter = JmxReporter.forRegistry(registry).build() + + override def start(): Unit = { + reporter.start() + } + + override def stop(): Unit = { + reporter.stop() + } + + override def report(): Unit = {} +} diff --git a/conf/metrics.properties.template b/conf/metrics.properties.template index e3b521369b9..a83a77ec6cb 100644 --- a/conf/metrics.properties.template +++ b/conf/metrics.properties.template @@ -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. Disabled by default; uncomment to enable. +#*.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-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, From c6392e214425faa42b88acaa582310c2a03567c9 Mon Sep 17 00:00:00 2001 From: senthh Date: Tue, 14 Jul 2026 17:27:18 +0530 Subject: [PATCH 2/8] Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- .../scala/org/apache/celeborn/common/metrics/sink/JmxSink.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 index 37963c59c5b..f9af5978af8 100644 --- 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 @@ -22,7 +22,7 @@ import java.util.Properties import com.codahale.metrics.MetricRegistry import com.codahale.metrics.jmx.JmxReporter -private class JmxSink(val property: Properties, val registry: MetricRegistry) extends Sink { +class JmxSink(val property: Properties, val registry: MetricRegistry) extends Sink { val reporter: JmxReporter = JmxReporter.forRegistry(registry).build() From 84ded1be24b8b0ab8f3a46a63ed90c2da698dc29 Mon Sep 17 00:00:00 2001 From: senthh Date: Wed, 15 Jul 2026 16:40:46 +0530 Subject: [PATCH 3/8] [CELEBORN-2380] Add JmxSinkSuite covering reflection-based loading and MBean lifecycle Add a JmxSinkSuite (extending the existing sink-level test coverage) that verifies JmxSink is loaded by MetricsSystem from the metrics configuration via its reflection-based (Properties, MetricRegistry) constructor, and that start()/stop() register and unregister metric MBeans in the platform MBeanServer under the "metrics" JMX domain. --- .../src/test/resources/metrics-jmx.properties | 17 ++++ .../common/metrics/sink/JmxSinkSuite.scala | 84 +++++++++++++++++++ 2 files changed, 101 insertions(+) create mode 100644 service/src/test/resources/metrics-jmx.properties create mode 100644 service/src/test/scala/org/apache/celeborn/server/common/metrics/sink/JmxSinkSuite.scala 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..33b9c86a9cd --- /dev/null +++ b/service/src/test/scala/org/apache/celeborn/server/common/metrics/sink/JmxSinkSuite.scala @@ -0,0 +1,84 @@ +/* + * 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 javax.management.ObjectName + +import scala.collection.JavaConverters._ + +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 + // JmxReporter publishes metrics under the "metrics" JMX domain by default. + val jmxDomainPattern = new ObjectName("metrics:*") + + val beforeStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala + + // 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) + val afterStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala + val registeredByJmxSink = afterStart -- beforeStart + + metricsSystem.stop() + val afterStop = mBeanServer.queryNames(jmxDomainPattern, null).asScala + + assert(metricsSystem.sinks.exists(_.isInstanceOf[JmxSink])) + assert( + registeredByJmxSink.nonEmpty, + "JmxSink.start() should register metric MBeans under the 'metrics' domain") + assert( + registeredByJmxSink.forall(name => !afterStop.contains(name)), + "JmxSink.stop() should unregister the MBeans it registered") + } +} From fe14206436b707f62190120436c6cbbd2a140bd6 Mon Sep 17 00:00:00 2001 From: senthh Date: Wed, 15 Jul 2026 16:44:18 +0530 Subject: [PATCH 4/8] [CELEBORN-2380] Update chart conf-hash test expectations for JmxSink metrics config Adding the commented JmxSink line to charts/celeborn/files/conf/metrics.properties changes the rendered config, and therefore the celeborn.apache.org/conf-hash checksum annotation on the master and worker StatefulSets. Update the expected hashes in the helm-unittest suites accordingly. --- charts/celeborn/tests/master/statefulset_test.yaml | 6 +++--- charts/celeborn/tests/worker/statefulset_test.yaml | 6 +++--- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/charts/celeborn/tests/master/statefulset_test.yaml b/charts/celeborn/tests/master/statefulset_test.yaml index a95718f93ac..33f1a83919a 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc - 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 118d5c045d52fbbd4e8f05cf0522044373165960aa57893f645a6f9e8844dd68 + value: 069125df96d33be7e8eb2b3e61c2f33c914be0b77efff204628b1b5dc4184464 - 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..849429be256 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc - 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 118d5c045d52fbbd4e8f05cf0522044373165960aa57893f645a6f9e8844dd68 + value: 069125df96d33be7e8eb2b3e61c2f33c914be0b77efff204628b1b5dc4184464 - it: Should add extra pod annotations if `worker.annotations` is specified template: worker/statefulset.yaml From 5ceeadc71936dcc2294c9b821db458252610846b Mon Sep 17 00:00:00 2001 From: senthh Date: Thu, 16 Jul 2026 23:38:50 +0530 Subject: [PATCH 5/8] [CELEBORN-2380] Address review: configurable JMX domain, test cleanup, spark-4.2 dep - Publish MBeans under a configurable, Celeborn-specific JMX domain (default "celeborn") via `*.sink.jmx.domain`, instead of JmxReporter's global default "metrics" domain, to avoid MBean collisions with other Dropwizard users in the same JVM. - Ensure JmxSinkSuite always stops the MetricsSystem/reporter via try/finally, and add coverage for the configured-domain path and report(). - Add metrics-jmx to dev/deps/dependencies-client-spark-4.2 (new profile on main). --- .../common/metrics/sink/JmxSink.scala | 19 +++++- conf/metrics.properties.template | 2 + dev/deps/dependencies-client-spark-4.2 | 1 + .../common/metrics/sink/JmxSinkSuite.scala | 65 +++++++++++++++---- 4 files changed, 73 insertions(+), 14 deletions(-) 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 index f9af5978af8..d54727b20e7 100644 --- 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 @@ -24,7 +24,19 @@ import com.codahale.metrics.jmx.JmxReporter class JmxSink(val property: Properties, val registry: MetricRegistry) extends Sink { - val reporter: JmxReporter = JmxReporter.forRegistry(registry).build() + // 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() @@ -36,3 +48,8 @@ class JmxSink(val property: Properties, val registry: MetricRegistry) extends Si 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 a83a77ec6cb..baa0d138ea1 100644 --- a/conf/metrics.properties.template +++ b/conf/metrics.properties.template @@ -20,4 +20,6 @@ *.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-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/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 index 33b9c86a9cd..f033cd711c1 100644 --- 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 @@ -18,10 +18,13 @@ 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 @@ -59,26 +62,62 @@ class JmxSinkSuite extends CelebornFunSuite { test("test jmx sink registers and unregisters MBeans lifecycle case") { val metricsSystem = newMetricsSystem() val mBeanServer = ManagementFactory.getPlatformMBeanServer - // JmxReporter publishes metrics under the "metrics" JMX domain by default. - val jmxDomainPattern = new ObjectName("metrics:*") + // 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] - // 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) - val afterStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala - val registeredByJmxSink = afterStart -- beforeStart + 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() - metricsSystem.stop() - val afterStop = mBeanServer.queryNames(jmxDomainPattern, null).asScala + val afterStart = mBeanServer.queryNames(jmxDomainPattern, null).asScala + registeredByJmxSink = (afterStart -- beforeStart).toSet - assert(metricsSystem.sinks.exists(_.isInstanceOf[JmxSink])) - assert( - registeredByJmxSink.nonEmpty, - "JmxSink.start() should register metric MBeans under the 'metrics' domain") + 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") { + val domain = "celeborn-jmx-test-domain" + 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:*") + + try { + sink.start() + sink.report() + assert( + mBeanServer.queryNames(jmxDomainPattern, null).asScala.nonEmpty, + s"JmxSink should register MBeans under the configured '$domain' domain") + } finally { + sink.stop() + } + + assert( + mBeanServer.queryNames(jmxDomainPattern, null).isEmpty, + "JmxSink.stop() should unregister the MBeans under the configured domain") + } } From 591e1088c4a6012c75fe6189f05659f7b48da117 Mon Sep 17 00:00:00 2001 From: senthh Date: Sun, 19 Jul 2026 09:17:13 +0530 Subject: [PATCH 6/8] [CELEBORN-2380] Make configured-domain JmxSink test robust against global MBeanServer Use a unique JMX domain per run and assert on the before/after MBean delta instead of absolute presence, so the test is not affected by MBeans already present in (or leaked by a prior failed run into) the JVM-global MBeanServer. --- .../server/common/metrics/sink/JmxSinkSuite.scala | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) 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 index f033cd711c1..370f3e28a3b 100644 --- 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 @@ -94,7 +94,10 @@ class JmxSinkSuite extends CelebornFunSuite { } test("test jmx sink honors configured domain case") { - val domain = "celeborn-jmx-test-domain" + // 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() @@ -106,18 +109,24 @@ class JmxSinkSuite extends CelebornFunSuite { 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( - mBeanServer.queryNames(jmxDomainPattern, null).asScala.nonEmpty, + registeredByJmxSink.nonEmpty, s"JmxSink should register MBeans under the configured '$domain' domain") } finally { sink.stop() } + val afterStop = mBeanServer.queryNames(jmxDomainPattern, null).asScala assert( - mBeanServer.queryNames(jmxDomainPattern, null).isEmpty, + registeredByJmxSink.forall(name => !afterStop.contains(name)), "JmxSink.stop() should unregister the MBeans under the configured domain") } } From 5df972cef28c81bec79b7d3ab8c70e91730e5093 Mon Sep 17 00:00:00 2001 From: senthh Date: Mon, 20 Jul 2026 14:20:27 +0530 Subject: [PATCH 7/8] [CELEBORN-2380] Address review: format commented JmxSink example with a space Format the disabled JmxSink example as `# *.sink.jmx.class=...` (space after the comment marker) in conf/metrics.properties.template and the Helm chart metrics.properties, per review. Update the master/worker StatefulSet conf-hash test expectations for the changed chart config. --- charts/celeborn/files/conf/metrics.properties | 2 +- charts/celeborn/tests/master/statefulset_test.yaml | 6 +++--- charts/celeborn/tests/worker/statefulset_test.yaml | 6 +++--- conf/metrics.properties.template | 2 +- 4 files changed, 8 insertions(+), 8 deletions(-) diff --git a/charts/celeborn/files/conf/metrics.properties b/charts/celeborn/files/conf/metrics.properties index a83a77ec6cb..d56d03e1f23 100644 --- a/charts/celeborn/files/conf/metrics.properties +++ b/charts/celeborn/files/conf/metrics.properties @@ -20,4 +20,4 @@ *.sink.loggerSink.class=org.apache.celeborn.common.metrics.sink.LoggerSink # Expose metrics as JMX MBeans. Disabled by default; uncomment to enable. -#*.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink +# *.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 33f1a83919a..7f8bde65c04 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc + value: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a - 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc + value: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 069125df96d33be7e8eb2b3e61c2f33c914be0b77efff204628b1b5dc4184464 + value: df60eb0f4e9bb5ad6e07f4eda40842e81265c7529faaf5c4371808617d7a0515 - 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 849429be256..6f1590e169f 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc + value: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a - 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: 402448664cbf4a61ccac0b5c456cb918305572d2442eafb232e1c7ed907410dc + value: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: 069125df96d33be7e8eb2b3e61c2f33c914be0b77efff204628b1b5dc4184464 + value: df60eb0f4e9bb5ad6e07f4eda40842e81265c7529faaf5c4371808617d7a0515 - it: Should add extra pod annotations if `worker.annotations` is specified template: worker/statefulset.yaml diff --git a/conf/metrics.properties.template b/conf/metrics.properties.template index baa0d138ea1..88cd9f7d365 100644 --- a/conf/metrics.properties.template +++ b/conf/metrics.properties.template @@ -22,4 +22,4 @@ # 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 +# *.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink From bb7b4c29a3c3be9e8a29447d85168106c67ec1f1 Mon Sep 17 00:00:00 2001 From: senthh Date: Tue, 21 Jul 2026 22:46:31 +0530 Subject: [PATCH 8/8] [CELEBORN-2380] Enable JmxSink by default in Helm chart metrics config Since the JmxSink ships in 0.7.0 (the main branch version), enable it by default in the Helm chart's metrics.properties alongside the other sinks instead of leaving it commented out. Update the master/worker StatefulSet conf-hash test expectations accordingly. --- charts/celeborn/files/conf/metrics.properties | 4 ++-- charts/celeborn/tests/master/statefulset_test.yaml | 6 +++--- charts/celeborn/tests/worker/statefulset_test.yaml | 6 +++--- 3 files changed, 8 insertions(+), 8 deletions(-) diff --git a/charts/celeborn/files/conf/metrics.properties b/charts/celeborn/files/conf/metrics.properties index d56d03e1f23..7ae31b95ed7 100644 --- a/charts/celeborn/files/conf/metrics.properties +++ b/charts/celeborn/files/conf/metrics.properties @@ -19,5 +19,5 @@ *.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. -# *.sink.jmx.class=org.apache.celeborn.common.metrics.sink.JmxSink +# 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 7f8bde65c04..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: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a + 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: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a + value: 035c23a53d85eb407eecf3fe864be5e623afae38e8ecdd4c10579e76d85f5c6b - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: df60eb0f4e9bb5ad6e07f4eda40842e81265c7529faaf5c4371808617d7a0515 + 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 6f1590e169f..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: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a + 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: 96b3ee0b0e28b11a66e44bfbaeb689df5ee0063ae0c78244e8fe7015efe1a77a + value: 035c23a53d85eb407eecf3fe864be5e623afae38e8ecdd4c10579e76d85f5c6b - equal: path: spec.template.metadata.annotations["celeborn.apache.org/conf-hash"] - value: df60eb0f4e9bb5ad6e07f4eda40842e81265c7529faaf5c4371808617d7a0515 + value: d15c180987eca7e3b13dc37e3277c057a4060799ddd45cf20490ebb91da21ee1 - it: Should add extra pod annotations if `worker.annotations` is specified template: worker/statefulset.yaml