This is an automated email from the ASF dual-hosted git repository.
ethanfeng pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new 1495fca69 [CELEBORN-1219] takeBuffer() avoid checking
source.metricsCollectCriticalEnabled twice
1495fca69 is described below
commit 1495fca69bf183a401911daaed92dce78f7d62d5
Author: Angerszhuuuu <[email protected]>
AuthorDate: Mon Jan 15 16:50:40 2024 +0800
[CELEBORN-1219] takeBuffer() avoid checking
source.metricsCollectCriticalEnabled twice
### What changes were proposed in this pull request?
takeBuffer() avoid checking source.metricsCollectCriticalEnabled twice
### Why are the changes needed?
### Does this PR introduce _any_ user-facing change?
### How was this patch tested?
Closes #2223 from AngersZhuuuu/CELEBORN-1219.
Authored-by: Angerszhuuuu <[email protected]>
Signed-off-by: mingji <[email protected]>
---
.../deploy/worker/storage/PartitionDataWriter.java | 23 +++++++++-------------
1 file changed, 9 insertions(+), 14 deletions(-)
diff --git
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionDataWriter.java
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionDataWriter.java
index 3185abb87..2058dd649 100644
---
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionDataWriter.java
+++
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionDataWriter.java
@@ -368,23 +368,18 @@ public abstract class PartitionDataWriter implements
DeviceObserver {
}
protected void takeBuffer() {
- // metrics start
- String metricsName = null;
- String fileAbsPath = null;
if (source.metricsCollectCriticalEnabled()) {
- metricsName = WorkerSource.TAKE_BUFFER_TIME();
- fileAbsPath = diskFileInfo.getFilePath();
+ String metricsName = WorkerSource.TAKE_BUFFER_TIME();
+ String fileAbsPath = diskFileInfo.getFilePath();
source.startTimer(metricsName, fileAbsPath);
- }
-
- // real action
- synchronized (flushLock) {
- flushBuffer = flusher.takeBuffer();
- }
-
- // metrics end
- if (source.metricsCollectCriticalEnabled()) {
+ synchronized (flushLock) {
+ flushBuffer = flusher.takeBuffer();
+ }
source.stopTimer(metricsName, fileAbsPath);
+ } else {
+ synchronized (flushLock) {
+ flushBuffer = flusher.takeBuffer();
+ }
}
}