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

ethanfeng pushed a commit to branch branch-0.4
in repository https://gitbox.apache.org/repos/asf/celeborn.git


The following commit(s) were added to refs/heads/branch-0.4 by this push:
     new 2d3ae1bec [CELEBORN-1453] Fix the thread safety bug in getMetrics 
(#2566)
2d3ae1bec is described below

commit 2d3ae1bec5473d097d1cfc473c8dd2954c35a0d3
Author: Fu Chen <[email protected]>
AuthorDate: Fri Jun 14 13:39:33 2024 +0800

    [CELEBORN-1453] Fix the thread safety bug in getMetrics (#2566)
    
    ### What changes were proposed in this pull request?
    Fix the thread safety bug in getMetrics of AbstractSource by changing the 
lock scope
    
    ### Why are the changes needed?
    When two threads access the getMetrics method in AbstractSource at the same 
time, one of the threads may get fewer metrics than the actual value, because 
the actual execution order may be like this: Thread A gets the lock, adds the 
metrics of the worker source to the innerMetrics queue and releases the lock, 
Thread B gets the lock, adds the metrics of the worker source to the 
innerMetrics queue and releases the lock, Thread A gets the lock, adds the 
metrics of other sources to the inn [...]
    
    ### Does this PR introduce _any_ user-facing change?
    No
    
    ### How was this patch tested?
    manual test
    
    Closes #2548 from littlexyw/get_metrics_fix.
    
    Authored-by: xinyuwang1 <[email protected]>
    
    Signed-off-by: zky.zhoukeyong <[email protected]>
    Co-authored-by: xinyuwang1 <[email protected]>
---
 .../common/metrics/source/AbstractSource.scala     | 28 +++++++++++-----------
 1 file changed, 14 insertions(+), 14 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..e77dd0685 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
@@ -384,26 +384,26 @@ abstract class AbstractSource(conf: CelebornConf, role: 
String)
   }
 
   override def getMetrics(): String = {
-    counters().foreach(c => recordCounter(c))
-    gauges().foreach(g => recordGauge(g))
-    histograms().foreach(h => {
-      recordHistogram(h)
-      h.asInstanceOf[CelebornHistogram].reservoir
-        .asInstanceOf[ResettableSlidingWindowReservoir].reset()
-    })
-    timers().foreach(t => {
-      recordTimer(t)
-      t.timer.asInstanceOf[CelebornTimer].reservoir
-        .asInstanceOf[ResettableSlidingWindowReservoir].reset()
-    })
-    val sb = new mutable.StringBuilder
     innerMetrics.synchronized {
+      counters().foreach(c => recordCounter(c))
+      gauges().foreach(g => recordGauge(g))
+      histograms().foreach(h => {
+        recordHistogram(h)
+        h.asInstanceOf[CelebornHistogram].reservoir
+          .asInstanceOf[ResettableSlidingWindowReservoir].reset()
+      })
+      timers().foreach(t => {
+        recordTimer(t)
+        t.timer.asInstanceOf[CelebornTimer].reservoir
+          .asInstanceOf[ResettableSlidingWindowReservoir].reset()
+      })
+      val sb = new mutable.StringBuilder
       while (!innerMetrics.isEmpty) {
         sb.append(innerMetrics.poll())
       }
       innerMetrics.clear()
+      sb.toString()
     }
-    sb.toString()
   }
 
   override def destroy(): Unit = {

Reply via email to