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

angerszhuuuu 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 a224f713b [CELEBORN-1226] Unify creation of thread using ThreadUtils
a224f713b is described below

commit a224f713b4730a87f053e34195e603e58d709780
Author: Angerszhuuuu <[email protected]>
AuthorDate: Mon Jan 22 12:09:22 2024 +0800

    [CELEBORN-1226] Unify creation of thread using ThreadUtils
    
    ### What changes were proposed in this pull request?
    Make all single thread use standard ThreadUtils to simplify the code
    
    ### Why are the changes needed?
    
    ### Does this PR introduce _any_ user-facing change?
    
    ### How was this patch tested?
    
    Closes #2229 from AngersZhuuuu/CELEBORN-1226.
    
    Authored-by: Angerszhuuuu <[email protected]>
    Signed-off-by: Angerszhuuuu <[email protected]>
---
 .../apache/celeborn/client/write/DataPusher.java   |   2 +
 .../deploy/worker/storage/CreditStreamManager.java |  29 +++--
 .../worker/storage/PartitionFilesSorter.java       | 118 ++++++++++-----------
 .../celeborn/service/deploy/worker/Worker.scala    |  10 +-
 .../service/deploy/worker/storage/Flusher.scala    |  32 +-----
 5 files changed, 82 insertions(+), 109 deletions(-)

diff --git 
a/client/src/main/java/org/apache/celeborn/client/write/DataPusher.java 
b/client/src/main/java/org/apache/celeborn/client/write/DataPusher.java
index a4932674e..25b26f4ad 100644
--- a/client/src/main/java/org/apache/celeborn/client/write/DataPusher.java
+++ b/client/src/main/java/org/apache/celeborn/client/write/DataPusher.java
@@ -34,6 +34,7 @@ import org.slf4j.LoggerFactory;
 import org.apache.celeborn.client.ShuffleClient;
 import org.apache.celeborn.common.CelebornConf;
 import org.apache.celeborn.common.exception.CelebornIOException;
+import org.apache.celeborn.common.util.ThreadExceptionHandler;
 
 public class DataPusher {
   private static final Logger logger = 
LoggerFactory.getLogger(DataPusher.class);
@@ -137,6 +138,7 @@ public class DataPusher {
           }
         };
     pushThread.setDaemon(true);
+    pushThread.setUncaughtExceptionHandler(new 
ThreadExceptionHandler("DataPusher-" + taskId));
     pushThread.start();
   }
 
diff --git 
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/CreditStreamManager.java
 
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/CreditStreamManager.java
index 22b2f10a4..7e2c6ee3a 100644
--- 
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/CreditStreamManager.java
+++ 
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/CreditStreamManager.java
@@ -35,6 +35,7 @@ import org.apache.celeborn.common.meta.DiskFileInfo;
 import org.apache.celeborn.common.meta.FileInfo;
 import org.apache.celeborn.common.meta.MapFileMeta;
 import org.apache.celeborn.common.util.JavaUtils;
+import org.apache.celeborn.common.util.ThreadUtils;
 import org.apache.celeborn.service.deploy.worker.memory.MemoryManager;
 
 public class CreditStreamManager {
@@ -50,7 +51,7 @@ public class CreditStreamManager {
   private final BlockingQueue<DelayedStreamId> recycleStreamIds = new 
DelayQueue<>();
 
   @GuardedBy("lock")
-  private volatile Thread recycleThread;
+  private volatile ExecutorService recycleThread;
 
   private final Object lock = new Object();
 
@@ -189,20 +190,18 @@ public class CreditStreamManager {
     synchronized (lock) {
       if (recycleThread == null) {
         recycleThread =
-            new Thread(
-                () -> {
-                  while (true) {
-                    try {
-                      DelayedStreamId delayedStreamId = 
recycleStreamIds.take();
-                      cleanResource(delayedStreamId.streamId);
-                    } catch (Throwable e) {
-                      logger.warn(e.getMessage(), e);
-                    }
-                  }
-                },
-                "recycle-thread");
-        recycleThread.setDaemon(true);
-        recycleThread.start();
+            
ThreadUtils.newDaemonSingleThreadExecutor("credit-stream-manager-recycle-thread");
+        recycleThread.submit(
+            () -> {
+              while (true) {
+                try {
+                  DelayedStreamId delayedStreamId = recycleStreamIds.take();
+                  cleanResource(delayedStreamId.streamId);
+                } catch (Throwable e) {
+                  logger.warn(e.getMessage(), e);
+                }
+              }
+            });
 
         logger.info("start stream recycle thread");
       }
diff --git 
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
 
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
index 967f0db60..5eaf05f99 100644
--- 
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
+++ 
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
@@ -94,7 +94,7 @@ public class PartitionFilesSorter extends 
ShuffleRecoverHelper {
   protected final AbstractSource source;
 
   private final ExecutorService fileSorterExecutors;
-  private final Thread fileSorterSchedulerThread;
+  private final ExecutorService fileSorterSchedulerThread;
   private final long indexCacheMaxWeight;
 
   public PartitionFilesSorter(
@@ -143,32 +143,31 @@ public class PartitionFilesSorter extends 
ShuffleRecoverHelper {
             .build();
 
     fileSorterSchedulerThread =
-        new Thread(
-            () -> {
-              try {
-                while (!shutdown) {
-                  FileSorter task = shuffleSortTaskDeque.take();
-                  memoryManager.reserveSortMemory(reservedMemoryPerPartition);
-                  while (!memoryManager.sortMemoryReady()) {
-                    Thread.sleep(20);
-                  }
-                  fileSorterExecutors.submit(
-                      () -> {
-                        try {
-                          task.sort();
-                        } catch (InterruptedException e) {
-                          logger.warn("File sorter thread was interrupted.");
-                        } finally {
-                          
memoryManager.releaseSortMemory(reservedMemoryPerPartition);
-                        }
-                      });
-                }
-              } catch (InterruptedException e) {
-                logger.warn("Sort scheduler thread is shutting down, detail: 
", e);
+        
ThreadUtils.newDaemonSingleThreadExecutor("worker-file-sorter-scheduler");
+    fileSorterSchedulerThread.submit(
+        () -> {
+          try {
+            while (!shutdown) {
+              FileSorter task = shuffleSortTaskDeque.take();
+              memoryManager.reserveSortMemory(reservedMemoryPerPartition);
+              while (!memoryManager.sortMemoryReady()) {
+                Thread.sleep(20);
               }
-            });
-    fileSorterSchedulerThread.setDaemon(true);
-    fileSorterSchedulerThread.start();
+              fileSorterExecutors.submit(
+                  () -> {
+                    try {
+                      task.sort();
+                    } catch (InterruptedException e) {
+                      logger.warn("File sorter thread was interrupted.");
+                    } finally {
+                      
memoryManager.releaseSortMemory(reservedMemoryPerPartition);
+                    }
+                  });
+            }
+          } catch (InterruptedException e) {
+            logger.warn("Sort scheduler thread is shutting down, detail: ", e);
+          }
+        });
   }
 
   public int getSortingCount() {
@@ -301,7 +300,7 @@ public class PartitionFilesSorter extends 
ShuffleRecoverHelper {
       long end = System.currentTimeMillis();
       logger.info("Await partition sorter executor complete cost " + (end - 
start) + "ms");
     } else {
-      fileSorterSchedulerThread.interrupt();
+      fileSorterSchedulerThread.shutdownNow();
       fileSorterExecutors.shutdownNow();
       cleaner.close();
       if (sortedFilesDb != null) {
@@ -761,47 +760,44 @@ class PartitionFilesCleaner {
       new LinkedBlockingQueue<>();
   private final Lock lock = new ReentrantLock();
   private final Condition notEmpty = lock.newCondition();
-  private final Thread cleaner;
+  private final ExecutorService cleaner =
+      
ThreadUtils.newDaemonSingleThreadExecutor("worker-partition-file-cleaner");
 
   PartitionFilesCleaner(PartitionFilesSorter partitionFilesSorter) {
-    cleaner =
-        new Thread(
-            () -> {
+    cleaner.submit(
+        () -> {
+          try {
+            while (!partitionFilesSorter.isShutdown()) {
+              lock.lockInterruptibly();
               try {
-                while (!partitionFilesSorter.isShutdown()) {
-                  lock.lockInterruptibly();
+                // CELEBORN-1210: use while instead of if in case of spurious 
wakeup.
+                while (fileSorters.isEmpty()) {
+                  notEmpty.await();
+                }
+                Iterator<PartitionFilesSorter.FileSorter> it = 
fileSorters.iterator();
+                while (it.hasNext()) {
+                  PartitionFilesSorter.FileSorter sorter = it.next();
                   try {
-                    // CELEBORN-1210: use while instead of if in case of 
spurious wakeup.
-                    while (fileSorters.isEmpty()) {
-                      notEmpty.await();
+                    if (((DiskFileInfo) 
sorter.getOriginFileInfo()).isStreamsEmpty()) {
+                      logger.debug(
+                          "Deleting the original files for shuffle key {}: {}",
+                          sorter.getShuffleKey(),
+                          ((DiskFileInfo) 
sorter.getOriginFileInfo()).getFilePath());
+                      sorter.deleteOriginFiles();
+                      it.remove();
                     }
-                    Iterator<PartitionFilesSorter.FileSorter> it = 
fileSorters.iterator();
-                    while (it.hasNext()) {
-                      PartitionFilesSorter.FileSorter sorter = it.next();
-                      try {
-                        if (((DiskFileInfo) 
sorter.getOriginFileInfo()).isStreamsEmpty()) {
-                          logger.debug(
-                              "Deleting the original files for shuffle key {}: 
{}",
-                              sorter.getShuffleKey(),
-                              ((DiskFileInfo) 
sorter.getOriginFileInfo()).getFilePath());
-                          sorter.deleteOriginFiles();
-                          it.remove();
-                        }
-                      } catch (IOException e) {
-                        logger.error("catch IOException when delete origin 
files", e);
-                      }
-                    }
-                  } finally {
-                    lock.unlock();
+                  } catch (IOException e) {
+                    logger.error("catch IOException when delete origin files", 
e);
                   }
                 }
-              } catch (InterruptedException e) {
-                logger.warn("partition file cleaner thread interrupted while 
wait new sorter.", e);
+              } finally {
+                lock.unlock();
               }
-            });
-    cleaner.setName("partition-files-cleaner");
-    cleaner.setDaemon(true);
-    cleaner.start();
+            }
+          } catch (InterruptedException e) {
+            logger.warn("Partition file cleaner thread interrupted while 
waiting new sorter.", e);
+          }
+        });
   }
 
   public void add(PartitionFilesSorter.FileSorter fileSorter) throws 
InterruptedException {
@@ -825,6 +821,6 @@ class PartitionFilesCleaner {
 
   public void close() {
     fileSorters.clear();
-    cleaner.interrupt();
+    cleaner.shutdownNow();
   }
 }
diff --git 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
index 643a077c8..6ca5e7d21 100644
--- 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
+++ 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
@@ -270,7 +270,8 @@ private[celeborn] class Worker(
   private val replicaFastFailDuration = conf.workerReplicateFastFailDuration
 
   private val cleanTaskQueue = new LinkedBlockingQueue[JHashSet[String]]
-  var cleaner: Thread = _
+  var cleaner: ExecutorService =
+    ThreadUtils.newDaemonSingleThreadExecutor("worker-cleaner")
 
   private val workerResourceConsumptionInterval = 
conf.workerResourceConsumptionInterval
   private val userResourceConsumptions =
@@ -414,7 +415,7 @@ private[celeborn] class Worker(
       replicaFastFailDuration,
       TimeUnit.MILLISECONDS)
 
-    cleaner = new Thread("Cleaner") {
+    cleaner.submit(new Runnable {
       override def run(): Unit = {
         while (true) {
           val expiredShuffleKeys = cleanTaskQueue.take()
@@ -426,16 +427,13 @@ private[celeborn] class Worker(
           }
         }
       }
-    }
+    })
 
     pushDataHandler.init(this)
     replicateHandler.init(this)
     fetchHandler.init(this)
     controller.init(this)
 
-    cleaner.setDaemon(true)
-    cleaner.start()
-
     logInfo("Worker started.")
     rpcEnv.awaitTermination()
   }
diff --git 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/Flusher.scala
 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/Flusher.scala
index de3571791..caa758ae8 100644
--- 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/Flusher.scala
+++ 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/Flusher.scala
@@ -19,7 +19,7 @@ package org.apache.celeborn.service.deploy.worker.storage
 
 import java.io.IOException
 import java.nio.channels.ClosedByInterruptException
-import java.util.concurrent.{LinkedBlockingQueue, TimeUnit}
+import java.util.concurrent.{ExecutorService, LinkedBlockingQueue, TimeUnit}
 import java.util.concurrent.atomic.{AtomicBoolean, AtomicLongArray}
 
 import scala.collection.JavaConverters._
@@ -31,6 +31,7 @@ import org.apache.celeborn.common.internal.Logging
 import org.apache.celeborn.common.meta.{DiskStatus, TimeWindow}
 import org.apache.celeborn.common.metrics.source.AbstractSource
 import org.apache.celeborn.common.protocol.StorageInfo
+import org.apache.celeborn.common.util.ThreadUtils
 import org.apache.celeborn.service.deploy.worker.WorkerSource
 import 
org.apache.celeborn.service.deploy.worker.congestcontrol.CongestionController
 import org.apache.celeborn.service.deploy.worker.memory.MemoryManager
@@ -44,7 +45,7 @@ abstract private[worker] class Flusher(
   protected lazy val flusherId: Int = System.identityHashCode(this)
   protected val workingQueues = new 
Array[LinkedBlockingQueue[FlushTask]](threadCount)
   protected val bufferQueue = new LinkedBlockingQueue[CompositeByteBuf]()
-  protected val workers = new Array[Thread](threadCount)
+  protected val workers = new Array[ExecutorService](threadCount)
   protected var nextWorkerIndex: Int = 0
 
   val lastBeginFlushTime: AtomicLongArray = new AtomicLongArray(threadCount)
@@ -58,7 +59,8 @@ abstract private[worker] class Flusher(
     }
     for (index <- 0 until threadCount) {
       workingQueues(index) = new LinkedBlockingQueue[FlushTask]()
-      workers(index) = new Thread(s"$this-$index") {
+      workers(index) = 
ThreadUtils.newDaemonSingleThreadExecutor(s"$this-$index")
+      workers(index).submit(new Runnable {
         override def run(): Unit = {
           while (!stopFlag.get()) {
             val task = workingQueues(index).take()
@@ -86,14 +88,7 @@ abstract private[worker] class Flusher(
             }
           }
         }
-      }
-      workers(index).setDaemon(true)
-      workers(index).setUncaughtExceptionHandler(new 
Thread.UncaughtExceptionHandler {
-        override def uncaughtException(t: Thread, e: Throwable): Unit = {
-          logError(s"$this thread terminated.", e)
-        }
       })
-      workers(index).start()
     }
   }
 
@@ -125,23 +120,6 @@ abstract private[worker] class Flusher(
     workingQueues(workerIndex).offer(task, timeoutMs, TimeUnit.MILLISECONDS)
   }
 
-  def bufferQueueInfo(): String = s"$this used buffers: ${bufferQueue.size()}"
-
-  def stopAndCleanFlusher(): Unit = {
-    stopFlag.set(true)
-    try {
-      workers.foreach(_.interrupt())
-    } catch {
-      case e: Exception =>
-        logError(s"Exception when interrupt worker: ${workers.mkString(",")}, 
$e")
-    }
-    workingQueues.foreach { queue =>
-      queue.asScala.foreach { task =>
-        returnBuffer(task.buffer)
-      }
-    }
-  }
-
   def processIOException(e: IOException, deviceErrorType: DiskStatus): Unit
 }
 

Reply via email to