This is an automated email from the ASF dual-hosted git repository.
chengpan 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 9a700cec9 [CELEBORN-1238] deviceCheckThreadPool is only initialized
when diskCheck is enabled
9a700cec9 is described below
commit 9a700cec90baf570fc227a3f40ce01115a07f5ea
Author: xianminglei <[email protected]>
AuthorDate: Fri Jan 19 17:36:28 2024 +0800
[CELEBORN-1238] deviceCheckThreadPool is only initialized when diskCheck is
enabled
### What changes were proposed in this pull request?
deviceCheckThreadPool is only initialized when diskCheck is enabled
### Why are the changes needed?
deviceCheckThreadPool is only initialized when diskCheck is enabled
### Does this PR introduce _any_ user-facing change?
No
### How was this patch tested?
Existing UTs.
Closes #2242 from leixm/issue_1238.
Authored-by: xianminglei <[email protected]>
Signed-off-by: Cheng Pan <[email protected]>
(cherry picked from commit ef47645b7a37b4cc3e6669978197ec88659d8baf)
Signed-off-by: Cheng Pan <[email protected]>
---
.../celeborn/service/deploy/worker/storage/DeviceMonitor.scala | 5 +++--
1 file changed, 3 insertions(+), 2 deletions(-)
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/DeviceMonitor.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/DeviceMonitor.scala
index 8acb8b2b7..ad961c0b0 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/DeviceMonitor.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/DeviceMonitor.scala
@@ -21,7 +21,7 @@ import java.io.{BufferedReader, File, FileInputStream,
InputStreamReader, IOExce
import java.nio.charset.Charset
import java.nio.file.{Files, Paths}
import java.util
-import java.util.concurrent.TimeUnit
+import java.util.concurrent.{ThreadPoolExecutor, TimeUnit}
import scala.collection.JavaConverters._
@@ -198,7 +198,7 @@ class LocalDeviceMonitor(
}
object DeviceMonitor extends Logging {
- val deviceCheckThreadPool =
ThreadUtils.newDaemonCachedThreadPool("device-check-thread", 5)
+ var deviceCheckThreadPool: ThreadPoolExecutor = _
def createDeviceMonitor(
conf: CelebornConf,
@@ -208,6 +208,7 @@ object DeviceMonitor extends Logging {
workerSource: AbstractSource): DeviceMonitor = {
try {
if (conf.workerDiskMonitorEnabled) {
+ deviceCheckThreadPool =
ThreadUtils.newDaemonCachedThreadPool("device-check-thread", 5)
val monitor =
new LocalDeviceMonitor(conf, deviceObserver, deviceInfos, diskInfos,
workerSource)
monitor.init()