This is an automated email from the ASF dual-hosted git repository.
angerszhuuuu pushed a commit to branch branch-0.4
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/branch-0.4 by this push:
new 8d5f30fca [CELEBORN-1292] Remove app level metrics from worker and
master
8d5f30fca is described below
commit 8d5f30fca87ba31815429e6a00d197666de21159
Author: Angerszhuuuu <[email protected]>
AuthorDate: Wed Feb 28 17:07:51 2024 +0800
[CELEBORN-1292] Remove app level metrics from worker and master
### What changes were proposed in this pull request?
Master side calculate sub resource consumption occupy cpu cause rpc time
out and miss prometheus metrics
<img width="1781" alt="截屏2024-02-28 12 04 19"
src="https://github.com/apache/incubator-celeborn/assets/46485123/ba49a4ac-ec49-4234-8758-c0db9242abf6">
Worker side generate too much metrics data
### Why are the changes needed?
Fix performance issue
### Does this PR introduce _any_ user-facing change?
Remove app level metrics
### How was this patch tested?
Closes #2342 from AngersZhuuuu/CELEBORN-1292.
Authored-by: Angerszhuuuu <[email protected]>
Signed-off-by: Angerszhuuuu <[email protected]>
(cherry picked from commit 3f884660b42f65b6fd54895bcb7bba89a54deac9)
Signed-off-by: Angerszhuuuu <[email protected]>
---
.../celeborn/service/deploy/master/Master.scala | 37 ++--------------------
.../celeborn/service/deploy/worker/Worker.scala | 4 ---
2 files changed, 3 insertions(+), 38 deletions(-)
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 4f9826a2a..ddcb245dd 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
@@ -783,8 +783,6 @@ private[celeborn] class Master(
nonEagerHandler.submit(new Runnable {
override def run(): Unit = {
statusSystem.handleAppLost(appId, requestId)
- // Resource consumption source should remove lose application gauges.
- removeAppResourceConsumption(appId)
logInfo(s"Removed application $appId")
if (hasHDFSStorage) {
checkAndCleanExpiredAppDirsOnHDFS(appId)
@@ -794,30 +792,6 @@ private[celeborn] class Master(
})
}
- private def removeAppResourceConsumption(appId: String): Unit = {
- removeResourceConsumptionGauge(
- ResourceConsumptionSource.DISK_FILE_COUNT,
- appId)
- removeResourceConsumptionGauge(
- ResourceConsumptionSource.DISK_BYTES_WRITTEN,
- appId)
- removeResourceConsumptionGauge(
- ResourceConsumptionSource.HDFS_FILE_COUNT,
- appId)
- removeResourceConsumptionGauge(
- ResourceConsumptionSource.HDFS_BYTES_WRITTEN,
- appId)
- }
-
- private def removeResourceConsumptionGauge(
- resourceConsumptionName: String,
- appId: String): Unit = {
- resourceConsumptionSource.removeGauge(
- resourceConsumptionName,
- resourceConsumptionSource.applicationLabel,
- appId)
- }
-
private def checkAndCleanExpiredAppDirsOnHDFS(expiredDir: String = ""): Unit
= {
if (hadoopFs == null) {
try {
@@ -883,10 +857,6 @@ private[celeborn] class Master(
private def handleResourceConsumption(userIdentifier: UserIdentifier):
ResourceConsumption = {
val userResourceConsumption =
computeUserResourceConsumption(userIdentifier)
gaugeResourceConsumption(userIdentifier)
- val subResourceConsumptions =
userResourceConsumption.subResourceConsumptions
- if (CollectionUtils.isNotEmpty(subResourceConsumptions)) {
- subResourceConsumptions.asScala.keys.foreach {
gaugeResourceConsumption(userIdentifier, _) }
- }
userResourceConsumption
}
@@ -937,13 +907,12 @@ private[celeborn] class Master(
}
}
+ // TODO: Support calculate topN app resource consumption.
private def computeUserResourceConsumption(
userIdentifier: UserIdentifier): ResourceConsumption = {
- val (resourceConsumption, subResourceConsumptions) =
statusSystem.workers.asScala.flatMap {
+ val resourceConsumption = statusSystem.workers.asScala.flatMap {
workerInfo =>
workerInfo.userResourceConsumption.asScala.get(userIdentifier)
- }.foldRight((ResourceConsumption(0, 0, 0, 0), Map.empty[String,
ResourceConsumption]))(
- _ addWithSubResourceConsumptions _)
- resourceConsumption.subResourceConsumptions =
subResourceConsumptions.asJava
+ }.foldRight(ResourceConsumption(0, 0, 0, 0))(_ add _)
resourceConsumption
}
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 941b662d3..873c91f37 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
@@ -531,10 +531,6 @@ private[celeborn] class Worker(
workerInfo.updateThenGetUserResourceConsumption(resourceConsumptionSnapshot.asJava)
resourceConsumptionSnapshot.foreach { case (userIdentifier,
userResourceConsumption) =>
gaugeResourceConsumption(userIdentifier)
- val subResourceConsumptions =
userResourceConsumption.subResourceConsumptions
- if (CollectionUtils.isNotEmpty(subResourceConsumptions)) {
- subResourceConsumptions.asScala.keys.foreach {
gaugeResourceConsumption(userIdentifier, _) }
- }
}
userResourceConsumptions
}