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
}