This is an automated email from the ASF dual-hosted git repository.
zhouky pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new 4f039d5f7 [CELEBORN-1453] Fix the thread safety bug in getMetrics
4f039d5f7 is described below
commit 4f039d5f71be743bd86439a61af0423f006bc190
Author: xinyuwang1 <[email protected]>
AuthorDate: Sat Jun 8 11:30:11 2024 +0800
[CELEBORN-1453] Fix the thread safety bug in getMetrics
### 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]>
---
.../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 98e264a93..72611598f 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
@@ -392,26 +392,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 = {