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

chengpan 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 ef47645b7 [CELEBORN-1238] deviceCheckThreadPool is only initialized 
when diskCheck is enabled
ef47645b7 is described below

commit ef47645b7a37b4cc3e6669978197ec88659d8baf
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]>
---
 .../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 289a41488..aa0239af4 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()

Reply via email to