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 aad392901 [CELEBORN-1259] Improve the default gracePeriod of
ThreadUtils#shutdown
aad392901 is described below
commit aad3929018dfc3ba6a533152a75e16b462a886ae
Author: SteNicholas <[email protected]>
AuthorDate: Thu Feb 1 18:13:36 2024 +0800
[CELEBORN-1259] Improve the default gracePeriod of ThreadUtils#shutdown
### What changes were proposed in this pull request?
Introduce `ThreadUtils#shutdown(executor)` method to improve the default
gracePeriod of `ThreadUtils#shutdown`.
### Why are the changes needed?
The default value of `gracePeriod` for `ThreadUtils#shutdown` is 30 seconds
at present. Meanwhile, the `gracePeriod` of most invoker for
`ThreadUtils#shutdown` is 800 milliseconds. Therefore, the default
`gracePeriod` of `ThreadUtils#shutdown` could be improved as 800 milliseconds.
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
No.
Closes #2276 from SteNicholas/CELEBORN-1259.
Authored-by: SteNicholas <[email protected]>
Signed-off-by: mingji <[email protected]>
---
client/src/main/java/org/apache/celeborn/client/ReviveManager.java | 4 +---
.../scala/org/apache/celeborn/client/ApplicationHeartbeater.scala | 3 +--
.../scala/org/apache/celeborn/client/ChangePartitionManager.scala | 3 +--
.../src/main/scala/org/apache/celeborn/client/CommitManager.scala | 3 +--
.../main/scala/org/apache/celeborn/client/LifecycleManager.scala | 2 +-
.../scala/org/apache/celeborn/client/ReleasePartitionManager.scala | 3 +--
.../main/java/org/apache/celeborn/common/client/MasterClient.java | 3 +--
.../main/scala/org/apache/celeborn/common/util/ThreadUtils.scala | 6 +++++-
.../scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala | 3 +--
.../celeborn/server/common/service/config/FsConfigServiceImpl.java | 4 +---
10 files changed, 14 insertions(+), 20 deletions(-)
diff --git a/client/src/main/java/org/apache/celeborn/client/ReviveManager.java
b/client/src/main/java/org/apache/celeborn/client/ReviveManager.java
index eb07aa375..875952f0d 100644
--- a/client/src/main/java/org/apache/celeborn/client/ReviveManager.java
+++ b/client/src/main/java/org/apache/celeborn/client/ReviveManager.java
@@ -22,8 +22,6 @@ import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
-import scala.concurrent.duration.Duration;
-
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -128,6 +126,6 @@ class ReviveManager {
}
public void close() {
- ThreadUtils.shutdown(batchReviveRequestScheduler, Duration.apply("800ms"));
+ ThreadUtils.shutdown(batchReviveRequestScheduler);
}
}
diff --git
a/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
b/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
index 20a02b207..5101d61f1 100644
---
a/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
+++
b/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
@@ -20,7 +20,6 @@ package org.apache.celeborn.client
import java.util.concurrent.{ScheduledFuture, TimeUnit}
import scala.collection.JavaConverters._
-import scala.concurrent.duration.DurationInt
import org.apache.celeborn.common.CelebornConf
import org.apache.celeborn.common.client.MasterClient
@@ -119,7 +118,7 @@ class ApplicationHeartbeater(
// Stop appHeartbeat first
logInfo(s"Stop Application heartbeat $appId")
appHeartbeat.cancel(true)
- ThreadUtils.shutdown(appHeartbeatHandlerThread, 800.millis)
+ ThreadUtils.shutdown(appHeartbeatHandlerThread)
if (applicationUnregisterEnabled) {
unregisterApplication()
}
diff --git
a/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
b/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
index b565da4a5..39bf7891a 100644
---
a/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
+++
b/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
@@ -22,7 +22,6 @@ import java.util.{Set => JSet}
import java.util.concurrent.{ConcurrentHashMap, ScheduledExecutorService,
ScheduledFuture, TimeUnit}
import scala.collection.JavaConverters._
-import scala.concurrent.duration.DurationInt
import org.apache.celeborn.common.CelebornConf
import org.apache.celeborn.common.internal.Logging
@@ -113,7 +112,7 @@ class ChangePartitionManager(
def stop(): Unit = {
batchHandleChangePartition.foreach(_.cancel(true))
- batchHandleChangePartitionSchedulerThread.foreach(ThreadUtils.shutdown(_,
800.millis))
+ batchHandleChangePartitionSchedulerThread.foreach(ThreadUtils.shutdown(_))
}
private val rpcContextRegisterFunc =
diff --git
a/client/src/main/scala/org/apache/celeborn/client/CommitManager.scala
b/client/src/main/scala/org/apache/celeborn/client/CommitManager.scala
index c1be16644..93178c047 100644
--- a/client/src/main/scala/org/apache/celeborn/client/CommitManager.scala
+++ b/client/src/main/scala/org/apache/celeborn/client/CommitManager.scala
@@ -23,7 +23,6 @@ import java.util.concurrent.atomic.{AtomicInteger, LongAdder}
import scala.collection.JavaConverters._
import scala.collection.mutable.ArrayBuffer
-import scala.concurrent.duration.DurationInt
import org.roaringbitmap.RoaringBitmap
@@ -172,7 +171,7 @@ class CommitManager(appUniqueId: String, val conf:
CelebornConf, lifecycleManage
def stop(): Unit = {
batchHandleCommitPartition.foreach(_.cancel(true))
- batchHandleCommitPartitionSchedulerThread.foreach(ThreadUtils.shutdown(_,
800.millis))
+ batchHandleCommitPartitionSchedulerThread.foreach(ThreadUtils.shutdown(_))
}
def registerShuffle(shuffleId: Int, numMappers: Int): Unit = {
diff --git
a/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
b/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
index ad2a036d0..48b227343 100644
--- a/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
+++ b/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
@@ -202,7 +202,7 @@ class LifecycleManager(val appUniqueId: String, val conf:
CelebornConf) extends
import scala.concurrent.duration._
checkForShuffleRemoval.cancel(true)
- ThreadUtils.shutdown(forwardMessageThread, 800.millis)
+ ThreadUtils.shutdown(forwardMessageThread)
commitManager.stop()
changePartitionManager.stop()
diff --git
a/client/src/main/scala/org/apache/celeborn/client/ReleasePartitionManager.scala
b/client/src/main/scala/org/apache/celeborn/client/ReleasePartitionManager.scala
index 32620d396..ef9b325d4 100644
---
a/client/src/main/scala/org/apache/celeborn/client/ReleasePartitionManager.scala
+++
b/client/src/main/scala/org/apache/celeborn/client/ReleasePartitionManager.scala
@@ -21,7 +21,6 @@ import java.util
import java.util.concurrent.{ScheduledExecutorService, ScheduledFuture,
TimeUnit}
import scala.collection.JavaConverters._
-import scala.concurrent.duration.DurationInt
import org.apache.celeborn.common.CelebornConf
import org.apache.celeborn.common.internal.Logging
@@ -107,7 +106,7 @@ class ReleasePartitionManager(
def stop(): Unit = {
batchHandleReleasePartition.foreach(_.cancel(true))
- batchHandleReleasePartitionSchedulerThread.foreach(ThreadUtils.shutdown(_,
800.millis))
+ batchHandleReleasePartitionSchedulerThread.foreach(ThreadUtils.shutdown(_))
}
def releasePartition(shuffleId: Int, partitionId: Int): Unit = {
diff --git
a/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java
b/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java
index 3c0c90f07..83a727477 100644
--- a/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java
+++ b/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java
@@ -29,7 +29,6 @@ import javax.annotation.Nullable;
import scala.Tuple2;
import scala.concurrent.Future;
-import scala.concurrent.duration.Duration;
import scala.reflect.ClassTag$;
import com.google.common.util.concurrent.Uninterruptibles;
@@ -125,7 +124,7 @@ public class MasterClient {
}
public void close() {
- ThreadUtils.shutdown(oneWayMessageSender, Duration.apply("800ms"));
+ ThreadUtils.shutdown(oneWayMessageSender);
}
@SuppressWarnings("UnstableApiUsage")
diff --git
a/common/src/main/scala/org/apache/celeborn/common/util/ThreadUtils.scala
b/common/src/main/scala/org/apache/celeborn/common/util/ThreadUtils.scala
index fe6605fb1..74ce0ca7d 100644
--- a/common/src/main/scala/org/apache/celeborn/common/util/ThreadUtils.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/util/ThreadUtils.scala
@@ -22,7 +22,7 @@ import java.util.concurrent.{ForkJoinPool => SForkJoinPool,
ForkJoinWorkerThread
import java.util.concurrent.locks.ReentrantLock
import scala.concurrent.{Awaitable, ExecutionContext,
ExecutionContextExecutor, Future}
-import scala.concurrent.duration.{Duration, FiniteDuration}
+import scala.concurrent.duration.{Duration, DurationInt, FiniteDuration}
import scala.language.higherKinds
import scala.util.control.NonFatal
@@ -345,6 +345,10 @@ object ThreadUtils {
}
// scalastyle:on awaitready
+ def shutdown(executor: ExecutorService): Unit = {
+ shutdown(executor, 800.millis)
+ }
+
def shutdown(
executor: ExecutorService,
gracePeriod: Duration = FiniteDuration(30, TimeUnit.SECONDS)): Unit = {
diff --git
a/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
b/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
index 9eede93ed..b74c8e6e7 100644
---
a/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
+++
b/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
@@ -24,7 +24,6 @@ import java.util.concurrent.atomic.AtomicInteger
import scala.collection.JavaConverters._
import scala.collection.mutable.ArrayBuffer
-import scala.concurrent.duration._
import scala.util.Random
import org.junit.Assert.{assertEquals, assertNotEquals, assertNotNull}
@@ -133,7 +132,7 @@ class WorkerInfoSuite extends CelebornFunSuite {
assertEquals(0, allocatedSlots.get())
assertEquals(0, worker.usedSlots())
- ThreadUtils.shutdown(es, 800.millisecond)
+ ThreadUtils.shutdown(es)
}
test("WorkerInfo not equals when host different.") {
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
index 3e160d8e7..cd519c14d 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
@@ -28,8 +28,6 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors;
-import scala.concurrent.duration.Duration;
-
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.yaml.snakeyaml.Yaml;
@@ -112,7 +110,7 @@ public class FsConfigServiceImpl implements ConfigService {
@Override
public void shutdown() {
- ThreadUtils.shutdown(configRefreshService, Duration.apply("800ms"));
+ ThreadUtils.shutdown(configRefreshService);
}
private File getConfigurationFile(Map<String, String> env) {