This is an automated email from the ASF dual-hosted git repository.
ethanfeng pushed a commit to branch branch-0.5
in repository https://gitbox.apache.org/repos/asf/celeborn.git
The following commit(s) were added to refs/heads/branch-0.5 by this push:
new 74d982b20 [CELEBORN-1182][FOLLOWUP] WorkerSource should use Counter to
support application dimension ActiveConnectionCount metric
74d982b20 is described below
commit 74d982b20f959339b1fb34e9399b97147031462a
Author: SteNicholas <[email protected]>
AuthorDate: Thu Jun 6 17:25:39 2024 +0800
[CELEBORN-1182][FOLLOWUP] WorkerSource should use Counter to support
application dimension ActiveConnectionCount metric
### What changes were proposed in this pull request?
`WorkerSource` should use Counter to support application dimension
`ActiveConnectionCount` metric.
Follow up #2167.
### Why are the changes needed?
`WorkerSource` uses `Gauge` for application dimension ActiveConnectionCount
metric via `appActiveConnections` at present, which has performance problem in
metrics REST API as follows:
```
"worker-JettyThreadPool-11242" #11242 daemon prio=5 os_prio=0
tid=0x00007f410800c000 nid=0x2d80 runnable [0x00007f3426de2000]
java.lang.Thread.State: RUNNABLE
at scala.collection.Iterator.foreach(Iterator.scala:941)
at scala.collection.Iterator.foreach$(Iterator.scala:941)
at scala.collection.AbstractIterator.foreach(Iterator.scala:1429)
at scala.collection.IterableLike.foreach(IterableLike.scala:74)
at scala.collection.IterableLike.foreach$(IterableLike.scala:73)
at scala.collection.AbstractIterable.foreach(Iterable.scala:56)
at scala.collection.TraversableOnce.count(TraversableOnce.scala:118)
at
scala.collection.TraversableOnce.count$(TraversableOnce.scala:116)
at scala.collection.AbstractTraversable.count(Traversable.scala:108)
at
org.apache.celeborn.service.deploy.worker.WorkerSource.$anonfun$recordAppActiveConnection$1(WorkerSource.scala:104)
at
org.apache.celeborn.service.deploy.worker.WorkerSource$$Lambda$787/1074905995.apply$mcI$sp(Unknown
Source)
at
scala.runtime.java8.JFunction0$mcI$sp.apply(JFunction0$mcI$sp.java:23)
at
org.apache.celeborn.common.metrics.source.GaugeSupplier$$anon$3.getValue(AbstractSource.scala:466)
at
org.apache.celeborn.common.metrics.source.AbstractSource.recordGauge(AbstractSource.scala:342)
at
org.apache.celeborn.common.metrics.source.AbstractSource.$anonfun$getMetrics$2(AbstractSource.scala:401)
at
org.apache.celeborn.common.metrics.source.AbstractSource.$anonfun$getMetrics$2$adapted(AbstractSource.scala:401)
at
org.apache.celeborn.common.metrics.source.AbstractSource$$Lambda$956/1021547679.apply(Unknown
Source)
at scala.collection.immutable.List.foreach(List.scala:392)
at
org.apache.celeborn.common.metrics.source.AbstractSource.getMetrics(AbstractSource.scala:401)
at
org.apache.celeborn.common.metrics.sink.AbstractServlet.$anonfun$getMetricsSnapshot$1(AbstractServlet.scala:34)
at
org.apache.celeborn.common.metrics.sink.AbstractServlet$$Lambda$954/1559941228.apply(Unknown
Source)
at
scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:238)
at
scala.collection.TraversableLike$$Lambda$33/829149076.apply(Unknown Source)
at scala.collection.Iterator.foreach(Iterator.scala:941)
at scala.collection.Iterator.foreach$(Iterator.scala:941)
at scala.collection.AbstractIterator.foreach(Iterator.scala:1429)
at scala.collection.IterableLike.foreach(IterableLike.scala:74)
at scala.collection.IterableLike.foreach$(IterableLike.scala:73)
at scala.collection.AbstractIterable.foreach(Iterable.scala:56)
at scala.collection.TraversableLike.map(TraversableLike.scala:238)
at scala.collection.TraversableLike.map$(TraversableLike.scala:231)
at scala.collection.AbstractTraversable.map(Traversable.scala:108)
at
org.apache.celeborn.common.metrics.sink.AbstractServlet.getMetricsSnapshot(AbstractServlet.scala:34)
at
org.apache.celeborn.common.metrics.sink.PrometheusServlet.$anonfun$createServletHandler$1(PrometheusServlet.scala:38)
at
org.apache.celeborn.common.metrics.sink.PrometheusServlet$$Lambda$721/2120532393.apply(Unknown
Source)
at
org.apache.celeborn.server.common.http.HttpUtils$$anon$1.doGet(HttpUtils.scala:51)
at javax.servlet.http.HttpServlet.service(HttpServlet.java:497)
at javax.servlet.http.HttpServlet.service(HttpServlet.java:584)
at
org.eclipse.jetty.servlet.ServletHolder.handle(ServletHolder.java:799)
at
org.eclipse.jetty.servlet.ServletHandler.doHandle(ServletHandler.java:554)
at
org.eclipse.jetty.server.handler.ScopedHandler.nextHandle(ScopedHandler.java:233)
at
org.eclipse.jetty.server.handler.ContextHandler.doHandle(ContextHandler.java:1440)
at
org.eclipse.jetty.server.handler.ScopedHandler.nextScope(ScopedHandler.java:188)
at
org.eclipse.jetty.servlet.ServletHandler.doScope(ServletHandler.java:505)
at
org.eclipse.jetty.server.handler.ScopedHandler.nextScope(ScopedHandler.java:186)
at
org.eclipse.jetty.server.handler.ContextHandler.doScope(ContextHandler.java:1355)
at
org.eclipse.jetty.server.handler.ScopedHandler.handle(ScopedHandler.java:141)
at
org.eclipse.jetty.server.handler.ContextHandlerCollection.handle(ContextHandlerCollection.java:234)
at
org.eclipse.jetty.server.handler.HandlerWrapper.handle(HandlerWrapper.java:127)
at org.eclipse.jetty.server.Server.handle(Server.java:516)
at
org.eclipse.jetty.server.HttpChannel.lambda$handle$1(HttpChannel.java:487)
at
org.eclipse.jetty.server.HttpChannel$$Lambda$636/1962809899.dispatch(Unknown
Source)
at
org.eclipse.jetty.server.HttpChannel.dispatch(HttpChannel.java:732)
at org.eclipse.jetty.server.HttpChannel.handle(HttpChannel.java:479)
at
org.eclipse.jetty.server.HttpConnection.onFillable(HttpConnection.java:277)
at
org.eclipse.jetty.io.AbstractConnection$ReadCallback.succeeded(AbstractConnection.java:311)
at org.eclipse.jetty.io.FillInterest.fillable(FillInterest.java:105)
at
org.eclipse.jetty.io.ChannelEndPoint$1.run(ChannelEndPoint.java:104)
at
org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.runTask(EatWhatYouKill.java:338)
at
org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.doProduce(EatWhatYouKill.java:315)
at
org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.tryProduce(EatWhatYouKill.java:173)
at
org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.run(EatWhatYouKill.java:131)
at
org.eclipse.jetty.util.thread.ReservedThreadExecutor$ReservedThread.run(ReservedThreadExecutor.java:409)
at
org.eclipse.jetty.util.thread.QueuedThreadPool.runJob(QueuedThreadPool.java:883)
at
org.eclipse.jetty.util.thread.QueuedThreadPool$Runner.run(QueuedThreadPool.java:1034)
at java.lang.Thread.run(Thread.java:748)
```
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
Cluster test.
```
$ curl http://bigdata-rss-worker:9096/metrics|grep
ActiveConnectionCount|grep application
% Total % Received % Xferd Average Speed Time Time Time
Current
Dload Upload Total Spent Left
Speed
100 54003 0 54003 0 0 2924k 0 --:--:-- --:--:-- --:--:--
3102k
metrics_ActiveConnectionCount_Count{applicationId="application_1688369676084_17462520_1",hostName="bigdata-rss-worker",role="Worker"}
15 1717590356773
metrics_ActiveConnectionCount_Count{applicationId="application_1650016801129_32165809_1",hostName="bigdata-rss-worker",role="Worker"}
7 171759035677
$ curl http://bigdata-rss-worker:9096/metrics|grep
ActiveConnectionCount|grep application
% Total % Received % Xferd Average Speed Time Time Time
Current
Dload Upload Total Spent Left
Speed
100 54025 0 54025 0 0 2891k 0 --:--:-- --:--:-- --:--:--
2931k
metrics_ActiveConnectionCount_Count{applicationId="application_1688369676084_17462520_1",hostName="bigdata-rss-worker",role="Worker"}
25 1717590431544
metrics_ActiveConnectionCount_Count{applicationId="application_1650016801129_32165809_1",hostName="bigdata-rss-worker",role="Worker"}
14 1717590431544
$ curl http://bigdata-rss-worker:9096/metrics|grep
ActiveConnectionCount|grep application
% Total % Received % Xferd Average Speed Time Time Time
Current
Dload Upload Total Spent Left
Speed
100 54014 0 54014 0 0 2727k 0 --:--:-- --:--:-- --:--:--
2776k
metrics_ActiveConnectionCount_Count{applicationId="application_1688369676084_17462520_1",hostName="bigdata-rss-worker",role="Worker"}
19 1717590480837
metrics_ActiveConnectionCount_Count{applicationId="application_1650016801129_32165809_1",hostName="bigdata-rss-worker",role="Worker"}
9 1717590480837
```
Closes #2546 from SteNicholas/CELEBORN-1182.
Authored-by: SteNicholas <[email protected]>
Signed-off-by: mingji <[email protected]>
(cherry picked from commit a5d3f0f30eb97767d98ae617f397ba9cda93718e)
Signed-off-by: mingji <[email protected]>
---
.../celeborn/common/metrics/source/AbstractSource.scala | 16 ++++++++++++----
.../celeborn/service/deploy/worker/WorkerSource.scala | 15 +++++++--------
2 files changed, 19 insertions(+), 12 deletions(-)
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 301d991a2..98e264a93 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
@@ -166,6 +166,14 @@ abstract class AbstractSource(conf: CelebornConf, role:
String)
}
}
+ def removeCounter(name: String, labels: Map[String, String]): Unit = {
+ val metricNameWithLabel = metricNameWithCustomizedLabels(name, labels)
+ val namedCounter = namedCounters.get(metricNameWithLabel)
+ if (namedCounter != null) {
+ removeMetric(metricNameWithLabel, namedCounter)
+ }
+ }
+
def removeGauge(name: String, labels: Map[String, String]): Unit = {
val labelString = MetricLabels.labelString(labels ++ staticLabels)
@@ -174,7 +182,7 @@ abstract class AbstractSource(conf: CelebornConf, role:
String)
val namedGauge = iter.next()
if (namedGauge.name.equals(name) &&
namedGauge.labelString.equals(labelString)) {
iter.remove()
- removeGaugeMetric(name, namedGauge)
+ removeMetric(name, namedGauge)
return
}
}
@@ -188,14 +196,14 @@ abstract class AbstractSource(conf: CelebornConf, role:
String)
val namedGauge = iter.next()
if (namedGauge.name.equals(name) &&
labels.toSet.subsetOf(namedGauge.labels.toSet)) {
iter.remove()
- removeGaugeMetric(name, namedGauge)
+ removeMetric(name, namedGauge)
return
}
}
}
- def removeGaugeMetric(name: String, namedGauge: NamedGauge[_]): Unit = {
- metricRegistry.remove(metricNameWithCustomizedLabelString(name,
namedGauge.labelString))
+ def removeMetric(name: String, metric: MetricLabels): Unit = {
+ metricRegistry.remove(metricNameWithCustomizedLabelString(name,
metric.labelString))
}
override def sample[T](metricsName: String, key: String)(f: => T): T = {
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/WorkerSource.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/WorkerSource.scala
index e34488bea..ac5e54cab 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/WorkerSource.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/WorkerSource.scala
@@ -28,7 +28,7 @@ import org.apache.celeborn.common.CelebornConf
import org.apache.celeborn.common.metrics.MetricsSystem
import org.apache.celeborn.common.metrics.source.AbstractSource
import org.apache.celeborn.common.network.client.TransportClient
-import org.apache.celeborn.common.util.{CollectionUtils, JavaUtils, Utils}
+import org.apache.celeborn.common.util.{JavaUtils, Utils}
class WorkerSource(conf: CelebornConf) extends AbstractSource(conf,
MetricsSystem.ROLE_WORKER) {
override val sourceName = "worker"
@@ -91,28 +91,27 @@ class WorkerSource(conf: CelebornConf) extends
AbstractSource(conf, MetricsSyste
}
def connectionInactive(client: TransportClient): Unit = {
- appActiveConnections.remove(client.getChannel.id().asLongText())
+ val applicationIds =
appActiveConnections.remove(client.getChannel.id().asLongText())
incCounter(ACTIVE_CONNECTION_COUNT, -1)
+ applicationIds.asScala.foreach(applicationId =>
+ incCounter(ACTIVE_CONNECTION_COUNT, -1, Map(applicationLabel ->
applicationId)))
}
def recordAppActiveConnection(client: TransportClient, shuffleKey: String):
Unit = {
val applicationIds =
appActiveConnections.get(client.getChannel.id().asLongText())
val applicationId = Utils.splitShuffleKey(shuffleKey)._1
if (applicationIds != null && !applicationIds.contains(applicationId)) {
+ addCounter(ACTIVE_CONNECTION_COUNT, Map(applicationLabel ->
applicationId))
+ incCounter(ACTIVE_CONNECTION_COUNT, 1, Map(applicationLabel ->
applicationId))
applicationIds.add(applicationId)
- addGauge(ACTIVE_CONNECTION_COUNT, Map(applicationLabel ->
applicationId)) { () =>
- appActiveConnections.asScala.count { case (_, applicationIds) =>
- applicationIds.contains(applicationId)
- }
- }
}
}
def removeAppActiveConnection(applicationId: String): Unit = {
+ removeCounter(ACTIVE_CONNECTION_COUNT, Map(applicationLabel ->
applicationId))
appActiveConnections.asScala.foreach { case (_, applicationIds) =>
if (applicationIds.contains(applicationId)) {
applicationIds.remove(applicationId)
- removeGauge(ACTIVE_CONNECTION_COUNT, Map(applicationLabel ->
applicationId))
}
}
}