This is an automated email from the ASF dual-hosted git repository.

rexxiong pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git


The following commit(s) were added to refs/heads/main by this push:
     new 04a1e9020 [CELEBORN-1122] Metrics supports json format
04a1e9020 is described below

commit 04a1e9020769764d1462d90fc7278a7949caa60e
Author: qinrui <[email protected]>
AuthorDate: Wed Dec 6 09:24:28 2023 +0800

    [CELEBORN-1122] Metrics supports json format
    
    ### What changes were proposed in this pull request?
    If the user does not use prometheus to collect monitoring metrics, but 
rather some other ones. Using metrics in JSON format would be more 
user-friendly.The PR supports JSON format for metrics.
    
    ### Why are the changes needed?
    Ditto.
    
    ### Does this PR introduce _any_ user-facing change?
    Metrics supports JSON format
    
    ### How was this patch tested?
    Cluster test.
    
    Closes #2089 from suizhe007/CELEBORN-1122.
    
    Authored-by: qinrui <[email protected]>
    Signed-off-by: Shuang <[email protected]>
---
 LICENSE-binary                                     |   5 +
 common/pom.xml                                     |  16 +
 .../org/apache/celeborn/common/CelebornConf.scala  |  18 ++
 .../celeborn/common/metrics/MetricsSystem.scala    |  40 ++-
 ...ometheusServlet.scala => AbstractServlet.scala} |  38 +--
 .../celeborn/common/metrics/sink/JsonServlet.scala | 355 +++++++++++++++++++++
 .../common/metrics/sink/PrometheusServlet.scala    |  28 +-
 .../common/metrics/source/AbstractSource.scala     |   8 +-
 conf/metrics.properties.template                   |   1 +
 dev/deps/dependencies-client-flink-1.14            |   5 +
 dev/deps/dependencies-client-flink-1.15            |   5 +
 dev/deps/dependencies-client-flink-1.17            |   5 +
 dev/deps/dependencies-client-flink-1.18            |   5 +
 dev/deps/dependencies-client-mr                    |   7 +-
 dev/deps/dependencies-client-spark-2.4             |   5 +
 dev/deps/dependencies-client-spark-3.0             |   5 +
 dev/deps/dependencies-client-spark-3.1             |   5 +
 dev/deps/dependencies-client-spark-3.2             |   5 +
 dev/deps/dependencies-client-spark-3.3             |   5 +
 dev/deps/dependencies-client-spark-3.4             |   5 +
 dev/deps/dependencies-client-spark-3.5             |   5 +
 dev/deps/dependencies-server                       |   5 +
 docs/configuration/metrics.md                      |   2 +
 .../celeborn/service/deploy/master/Master.scala    |   3 +-
 pom.xml                                            |  21 ++
 project/CelebornBuild.scala                        |  10 +-
 .../celeborn/server/common/HttpService.scala       |   2 +-
 .../server/common/http/HttpRequestHandler.scala    |  26 +-
 .../celeborn/service/deploy/worker/Worker.scala    |   2 +-
 29 files changed, 565 insertions(+), 77 deletions(-)

diff --git a/LICENSE-binary b/LICENSE-binary
index 61838e78a..c2fd1d57c 100644
--- a/LICENSE-binary
+++ b/LICENSE-binary
@@ -280,6 +280,11 @@ org.scala-lang:scala-reflect
 org.slf4j:jcl-over-slf4j
 org.yaml:snakeyaml
 org.rocksdb:rocksdbjni
+com.fasterxml.jackson.module:jackson-module-scala
+com.fasterxml.jackson.core:jackson-databind
+com.fasterxml.jackson.core:jackson-annotations
+com.fasterxml.jackson.core:jackson-core
+com.thoughtworks.paranamer:paranamer
 
 
 
------------------------------------------------------------------------------------
diff --git a/common/pom.xml b/common/pom.xml
index 7af0e9b9e..6419ce372 100644
--- a/common/pom.xml
+++ b/common/pom.xml
@@ -111,6 +111,22 @@
       <groupId>org.roaringbitmap</groupId>
       <artifactId>RoaringBitmap</artifactId>
     </dependency>
+    <dependency>
+      <groupId>com.fasterxml.jackson.module</groupId>
+      <artifactId>jackson-module-scala_${scala.binary.version}</artifactId>
+    </dependency>
+    <dependency>
+      <groupId>com.fasterxml.jackson.core</groupId>
+      <artifactId>jackson-databind</artifactId>
+    </dependency>
+    <dependency>
+      <groupId>com.fasterxml.jackson.core</groupId>
+      <artifactId>jackson-annotations</artifactId>
+    </dependency>
+    <dependency>
+      <groupId>com.fasterxml.jackson.core</groupId>
+      <artifactId>jackson-core</artifactId>
+    </dependency>
     <!-- Test dependencies -->
     <dependency>
       <groupId>org.mockito</groupId>
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala 
b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
index e56d3aa1a..0960529e6 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -721,6 +721,7 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable 
with Logging with Se
   def metricsAppTopDiskUsageInterval: Long = 
get(METRICS_APP_TOP_DISK_USAGE_INTERVAL)
   def metricsWorkerForceAppendPauseSpentTimeThreshold: Int =
     get(METRICS_WORKER_PAUSE_SPENT_TIME_FORCE_APPEND_THRESHOLD)
+  def metricsJsonPrettyEnabled: Boolean = get(METRICS_JSON_PRETTY_ENABLED)
 
   // //////////////////////////////////////////////////////
   //                      Quota                         //
@@ -3900,6 +3901,23 @@ object CelebornConf extends Logging {
       .checkValue(path => path.startsWith("/"), "Context path must start with 
'/'")
       .createWithDefault("/metrics/prometheus")
 
+  val METRICS_JSON_PATH: ConfigEntry[String] =
+    buildConf("celeborn.metrics.json.path")
+      .categories("metrics")
+      .doc("URI context path of json metrics HTTP server.")
+      .version("0.4.0")
+      .stringConf
+      .checkValue(path => path.startsWith("/"), "Context path must start with 
'/'")
+      .createWithDefault("/metrics/json")
+
+  val METRICS_JSON_PRETTY_ENABLED: ConfigEntry[Boolean] =
+    buildConf("celeborn.metrics.json.pretty.enabled")
+      .categories("metrics")
+      .doc("When true, view metrics in json pretty format")
+      .version("0.4.0")
+      .booleanConf
+      .createWithDefault(true)
+
   val QUOTA_ENABLED: ConfigEntry[Boolean] =
     buildConf("celeborn.quota.enabled")
       .categories("quota")
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/metrics/MetricsSystem.scala 
b/common/src/main/scala/org/apache/celeborn/common/metrics/MetricsSystem.scala
index 5fdd296a5..82311bea9 100644
--- 
a/common/src/main/scala/org/apache/celeborn/common/metrics/MetricsSystem.scala
+++ 
b/common/src/main/scala/org/apache/celeborn/common/metrics/MetricsSystem.scala
@@ -27,30 +27,34 @@ import scala.util.matching.Regex
 import com.codahale.metrics.{Metric, MetricFilter, MetricRegistry}
 
 import org.apache.celeborn.common.CelebornConf
+import org.apache.celeborn.common.CelebornConf.{METRICS_JSON_PATH, 
METRICS_PROMETHEUS_PATH}
 import org.apache.celeborn.common.internal.Logging
-import org.apache.celeborn.common.metrics.sink.{PrometheusHttpRequestHandler, 
PrometheusServlet, Sink}
+import org.apache.celeborn.common.metrics.sink.{JsonServlet, 
PrometheusServlet, ServletHttpRequestHandler, Sink}
 import org.apache.celeborn.common.metrics.source.Source
 import org.apache.celeborn.common.util.Utils
 
 class MetricsSystem(
     val instance: String,
-    conf: CelebornConf,
-    val servletPath: String) extends Logging {
+    conf: CelebornConf) extends Logging {
   private[this] val metricsConfig = new MetricsConfig(conf)
 
   private val sinks = new ArrayBuffer[Sink]
   private val sources = new CopyOnWriteArrayList[Source]
   private val registry = new MetricRegistry()
+  private val prometheusServletPath = conf.get(METRICS_PROMETHEUS_PATH)
+  private val jsonServletPath = conf.get(METRICS_JSON_PATH)
 
   private var prometheusServlet: Option[PrometheusServlet] = None
+  private var jsonServlet: Option[JsonServlet] = None
 
   var running: Boolean = false
 
   metricsConfig.initialize()
 
-  def getPrometheusHandler: PrometheusHttpRequestHandler = {
+  def getServletHandlers: Array[ServletHttpRequestHandler] = {
     require(running, "Can only call getServletHandlers on a running 
MetricsSystem")
-    prometheusServlet.map(_.getHandler(conf)).orNull
+    prometheusServlet.map(_.getHandlers(conf)).getOrElse(Array()) ++
+      jsonServlet.map(_.getHandlers(conf)).getOrElse(Array())
   }
 
   def start(registerStaticSources: Boolean = true) {
@@ -132,8 +136,25 @@ class MetricsSystem(
                 classOf[MetricRegistry],
                 classOf[Seq[Source]],
                 classOf[String])
-              .newInstance(kv._2, registry, sources.asScala, servletPath)
-            prometheusServlet = Some(servlet.asInstanceOf[PrometheusServlet])
+            prometheusServlet = Some(servlet.newInstance(
+              kv._2,
+              registry,
+              sources.asScala,
+              prometheusServletPath).asInstanceOf[PrometheusServlet])
+          } else if (kv._1 == "jsonServlet") {
+            val servlet = Utils.classForName(classPath)
+              .getConstructor(
+                classOf[Properties],
+                classOf[MetricRegistry],
+                classOf[Seq[Source]],
+                classOf[String],
+                classOf[Boolean])
+            jsonServlet = Some(servlet.newInstance(
+              kv._2,
+              registry,
+              sources.asScala,
+              jsonServletPath,
+              
conf.metricsJsonPrettyEnabled.asInstanceOf[Object]).asInstanceOf[JsonServlet])
           } else {
             val sink = Utils.classForName(classPath)
               .getConstructor(classOf[Properties], classOf[MetricRegistry])
@@ -171,8 +192,7 @@ object MetricsSystem {
 
   def createMetricsSystem(
       instance: String,
-      conf: CelebornConf,
-      servletPath: String): MetricsSystem = {
-    new MetricsSystem(instance, conf, servletPath)
+      conf: CelebornConf): MetricsSystem = {
+    new MetricsSystem(instance, conf)
   }
 }
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/metrics/sink/PrometheusServlet.scala
 
b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/AbstractServlet.scala
similarity index 62%
copy from 
common/src/main/scala/org/apache/celeborn/common/metrics/sink/PrometheusServlet.scala
copy to 
common/src/main/scala/org/apache/celeborn/common/metrics/sink/AbstractServlet.scala
index a77a8375e..325826164 100644
--- 
a/common/src/main/scala/org/apache/celeborn/common/metrics/sink/PrometheusServlet.scala
+++ 
b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/AbstractServlet.scala
@@ -14,28 +14,20 @@
  * 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 io.netty.channel.ChannelHandler.Sharable
-
 import org.apache.celeborn.common.CelebornConf
 import org.apache.celeborn.common.internal.Logging
 import org.apache.celeborn.common.metrics.source.Source
 
-class PrometheusServlet(
-    val property: Properties,
-    val registry: MetricRegistry,
-    val sources: Seq[Source],
-    val servletPath: String) extends Sink with Logging {
-
-  def getHandler(conf: CelebornConf): PrometheusHttpRequestHandler = {
-    new PrometheusHttpRequestHandler(servletPath, this)
+abstract class AbstractServlet(sources: Seq[Source]) extends Sink with Logging 
{
+  def getHandlers(conf: CelebornConf): Array[ServletHttpRequestHandler] = {
+    Array[ServletHttpRequestHandler](
+      createHttpRequestHandler())
   }
 
+  def createHttpRequestHandler(): ServletHttpRequestHandler
+
   def getMetricsSnapshot: String = {
     sources.map(_.getMetrics).mkString
   }
@@ -47,16 +39,10 @@ class PrometheusServlet(
   override def report(): Unit = {}
 }
 
-@Sharable
-class PrometheusHttpRequestHandler(
-    path: String,
-    prometheusServlet: PrometheusServlet) extends Logging {
-
-  def handleRequest(uri: String): String = {
-    if (uri == path) {
-      prometheusServlet.getMetricsSnapshot
-    } else {
-      s"Unknown path $uri!"
-    }
-  }
+abstract class ServletHttpRequestHandler(path: String) extends Logging {
+
+  def handleRequest(uri: String): String
+
+  def getServletPath(): String = path
+
 }
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/metrics/sink/JsonServlet.scala
 
b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/JsonServlet.scala
new file mode 100644
index 000000000..7a2b8b52c
--- /dev/null
+++ 
b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/JsonServlet.scala
@@ -0,0 +1,355 @@
+/*
+ * 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 scala.collection.mutable.ArrayBuffer
+
+import com.codahale.metrics.MetricRegistry
+import com.fasterxml.jackson.databind.ObjectMapper
+import com.fasterxml.jackson.module.scala.{ClassTagExtensions, 
DefaultScalaModule}
+import io.netty.channel.ChannelHandler.Sharable
+
+import org.apache.celeborn.common.metrics.{CelebornHistogram, CelebornTimer, 
ResettableSlidingWindowReservoir}
+import org.apache.celeborn.common.metrics.source.{AbstractSource, 
NamedCounter, NamedGauge, NamedHistogram, NamedTimer, Source}
+
+object JsonConverter {
+  val mapper = new ObjectMapper() with ClassTagExtensions
+  mapper.registerModule(DefaultScalaModule)
+
+  def toJson(value: Any): String = {
+    mapper.writeValueAsString(value)
+  }
+
+  def toPrettyJson(value: Any): String = {
+    mapper.writerWithDefaultPrettyPrinter().writeValueAsString(value)
+  }
+}
+
+case class MetricData(
+    name: String,
+    value: Any,
+    timestampMs: Long,
+    labelNames: ArrayBuffer[String],
+    labelValues: ArrayBuffer[String])
+
+class JsonServlet(
+    val property: Properties,
+    val registry: MetricRegistry,
+    val sources: Seq[Source],
+    val servletPath: String,
+    val jsonPrettyEnabled: Boolean) extends AbstractServlet(sources) {
+
+  override def getMetricsSnapshot: String = {
+    val metricDatas = new ArrayBuffer[MetricData]
+    try {
+      sources.map(source => metricDatas ++= getMetrics(source))
+      if (jsonPrettyEnabled) {
+        JsonConverter.toPrettyJson(metricDatas.map(_.asInstanceOf[Any]))
+      } else {
+        JsonConverter.toJson(metricDatas.map(_.asInstanceOf[Any]))
+      }
+    } catch {
+      case e: Throwable =>
+        logError("failed to get json data for metrics", e)
+        JsonConverter.toJson(new ArrayBuffer())
+    }
+  }
+
+  override def createHttpRequestHandler(): ServletHttpRequestHandler = {
+    new JsonHttpRequestHandler(servletPath, this)
+  }
+
+  override def stop(): Unit = {}
+
+  def getMetrics(source: Source): ArrayBuffer[MetricData] = {
+    val metricDatas = new ArrayBuffer[MetricData]
+    val absSource = source.asInstanceOf[AbstractSource]
+    absSource.counters().foreach(c => recordCounter(absSource, c, metricDatas))
+    absSource.gauges().foreach(g => recordGauge(absSource, g, metricDatas))
+    absSource.histograms().foreach(h => {
+      recordHistogram(absSource, h, metricDatas)
+      h.asInstanceOf[CelebornHistogram].reservoir
+        .asInstanceOf[ResettableSlidingWindowReservoir].reset()
+    })
+    absSource.timers().foreach(t => {
+      recordTimer(absSource, t, metricDatas)
+      t.timer.asInstanceOf[CelebornTimer].reservoir
+        .asInstanceOf[ResettableSlidingWindowReservoir].reset()
+    })
+    metricDatas
+  }
+
+  def recordCounter(
+      absSource: AbstractSource,
+      nc: NamedCounter,
+      metricDatas: ArrayBuffer[MetricData]): Unit = {
+    val timestamp = System.currentTimeMillis
+    val labelNames = new ArrayBuffer[String]
+    val labelValues = new ArrayBuffer[String]
+    nc.labels.map { case (k, v) =>
+      labelNames += k
+      labelValues += v
+    }
+    val metricData =
+      MetricData(
+        nc.name,
+        nc.counter.getCount,
+        timestamp,
+        labelNames,
+        labelValues)
+    updateInnerMetrics(absSource, metricData, metricDatas)
+  }
+
+  def recordGauge(
+      absSource: AbstractSource,
+      ng: NamedGauge[_],
+      metricDatas: ArrayBuffer[MetricData]): Unit = {
+    val timestamp = System.currentTimeMillis
+    val labelNames = new ArrayBuffer[String]
+    val labelValues = new ArrayBuffer[String]
+    ng.labels.map { case (k, v) =>
+      labelNames += k
+      labelValues += v
+    }
+    val metricData = MetricData(ng.name, ng.gauge.getValue, timestamp, 
labelNames, labelValues)
+    updateInnerMetrics(absSource, metricData, metricDatas)
+  }
+
+  def recordHistogram(
+      absSource: AbstractSource,
+      nh: NamedHistogram,
+      metricDatas: ArrayBuffer[MetricData]): Unit = {
+    val timestamp = System.currentTimeMillis
+    val labelNames = new ArrayBuffer[String]
+    val labelValues = new ArrayBuffer[String]
+    nh.labels.map { case (k, v) =>
+      labelNames += k
+      labelValues += v
+    }
+    val snapshot = nh.histogram.getSnapshot
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_Count",
+        nh.histogram.getCount,
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_Max",
+        absSource.reportNanosAsMills(snapshot.getMax),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_Mean",
+        absSource.reportNanosAsMills(snapshot.getMean),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_Min",
+        absSource.reportNanosAsMills(snapshot.getMin),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_50thPercentile",
+        absSource.reportNanosAsMills(snapshot.getMedian),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_75thPercentile",
+        absSource.reportNanosAsMills(snapshot.get75thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_95thPercentile",
+        absSource.reportNanosAsMills(snapshot.get95thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_98thPercentile",
+        absSource.reportNanosAsMills(snapshot.get98thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_99thPercentile",
+        absSource.reportNanosAsMills(snapshot.get99thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nh.name}_999thPercentile",
+        absSource.reportNanosAsMills(snapshot.get999thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+  }
+
+  def recordTimer(
+      absSource: AbstractSource,
+      nt: NamedTimer,
+      metricDatas: ArrayBuffer[MetricData]): Unit = {
+    val timestamp = System.currentTimeMillis
+    val labelNames = new ArrayBuffer[String]
+    val labelValues = new ArrayBuffer[String]
+    nt.labels.map { case (k, v) =>
+      labelNames += k
+      labelValues += v
+    }
+    val snapshot = nt.timer.getSnapshot
+    updateInnerMetrics(
+      absSource,
+      MetricData(s"${nt.name}_Count", nt.timer.getCount, timestamp, 
labelNames, labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_Max",
+        absSource.reportNanosAsMills(snapshot.getMax),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_Mean",
+        absSource.reportNanosAsMills(snapshot.getMean),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_Min",
+        absSource.reportNanosAsMills(snapshot.getMin),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_50thPercentile",
+        absSource.reportNanosAsMills(snapshot.getMedian),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_75thPercentile",
+        absSource.reportNanosAsMills(snapshot.get75thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_95thPercentile",
+        absSource.reportNanosAsMills(snapshot.get95thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_98thPercentile",
+        absSource.reportNanosAsMills(snapshot.get98thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_99thPercentile",
+        absSource.reportNanosAsMills(snapshot.get99thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+    updateInnerMetrics(
+      absSource,
+      MetricData(
+        s"${nt.name}_999thPercentile",
+        absSource.reportNanosAsMills(snapshot.get999thPercentile),
+        timestamp,
+        labelNames,
+        labelValues),
+      metricDatas)
+  }
+
+  private def updateInnerMetrics(
+      absSource: AbstractSource,
+      metricData: MetricData,
+      metricDatas: ArrayBuffer[MetricData]): Unit = {
+    if (metricDatas.size < absSource.metricsCapacity) {
+      metricDatas += metricData
+    }
+  }
+}
+
+@Sharable
+class JsonHttpRequestHandler(path: String, jsonServlet: JsonServlet)
+  extends ServletHttpRequestHandler(path) {
+
+  override def handleRequest(uri: String): String = {
+    jsonServlet.getMetricsSnapshot
+  }
+}
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/metrics/sink/PrometheusServlet.scala
 
b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/PrometheusServlet.scala
index a77a8375e..4b4548421 100644
--- 
a/common/src/main/scala/org/apache/celeborn/common/metrics/sink/PrometheusServlet.scala
+++ 
b/common/src/main/scala/org/apache/celeborn/common/metrics/sink/PrometheusServlet.scala
@@ -22,41 +22,25 @@ import java.util.Properties
 import com.codahale.metrics.MetricRegistry
 import io.netty.channel.ChannelHandler.Sharable
 
-import org.apache.celeborn.common.CelebornConf
-import org.apache.celeborn.common.internal.Logging
 import org.apache.celeborn.common.metrics.source.Source
 
 class PrometheusServlet(
     val property: Properties,
     val registry: MetricRegistry,
     val sources: Seq[Source],
-    val servletPath: String) extends Sink with Logging {
+    val servletPath: String) extends AbstractServlet(sources) {
 
-  def getHandler(conf: CelebornConf): PrometheusHttpRequestHandler = {
+  override def createHttpRequestHandler(): ServletHttpRequestHandler = {
     new PrometheusHttpRequestHandler(servletPath, this)
   }
-
-  def getMetricsSnapshot: String = {
-    sources.map(_.getMetrics).mkString
-  }
-
-  override def start(): Unit = {}
-
-  override def stop(): Unit = {}
-
-  override def report(): Unit = {}
 }
 
 @Sharable
 class PrometheusHttpRequestHandler(
     path: String,
-    prometheusServlet: PrometheusServlet) extends Logging {
-
-  def handleRequest(uri: String): String = {
-    if (uri == path) {
-      prometheusServlet.getMetricsSnapshot
-    } else {
-      s"Unknown path $uri!"
-    }
+    prometheusServlet: PrometheusServlet) extends 
ServletHttpRequestHandler(path) {
+
+  override def handleRequest(uri: String): String = {
+    prometheusServlet.getMetricsSnapshot
   }
 }
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/metrics/source/AbstractSource.scala
 
b/common/src/main/scala/org/apache/celeborn/common/metrics/source/AbstractSource.scala
index 4f53cb70b..5ff032c57 100644
--- 
a/common/src/main/scala/org/apache/celeborn/common/metrics/source/AbstractSource.scala
+++ 
b/common/src/main/scala/org/apache/celeborn/common/metrics/source/AbstractSource.scala
@@ -138,7 +138,7 @@ abstract class AbstractSource(conf: CelebornConf, role: 
String)
     }
   }
 
-  protected def counters(): List[NamedCounter] = {
+  def counters(): List[NamedCounter] = {
     namedCounters.values().asScala.toList
   }
 
@@ -146,11 +146,11 @@ abstract class AbstractSource(conf: CelebornConf, role: 
String)
     namedGauges.asScala.toList
   }
 
-  protected def histograms(): List[NamedHistogram] = {
+  def histograms(): List[NamedHistogram] = {
     List.empty[NamedHistogram]
   }
 
-  protected def timers(): List[NamedTimer] = {
+  def timers(): List[NamedTimer] = {
     namedTimers.values().asScala.toList.map(_._1)
   }
 
@@ -401,7 +401,7 @@ abstract class AbstractSource(conf: CelebornConf, role: 
String)
     s"metrics_${key.replaceAll("[^a-zA-Z0-9]", "_")}_"
   }
 
-  protected def reportNanosAsMills(value: Double): Double = {
+  def reportNanosAsMills(value: Double): Double = {
     BigDecimal(value / 1000000).setScale(2, 
BigDecimal.RoundingMode.HALF_UP).toDouble
   }
 
diff --git a/conf/metrics.properties.template b/conf/metrics.properties.template
index 7c5630afd..56ae280cf 100644
--- a/conf/metrics.properties.template
+++ b/conf/metrics.properties.template
@@ -16,3 +16,4 @@
 #
 
 
*.sink.prometheusServlet.class=org.apache.celeborn.common.metrics.sink.PrometheusServlet
+*.sink.jsonServlet.class=org.apache.celeborn.common.metrics.sink.JsonServlet
diff --git a/dev/deps/dependencies-client-flink-1.14 
b/dev/deps/dependencies-client-flink-1.14
index 111bf56cf..64a4ba77d 100644
--- a/dev/deps/dependencies-client-flink-1.14
+++ b/dev/deps/dependencies-client-flink-1.14
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-flink-1.15 
b/dev/deps/dependencies-client-flink-1.15
index 111bf56cf..64a4ba77d 100644
--- a/dev/deps/dependencies-client-flink-1.15
+++ b/dev/deps/dependencies-client-flink-1.15
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-flink-1.17 
b/dev/deps/dependencies-client-flink-1.17
index 111bf56cf..64a4ba77d 100644
--- a/dev/deps/dependencies-client-flink-1.17
+++ b/dev/deps/dependencies-client-flink-1.17
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-flink-1.18 
b/dev/deps/dependencies-client-flink-1.18
index 111bf56cf..64a4ba77d 100644
--- a/dev/deps/dependencies-client-flink-1.18
+++ b/dev/deps/dependencies-client-flink-1.18
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-mr b/dev/deps/dependencies-client-mr
index d1e8720c4..362227a6f 100644
--- a/dev/deps/dependencies-client-mr
+++ b/dev/deps/dependencies-client-mr
@@ -68,12 +68,15 @@ 
hadoop-yarn-server-nodemanager/3.3.6//hadoop-yarn-server-nodemanager-3.3.6.jar
 hadoop-yarn-server-web-proxy/3.3.6//hadoop-yarn-server-web-proxy-3.3.6.jar
 httpclient/4.5.13//httpclient-4.5.13.jar
 httpcore/4.4.13//httpcore-4.4.13.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
 jackson-core-asl/1.9.13//jackson-core-asl-1.9.13.jar
-jackson-core/2.12.7//jackson-core-2.12.7.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
 jackson-jaxrs-base/2.12.7//jackson-jaxrs-base-2.12.7.jar
 jackson-jaxrs-json-provider/2.12.7//jackson-jaxrs-json-provider-2.12.7.jar
 jackson-mapper-asl/1.9.13//jackson-mapper-asl-1.9.13.jar
 
jackson-module-jaxb-annotations/2.12.7//jackson-module-jaxb-annotations-2.12.7.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jakarta.activation-api/1.2.1//jakarta.activation-api-1.2.1.jar
 jakarta.xml.bind-api/2.3.2//jakarta.xml.bind-api-2.3.2.jar
 
javax-websocket-client-impl/9.4.51.v20230217//javax-websocket-client-impl-9.4.51.v20230217.jar
@@ -174,7 +177,7 @@ netty/3.10.6.Final//netty-3.10.6.Final.jar
 nimbus-jose-jwt/9.8.1//nimbus-jose-jwt-9.8.1.jar
 okhttp/4.9.3//okhttp-4.9.3.jar
 okio/2.8.0//okio-2.8.0.jar
-paranamer/2.3//paranamer-2.3.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-spark-2.4 
b/dev/deps/dependencies-client-spark-2.4
index a002ea22b..2416aff41 100644
--- a/dev/deps/dependencies-client-spark-2.4
+++ b/dev/deps/dependencies-client-spark-2.4
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.11/2.15.3//jackson-module-scala_2.11-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-spark-3.0 
b/dev/deps/dependencies-client-spark-3.0
index 00adfa190..0dcc4c9d7 100644
--- a/dev/deps/dependencies-client-spark-3.0
+++ b/dev/deps/dependencies-client-spark-3.0
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-spark-3.1 
b/dev/deps/dependencies-client-spark-3.1
index b9c2c81fd..4342df2da 100644
--- a/dev/deps/dependencies-client-spark-3.1
+++ b/dev/deps/dependencies-client-spark-3.1
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-spark-3.2 
b/dev/deps/dependencies-client-spark-3.2
index 2c3cf1672..5b18f5160 100644
--- a/dev/deps/dependencies-client-spark-3.2
+++ b/dev/deps/dependencies-client-spark-3.2
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-spark-3.3 
b/dev/deps/dependencies-client-spark-3.3
index 111bf56cf..64a4ba77d 100644
--- a/dev/deps/dependencies-client-spark-3.3
+++ b/dev/deps/dependencies-client-spark-3.3
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-spark-3.4 
b/dev/deps/dependencies-client-spark-3.4
index dd2ad1a3e..ba9f16bfd 100644
--- a/dev/deps/dependencies-client-spark-3.4
+++ b/dev/deps/dependencies-client-spark-3.4
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-client-spark-3.5 
b/dev/deps/dependencies-client-spark-3.5
index d3a66fca3..8f9977c89 100644
--- a/dev/deps/dependencies-client-spark-3.5
+++ b/dev/deps/dependencies-client-spark-3.5
@@ -24,6 +24,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
 jsr305/1.3.9//jsr305-1.3.9.jar
 jul-to-slf4j/1.7.36//jul-to-slf4j-1.7.36.jar
@@ -66,6 +70,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/dev/deps/dependencies-server b/dev/deps/dependencies-server
index efc6c68c4..fe3785007 100644
--- a/dev/deps/dependencies-server
+++ b/dev/deps/dependencies-server
@@ -25,6 +25,10 @@ failureaccess/1.0.1//failureaccess-1.0.1.jar
 guava/32.1.3-jre//guava-32.1.3-jre.jar
 hadoop-client-api/3.3.6//hadoop-client-api-3.3.6.jar
 hadoop-client-runtime/3.3.6//hadoop-client-runtime-3.3.6.jar
+jackson-annotations/2.15.3//jackson-annotations-2.15.3.jar
+jackson-core/2.15.3//jackson-core-2.15.3.jar
+jackson-databind/2.15.3//jackson-databind-2.15.3.jar
+jackson-module-scala_2.12/2.15.3//jackson-module-scala_2.12-2.15.3.jar
 javassist/3.28.0-GA//javassist-3.28.0-GA.jar
 javax.servlet-api/3.1.0//javax.servlet-api-3.1.0.jar
 jcl-over-slf4j/1.7.36//jcl-over-slf4j-1.7.36.jar
@@ -73,6 +77,7 @@ 
netty-transport-rxtx/4.1.93.Final//netty-transport-rxtx-4.1.93.Final.jar
 netty-transport-sctp/4.1.93.Final//netty-transport-sctp-4.1.93.Final.jar
 netty-transport-udt/4.1.93.Final//netty-transport-udt-4.1.93.Final.jar
 netty-transport/4.1.93.Final//netty-transport-4.1.93.Final.jar
+paranamer/2.8//paranamer-2.8.jar
 protobuf-java/3.19.2//protobuf-java-3.19.2.jar
 ratis-client/2.5.1//ratis-client-2.5.1.jar
 ratis-common/2.5.1//ratis-common-2.5.1.jar
diff --git a/docs/configuration/metrics.md b/docs/configuration/metrics.md
index fb2d38c9e..fd1beadfa 100644
--- a/docs/configuration/metrics.md
+++ b/docs/configuration/metrics.md
@@ -27,6 +27,8 @@ license: |
 | celeborn.metrics.conf | &lt;undefined&gt; | Custom metrics configuration 
file path. Default use `metrics.properties` in classpath. | 0.3.0 | 
 | celeborn.metrics.enabled | true | When true, enable metrics system. | 0.2.0 
| 
 | celeborn.metrics.extraLabels |  | If default metric labels are not enough, 
extra metric labels can be customized. Labels' pattern is: 
`<label1_key>=<label1_value>[,<label2_key>=<label2_value>]*`; e.g. 
`env=prod,version=1` | 0.3.0 | 
+| celeborn.metrics.json.path | /metrics/json | URI context path of json 
metrics HTTP server. | 0.4.0 | 
+| celeborn.metrics.json.pretty.enabled | true | When true, view metrics in 
json pretty format | 0.4.0 | 
 | celeborn.metrics.prometheus.path | /metrics/prometheus | URI context path of 
prometheus metrics HTTP server. | 0.4.0 | 
 | celeborn.metrics.sample.rate | 1.0 | It controls if Celeborn collect timer 
metrics for some operations. Its value should be in [0.0, 1.0]. | 0.2.0 | 
 | celeborn.metrics.timer.slidingWindow.size | 4096 | The sliding window size 
of timer metric. | 0.2.0 | 
diff --git 
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala 
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
index e17bb46bf..85e7a2da9 100644
--- 
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
+++ 
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
@@ -31,7 +31,6 @@ import org.apache.ratis.proto.RaftProtos
 import org.apache.ratis.proto.RaftProtos.RaftPeerRole
 
 import org.apache.celeborn.common.CelebornConf
-import org.apache.celeborn.common.CelebornConf.METRICS_PROMETHEUS_PATH
 import org.apache.celeborn.common.client.MasterClient
 import org.apache.celeborn.common.identity.UserIdentifier
 import org.apache.celeborn.common.internal.Logging
@@ -58,7 +57,7 @@ private[celeborn] class Master(
   override def serviceName: String = Service.MASTER
 
   override val metricsSystem: MetricsSystem =
-    MetricsSystem.createMetricsSystem(serviceName, conf, 
conf.get(METRICS_PROMETHEUS_PATH))
+    MetricsSystem.createMetricsSystem(serviceName, conf)
 
   override val rpcEnv: RpcEnv = RpcEnv.create(
     RpcNameConstants.MASTER_SYS,
diff --git a/pom.xml b/pom.xml
index cbf1f3c6b..b64d25fd9 100644
--- a/pom.xml
+++ b/pom.xml
@@ -96,6 +96,7 @@
     <zstd-jni.version>1.5.2-1</zstd-jni.version>
     <kubernetes-client.version>6.7.0</kubernetes-client.version>
     <rocksdbjni.version>8.5.3</rocksdbjni.version>
+    <jackson.version>2.15.3</jackson.version>
 
     <shading.prefix>org.apache.celeborn.shaded</shading.prefix>
 
@@ -464,6 +465,26 @@
         <artifactId>rocksdbjni</artifactId>
         <version>${rocksdbjni.version}</version>
       </dependency>
+      <dependency>
+        <groupId>com.fasterxml.jackson.module</groupId>
+        <artifactId>jackson-module-scala_${scala.binary.version}</artifactId>
+        <version>${jackson.version}</version>
+      </dependency>
+      <dependency>
+        <groupId>com.fasterxml.jackson.core</groupId>
+        <artifactId>jackson-databind</artifactId>
+        <version>${jackson.version}</version>
+      </dependency>
+      <dependency>
+        <groupId>com.fasterxml.jackson.core</groupId>
+        <artifactId>jackson-annotations</artifactId>
+        <version>${jackson.version}</version>
+      </dependency>
+      <dependency>
+        <groupId>com.fasterxml.jackson.core</groupId>
+        <artifactId>jackson-core</artifactId>
+        <version>${jackson.version}</version>
+      </dependency>
     </dependencies>
   </dependencyManagement>
 
diff --git a/project/CelebornBuild.scala b/project/CelebornBuild.scala
index 96f590768..feef05ad9 100644
--- a/project/CelebornBuild.scala
+++ b/project/CelebornBuild.scala
@@ -57,6 +57,7 @@ object Dependencies {
   val ratisVersion = "2.5.1"
   val roaringBitmapVersion = "0.9.32"
   val rocksdbJniVersion = "8.5.3"
+  val jacksonVersion = "2.15.3"
   val scalatestMockitoVersion = "1.17.14"
   val scalatestVersion = "3.2.16"
   val slf4jVersion = "1.7.36"
@@ -110,6 +111,10 @@ object Dependencies {
     ExclusionRule("org.slf4j", "slf4j-simple"))
   val roaringBitmap = "org.roaringbitmap" % "RoaringBitmap" % 
roaringBitmapVersion
   val rocksdbJni = "org.rocksdb" % "rocksdbjni" % rocksdbJniVersion
+  val jacksonDatabind = "com.fasterxml.jackson.core" % "jackson-databind" % 
jacksonVersion
+  val jacksonCore = "com.fasterxml.jackson.core" % "jackson-core" % 
jacksonVersion
+  val jacksonAnnotations = "com.fasterxml.jackson.core" % 
"jackson-annotations" % jacksonVersion
+  val jacksonModule = "com.fasterxml.jackson.module" %% "jackson-module-scala" 
% jacksonVersion
   val scalaReflect = "org.scala-lang" % "scala-reflect" % projectScalaVersion
   val slf4jApi = "org.slf4j" % "slf4j-api" % slf4jVersion
   val slf4jJulToSlf4j = "org.slf4j" % "jul-to-slf4j" % slf4jVersion
@@ -262,7 +267,6 @@ object Utils {
       }
       profiles
   }
-
   val SPARK_VERSION = profiles.filter(_.startsWith("spark")).headOption
 
   lazy val sparkClientProjects = SPARK_VERSION match {
@@ -336,6 +340,10 @@ object CelebornCommon {
         Dependencies.slf4jJulToSlf4j,
         Dependencies.slf4jApi,
         Dependencies.snakeyaml,
+        Dependencies.jacksonModule,
+        Dependencies.jacksonCore,
+        Dependencies.jacksonDatabind,
+        Dependencies.jacksonAnnotations,
         Dependencies.log4jSlf4jImpl % "test",
         Dependencies.log4j12Api % "test"
       ) ++ commonUnitTestDependencies,
diff --git 
a/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala 
b/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
index 5f9ff1f5a..4918bcc23 100644
--- a/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
+++ b/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
@@ -79,7 +79,7 @@ abstract class HttpService extends Service with Logging {
   def startHttpServer(): Unit = {
     val handlers =
       if (metricsSystem.running) {
-        new HttpRequestHandler(this, metricsSystem.getPrometheusHandler)
+        new HttpRequestHandler(this, metricsSystem.getServletHandlers)
       } else {
         new HttpRequestHandler(this, null)
       }
diff --git 
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
 
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
index da59f9f5c..531f19454 100644
--- 
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
+++ 
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
@@ -24,7 +24,7 @@ import io.netty.handler.codec.http._
 import io.netty.util.CharsetUtil
 
 import org.apache.celeborn.common.internal.Logging
-import org.apache.celeborn.common.metrics.sink.PrometheusHttpRequestHandler
+import org.apache.celeborn.common.metrics.sink.{JsonHttpRequestHandler, 
ServletHttpRequestHandler}
 import org.apache.celeborn.server.common.HttpService
 
 /**
@@ -36,7 +36,7 @@ import org.apache.celeborn.server.common.HttpService
 @Sharable
 class HttpRequestHandler(
     service: HttpService,
-    prometheusHttpRequestHandler: PrometheusHttpRequestHandler)
+    servletHttpRequestHandlers: Array[ServletHttpRequestHandler])
   extends SimpleChannelInboundHandler[FullHttpRequest] with Logging {
 
   override def channelReadComplete(ctx: ChannelHandlerContext): Unit = {
@@ -47,21 +47,31 @@ class HttpRequestHandler(
     val uri = req.uri()
     val (path, parameters) = HttpUtils.parseUri(uri)
     val msg = HttpUtils.handleRequest(service, path, parameters)
-    val response = msg match {
+    val textType = "text/plain; charset=UTF-8"
+    val jsonType = "application/json"
+    val (response, contentType) = msg match {
       case Invalid.invalid =>
-        if (prometheusHttpRequestHandler != null) {
-          prometheusHttpRequestHandler.handleRequest(uri)
+        if (servletHttpRequestHandlers != null) {
+          servletHttpRequestHandlers.find(servlet =>
+            uri == servlet.getServletPath()).map {
+            case jsonHandler: JsonHttpRequestHandler =>
+              (jsonHandler.handleRequest(uri), jsonType)
+            case handler: ServletHttpRequestHandler =>
+              (handler.handleRequest(uri), textType)
+          }.getOrElse((s"Unknown path $uri!", textType))
         } else {
-          s"${Invalid.description(service.serviceName)} 
${HttpUtils.help(service.serviceName)}"
+          (
+            s"${Invalid.description(service.serviceName)} 
${HttpUtils.help(service.serviceName)}",
+            textType)
         }
-      case _ => msg
+      case _ => (msg, textType)
     }
 
     val res = new DefaultFullHttpResponse(
       HttpVersion.HTTP_1_1,
       HttpResponseStatus.OK,
       Unpooled.copiedBuffer(response, CharsetUtil.UTF_8))
-    res.headers().set(HttpHeaderNames.CONTENT_TYPE, "text/plain; 
charset=UTF-8")
+    res.headers().set(HttpHeaderNames.CONTENT_TYPE, contentType)
     ctx.writeAndFlush(res).addListener(ChannelFutureListener.CLOSE)
   }
 }
diff --git 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
index ecf867b3e..aa6e7aa24 100644
--- 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
+++ 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
@@ -64,7 +64,7 @@ private[celeborn] class Worker(
   override def serviceName: String = Service.WORKER
 
   override val metricsSystem: MetricsSystem =
-    MetricsSystem.createMetricsSystem(serviceName, conf, 
conf.get(METRICS_PROMETHEUS_PATH))
+    MetricsSystem.createMetricsSystem(serviceName, conf)
 
   val rpcEnv: RpcEnv = RpcEnv.create(
     RpcNameConstants.WORKER_SYS,

Reply via email to