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
   }

Reply via email to