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 33cf343d2 [CELEBORN-666][REFACTOR] Unify exclude and blacklist related 
configuration
33cf343d2 is described below

commit 33cf343d201fda18db397119ab7cffe0a19d7db9
Author: Angerszhuuuu <[email protected]>
AuthorDate: Wed Jun 28 10:59:58 2023 +0800

    [CELEBORN-666][REFACTOR] Unify exclude and blacklist related configuration
    
    ### What changes were proposed in this pull request?
    Unify exclude and blacklist related configuration
    
    ### Why are the changes needed?
    
    ### Does this PR introduce _any_ user-facing change?
    
    ### How was this patch tested?
    
    Closes #1633 from AngersZhuuuu/CELEBORN-666-NEW.
    
    Authored-by: Angerszhuuuu <[email protected]>
    Signed-off-by: Angerszhuuuu <[email protected]>
---
 .../apache/celeborn/client/ShuffleClientImpl.java  | 33 +++++++++++-----------
 .../celeborn/client/read/RssInputStream.java       |  6 ++--
 .../celeborn/client/WorkerStatusTracker.scala      |  8 +++---
 .../org/apache/celeborn/common/CelebornConf.scala  | 19 +++++++------
 docs/configuration/client.md                       |  6 ++--
 .../celeborn/tests/spark/PushDataTimeoutTest.scala |  6 ++--
 6 files changed, 40 insertions(+), 38 deletions(-)

diff --git 
a/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java 
b/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
index 70dd9cb59..94d3df4a2 100644
--- a/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
+++ b/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
@@ -100,8 +100,8 @@ public class ShuffleClientImpl extends ShuffleClient {
   // key: shuffleId-mapId-attemptId
   protected final Map<String, PushState> pushStates = 
JavaUtils.newConcurrentHashMap();
 
-  private final boolean shuffleClientPushBlacklistEnabled;
-  private final Set<String> blacklist = ConcurrentHashMap.newKeySet();
+  private final boolean pushExcludeWorkerOnFailureEnabled;
+  private final Set<String> pushExcludedWorkers = 
ConcurrentHashMap.newKeySet();
   private final ConcurrentHashMap<String, Long> fetchExcludedWorkers =
       JavaUtils.newConcurrentHashMap();
 
@@ -163,7 +163,7 @@ public class ShuffleClientImpl extends ShuffleClient {
     maxReviveTimes = conf.clientPushMaxReviveTimes();
     testRetryRevive = conf.testRetryRevive();
     pushBufferMaxSize = conf.clientPushBufferMaxSize();
-    shuffleClientPushBlacklistEnabled = conf.clientPushBlacklistEnabled();
+    pushExcludeWorkerOnFailureEnabled = 
conf.clientPushExcludeWorkerOnFailureEnabled();
     if (conf.clientPushReplicateEnabled()) {
       pushDataTimeout = conf.pushDataTimeoutMs() * 2;
     } else {
@@ -197,10 +197,11 @@ public class ShuffleClientImpl extends ShuffleClient {
   private boolean checkPushBlacklisted(
       PartitionLocation location, RpcResponseCallback wrappedCallback) {
     // If shuffleClientBlacklistEnabled = false, blacklist should be empty.
-    if (blacklist.contains(location.hostAndPushPort())) {
+    if (pushExcludedWorkers.contains(location.hostAndPushPort())) {
       wrappedCallback.onFailure(new 
CelebornIOException(StatusCode.PUSH_DATA_MASTER_BLACKLISTED));
       return true;
-    } else if (location.hasPeer() && 
blacklist.contains(location.getPeer().hostAndPushPort())) {
+    } else if (location.hasPeer()
+        && pushExcludedWorkers.contains(location.getPeer().hostAndPushPort())) 
{
       wrappedCallback.onFailure(new 
CelebornIOException(StatusCode.PUSH_DATA_SLAVE_BLACKLISTED));
       return true;
     } else {
@@ -526,7 +527,7 @@ public class ShuffleClientImpl extends ShuffleClient {
           for (int i = 0; i < response.getPartitionLocationsList().size(); 
i++) {
             PartitionLocation partitionLoc =
                 
PbSerDeUtils.fromPbPartitionLocation(response.getPartitionLocationsList().get(i));
-            blacklist.remove(partitionLoc.hostAndPushPort());
+            pushExcludedWorkers.remove(partitionLoc.hostAndPushPort());
 
             result.put(partitionLoc.getId(), partitionLoc);
           }
@@ -621,19 +622,19 @@ public class ShuffleClientImpl extends ShuffleClient {
 
   void excludeWorkerByCause(StatusCode cause, PartitionLocation oldLocation) {
     // Add ShuffleClient side blacklist
-    if (shuffleClientPushBlacklistEnabled && oldLocation != null) {
+    if (pushExcludeWorkerOnFailureEnabled && oldLocation != null) {
       if (cause == StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_MASTER) {
-        blacklist.add(oldLocation.hostAndPushPort());
+        pushExcludedWorkers.add(oldLocation.hostAndPushPort());
       } else if (cause == StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_MASTER) {
-        blacklist.add(oldLocation.hostAndPushPort());
+        pushExcludedWorkers.add(oldLocation.hostAndPushPort());
       } else if (cause == StatusCode.PUSH_DATA_TIMEOUT_MASTER) {
-        blacklist.add(oldLocation.hostAndPushPort());
+        pushExcludedWorkers.add(oldLocation.hostAndPushPort());
       } else if (cause == StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_SLAVE) {
-        blacklist.add(oldLocation.getPeer().hostAndPushPort());
+        pushExcludedWorkers.add(oldLocation.getPeer().hostAndPushPort());
       } else if (cause == StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_SLAVE) {
-        blacklist.add(oldLocation.getPeer().hostAndPushPort());
+        pushExcludedWorkers.add(oldLocation.getPeer().hostAndPushPort());
       } else if (cause == StatusCode.PUSH_DATA_TIMEOUT_SLAVE) {
-        blacklist.add(oldLocation.getPeer().hostAndPushPort());
+        pushExcludedWorkers.add(oldLocation.getPeer().hostAndPushPort());
       }
     }
   }
@@ -704,14 +705,14 @@ public class ShuffleClientImpl extends ShuffleClient {
         int partitionId = partitionInfo.getPartitionId();
         int statusCode = partitionInfo.getStatus();
         if (partitionInfo.getOldAvailable()) {
-          blacklist.remove(oldLocMap.get(partitionId).hostAndPushPort());
+          
pushExcludedWorkers.remove(oldLocMap.get(partitionId).hostAndPushPort());
         }
 
         if (StatusCode.SUCCESS.getValue() == statusCode) {
           PartitionLocation loc =
               
PbSerDeUtils.fromPbPartitionLocation(partitionInfo.getPartition());
           partitionLocationMap.put(partitionId, loc);
-          blacklist.remove(loc.hostAndPushPort());
+          pushExcludedWorkers.remove(loc.hostAndPushPort());
         } else if (StatusCode.STAGE_ENDED.getValue() == statusCode) {
           stageEnded(shuffleId);
           return results;
@@ -1618,7 +1619,7 @@ public class ShuffleClientImpl extends ShuffleClient {
     if (null != driverRssMetaService) {
       driverRssMetaService = null;
     }
-    blacklist.clear();
+    pushExcludedWorkers.clear();
     fetchExcludedWorkers.clear();
     logger.warn("Shuffle client has been shutdown!");
   }
diff --git 
a/client/src/main/java/org/apache/celeborn/client/read/RssInputStream.java 
b/client/src/main/java/org/apache/celeborn/client/read/RssInputStream.java
index 5f88aef9c..cf95cbd2d 100644
--- a/client/src/main/java/org/apache/celeborn/client/read/RssInputStream.java
+++ b/client/src/main/java/org/apache/celeborn/client/read/RssInputStream.java
@@ -130,7 +130,7 @@ public abstract class RssInputStream extends InputStream {
     private final boolean rangeReadFilter;
 
     private boolean pushReplicateEnabled;
-    private boolean fetchBlacklistEnabled;
+    private boolean fetchExcludeWorkerOnFailureEnabled;
     private long fetchExcludedWorkerExpireTimeout;
     private final ConcurrentHashMap<String, Long> fetchExcludedWorkers;
 
@@ -155,7 +155,7 @@ public abstract class RssInputStream extends InputStream {
       this.endMapIndex = endMapIndex;
       this.rangeReadFilter = conf.shuffleRangeReadFilterEnabled();
       this.pushReplicateEnabled = conf.clientPushReplicateEnabled();
-      this.fetchBlacklistEnabled = 
conf.clientFetchExcludeWorkerOnFailureEnabled();
+      this.fetchExcludeWorkerOnFailureEnabled = 
conf.clientFetchExcludeWorkerOnFailureEnabled();
       this.fetchExcludedWorkerExpireTimeout = 
conf.clientFetchExcludedWorkerExpireTimeout();
       this.fetchExcludedWorkers = fetchExcludedWorkers;
 
@@ -241,7 +241,7 @@ public abstract class RssInputStream extends InputStream {
     }
 
     private void excludeFailedLocation(PartitionLocation location, Exception 
e) {
-      if (pushReplicateEnabled && fetchBlacklistEnabled && isCriticalCause(e)) 
{
+      if (pushReplicateEnabled && fetchExcludeWorkerOnFailureEnabled && 
isCriticalCause(e)) {
         fetchExcludedWorkers.put(location.hostAndFetchPort(), 
System.currentTimeMillis());
       }
     }
diff --git 
a/client/src/main/scala/org/apache/celeborn/client/WorkerStatusTracker.scala 
b/client/src/main/scala/org/apache/celeborn/client/WorkerStatusTracker.scala
index 74cdeca2a..5f671080a 100644
--- a/client/src/main/scala/org/apache/celeborn/client/WorkerStatusTracker.scala
+++ b/client/src/main/scala/org/apache/celeborn/client/WorkerStatusTracker.scala
@@ -85,26 +85,26 @@ class WorkerStatusTracker(
         case StatusCode.PUSH_DATA_WRITE_FAIL_MASTER =>
           blacklistWorker(oldPartition, StatusCode.PUSH_DATA_WRITE_FAIL_MASTER)
         case StatusCode.PUSH_DATA_WRITE_FAIL_SLAVE
-            if oldPartition.hasPeer && conf.clientBlacklistSlaveEnabled =>
+            if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled 
=>
           blacklistWorker(oldPartition.getPeer, 
StatusCode.PUSH_DATA_WRITE_FAIL_SLAVE)
         case StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_MASTER =>
           blacklistWorker(oldPartition, 
StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_MASTER)
         case StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_SLAVE
-            if oldPartition.hasPeer && conf.clientBlacklistSlaveEnabled =>
+            if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled 
=>
           blacklistWorker(
             oldPartition.getPeer,
             StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_SLAVE)
         case StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_MASTER =>
           blacklistWorker(oldPartition, 
StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_MASTER)
         case StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_SLAVE
-            if oldPartition.hasPeer && conf.clientBlacklistSlaveEnabled =>
+            if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled 
=>
           blacklistWorker(
             oldPartition.getPeer,
             StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_SLAVE)
         case StatusCode.PUSH_DATA_TIMEOUT_MASTER =>
           blacklistWorker(oldPartition, StatusCode.PUSH_DATA_TIMEOUT_MASTER)
         case StatusCode.PUSH_DATA_TIMEOUT_SLAVE
-            if oldPartition.hasPeer && conf.clientBlacklistSlaveEnabled =>
+            if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled 
=>
           blacklistWorker(
             oldPartition.getPeer,
             StatusCode.PUSH_DATA_TIMEOUT_SLAVE)
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala 
b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
index c86e9da7f..1b2af255c 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -695,8 +695,7 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable 
with Logging with Se
   def appHeartbeatIntervalMs: Long = get(APPLICATION_HEARTBEAT_INTERVAL)
   def clientCheckedUseAllocatedWorkers: Boolean = 
get(CLIENT_CHECKED_USE_ALLOCATED_WORKERS)
   def clientExcludedWorkerExpireTimeout: Long = 
get(CLIENT_EXCLUDED_WORKER_EXPIRE_TIMEOUT)
-  def clientBlacklistSlaveEnabled: Boolean = 
get(CLIENT_BLACKLIST_SLAVE_ENABLED)
-  def clientPushBlacklistEnabled: Boolean = get(CLIENT_PUSH_BLACKLIST_ENABLED)
+  def clientExcludeSlaveOnFailureEnabled: Boolean = 
get(CLIENT_EXCLUDE_SLAVE_ON_FAILURE_ENABLED)
 
   // //////////////////////////////////////////////////////
   //               Shuffle Compression                   //
@@ -748,6 +747,8 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable 
with Logging with Se
   def clientPushBufferInitialSize: Int = 
get(CLIENT_PUSH_BUFFER_INITIAL_SIZE).toInt
   def clientPushBufferMaxSize: Int = get(CLIENT_PUSH_BUFFER_MAX_SIZE).toInt
   def clientPushQueueCapacity: Int = get(CLIENT_PUSH_QUEUE_CAPACITY)
+  def clientPushExcludeWorkerOnFailureEnabled: Boolean =
+    get(CLIENT_PUSH_EXCLUDE_WORKER_ON_FAILURE_ENABLED)
   def clientPushMaxReqsInFlight: Int = get(CLIENT_PUSH_MAX_REQS_IN_FLIGHT)
   def clientPushMaxReviveTimes: Int = get(CLIENT_PUSH_MAX_REVIVE_TIMES)
   def clientPushReviveInterval: Long = get(CLIENT_PUSH_REVIVE_INTERVAL)
@@ -2574,11 +2575,11 @@ object CelebornConf extends Logging {
       .timeConf(TimeUnit.MILLISECONDS)
       .createWithDefaultString("10s")
 
-  val CLIENT_BLACKLIST_SLAVE_ENABLED: ConfigEntry[Boolean] =
-    buildConf("celeborn.client.blacklistSlave.enabled")
+  val CLIENT_EXCLUDE_SLAVE_ON_FAILURE_ENABLED: ConfigEntry[Boolean] =
+    buildConf("celeborn.client.excludeSlaveOnFailure.enabled")
       .categories("client")
       .version("0.3.0")
-      .doc("When true, Celeborn will add partition's peer worker into 
blacklist " +
+      .doc("When true, Celeborn will exclude partition's peer worker on 
failure " +
         "when push data to slave failed.")
       .booleanConf
       .createWithDefault(true)
@@ -2693,10 +2694,10 @@ object CelebornConf extends Logging {
       .intConf
       .createWithDefault(2048)
 
-  val CLIENT_PUSH_BLACKLIST_ENABLED: ConfigEntry[Boolean] =
-    buildConf("celeborn.client.push.blacklist.enabled")
+  val CLIENT_PUSH_EXCLUDE_WORKER_ON_FAILURE_ENABLED: ConfigEntry[Boolean] =
+    buildConf("celeborn.client.push.excludeWorkerOnFailure.enabled")
       .categories("client")
-      .doc("Whether to enable shuffle client-side push blacklist of workers.")
+      .doc("Whether to enable shuffle client-side push exclude workers on 
failures.")
       .version("0.3.0")
       .booleanConf
       .createWithDefault(false)
@@ -2874,7 +2875,7 @@ object CelebornConf extends Logging {
     buildConf("celeborn.client.fetch.excludedWorker.expireTimeout")
       .categories("client")
       .doc("ShuffleClient is a static object, it will be used in the whole 
lifecycle of Executor," +
-        "We give a expire time for blacklisted worker to avoid a transient 
worker issues.")
+        "We give a expire time for excluded workers to avoid a transient 
worker issues.")
       .version("0.3.0")
       .fallbackConf(CLIENT_EXCLUDED_WORKER_EXPIRE_TIMEOUT)
 
diff --git a/docs/configuration/client.md b/docs/configuration/client.md
index 9d9e891e8..0bc1f79bb 100644
--- a/docs/configuration/client.md
+++ b/docs/configuration/client.md
@@ -20,12 +20,12 @@ license: |
 | Key | Default | Description | Since |
 | --- | ------- | ----------- | ----- |
 | celeborn.client.application.heartbeatInterval | 10s | Interval for client to 
send heartbeat message to master. | 0.3.0 | 
-| celeborn.client.blacklistSlave.enabled | true | When true, Celeborn will add 
partition's peer worker into blacklist when push data to slave failed. | 0.3.0 
| 
 | celeborn.client.closeIdleConnections | true | Whether client will close idle 
connections. | 0.3.0 | 
 | celeborn.client.commitFiles.ignoreExcludedWorker | false | When true, 
LifecycleManager will skip workers which are in the excluded list. | 0.3.0 | 
+| celeborn.client.excludeSlaveOnFailure.enabled | true | When true, Celeborn 
will exclude partition's peer worker on failure when push data to slave failed. 
| 0.3.0 | 
 | celeborn.client.excludedWorker.expireTimeout | 180s | Timeout time for 
LifecycleManager to clear reserved excluded worker. Default to be 1.5 * 
`celeborn.master.heartbeat.worker.timeout`to cover worker heartbeat timeout 
check period | 0.3.0 | 
 | celeborn.client.fetch.excludeWorkerOnFailure.enabled | false | Whether to 
enable shuffle client-side fetch exclude workers on failure. | 0.3.0 | 
-| celeborn.client.fetch.excludedWorker.expireTimeout | &lt;value of 
celeborn.client.excludedWorker.expireTimeout&gt; | ShuffleClient is a static 
object, it will be used in the whole lifecycle of Executor,We give a expire 
time for blacklisted worker to avoid a transient worker issues. | 0.3.0 | 
+| celeborn.client.fetch.excludedWorker.expireTimeout | &lt;value of 
celeborn.client.excludedWorker.expireTimeout&gt; | ShuffleClient is a static 
object, it will be used in the whole lifecycle of Executor,We give a expire 
time for excluded workers to avoid a transient worker issues. | 0.3.0 | 
 | celeborn.client.fetch.maxReqsInFlight | 3 | Amount of in-flight chunk fetch 
request. | 0.3.0 | 
 | celeborn.client.fetch.maxRetriesForEachReplica | 3 | Max retry times of 
fetch chunk on each replica | 0.3.0 | 
 | celeborn.client.fetch.timeout | 600s | Timeout for a task to open stream and 
fetch chunk. | 0.3.0 | 
@@ -37,9 +37,9 @@ license: |
 | celeborn.client.flink.resultPartition.memory | 64m | Memory reserved for a 
result partition. | 0.3.0 | 
 | celeborn.client.flink.resultPartition.minMemory | 8m | Min memory reserved 
for a result partition. | 0.3.0 | 
 | celeborn.client.flink.resultPartition.supportFloatingBuffer | true | Whether 
to support floating buffer for result partitions. | 0.3.0 | 
-| celeborn.client.push.blacklist.enabled | false | Whether to enable shuffle 
client-side push blacklist of workers. | 0.3.0 | 
 | celeborn.client.push.buffer.initial.size | 8k |  | 0.3.0 | 
 | celeborn.client.push.buffer.max.size | 64k | Max size of reducer partition 
buffer memory for shuffle hash writer. The pushed data will be buffered in 
memory before sending to Celeborn worker. For performance consideration keep 
this buffer size higher than 32K. Example: If reducer amount is 2000, buffer 
size is 64K, then each task will consume up to `64KiB * 2000 = 125MiB` heap 
memory. | 0.3.0 | 
+| celeborn.client.push.excludeWorkerOnFailure.enabled | false | Whether to 
enable shuffle client-side push exclude workers on failures. | 0.3.0 | 
 | celeborn.client.push.limit.inFlight.sleepInterval | 50ms | Sleep interval 
when check netty in-flight requests to be done. | 0.3.0 | 
 | celeborn.client.push.limit.inFlight.timeout | &lt;undefined&gt; | Timeout 
for netty in-flight requests to be done.Default value should be 
`celeborn.client.push.timeout * 2`. | 0.3.0 | 
 | celeborn.client.push.limit.strategy | SIMPLE | The strategy used to control 
the push speed. Valid strategies are SIMPLE and SLOWSTART. the SLOWSTART 
strategy is usually cooperate with congest control mechanism in the worker 
side. | 0.3.0 | 
diff --git 
a/tests/spark-it/src/test/scala/org/apache/celeborn/tests/spark/PushDataTimeoutTest.scala
 
b/tests/spark-it/src/test/scala/org/apache/celeborn/tests/spark/PushDataTimeoutTest.scala
index 2ad268592..68e381fc6 100644
--- 
a/tests/spark-it/src/test/scala/org/apache/celeborn/tests/spark/PushDataTimeoutTest.scala
+++ 
b/tests/spark-it/src/test/scala/org/apache/celeborn/tests/spark/PushDataTimeoutTest.scala
@@ -65,7 +65,7 @@ class PushDataTimeoutTest extends AnyFunSuite
         .set(s"spark.${CelebornConf.CLIENT_PUSH_DATA_TIMEOUT.key}", "5s")
         .set(s"spark.celeborn.data.push.timeoutCheck.interval", "2s")
         .set(s"spark.${CelebornConf.CLIENT_PUSH_REPLICATE_ENABLED.key}", 
enabled.toString)
-        .set(s"spark.${CelebornConf.CLIENT_BLACKLIST_SLAVE_ENABLED.key}", 
"false")
+        
.set(s"spark.${CelebornConf.CLIENT_EXCLUDE_SLAVE_ON_FAILURE_ENABLED.key}", 
"false")
         // make sure PushDataHandler.handlePushData be triggered
         .set(s"spark.${CelebornConf.CLIENT_PUSH_BUFFER_MAX_SIZE.key}", "5")
 
@@ -97,7 +97,7 @@ class PushDataTimeoutTest extends AnyFunSuite
         .set(s"spark.${CelebornConf.CLIENT_PUSH_DATA_TIMEOUT.key}", "5s")
         .set(s"spark.celeborn.data.push.timeoutCheck.interval", "2s")
         .set(s"spark.${CelebornConf.CLIENT_PUSH_REPLICATE_ENABLED.key}", 
enabled.toString)
-        .set(s"spark.${CelebornConf.CLIENT_BLACKLIST_SLAVE_ENABLED.key}", 
"false")
+        
.set(s"spark.${CelebornConf.CLIENT_EXCLUDE_SLAVE_ON_FAILURE_ENABLED.key}", 
"false")
 
       val sparkSession = SparkSession.builder().config(sparkConf).getOrCreate()
       val sqlResult = runsql(sparkSession)
@@ -123,7 +123,7 @@ class PushDataTimeoutTest extends AnyFunSuite
   test("celeborn spark integration test - pushdata timeout will add to 
blacklist") {
     val sparkConf = new 
SparkConf().setAppName("rss-demo").setMaster("local[2]")
       .set(s"spark.${CelebornConf.CLIENT_PUSH_DATA_TIMEOUT.key}", "5s")
-      .set(s"spark.${CelebornConf.CLIENT_BLACKLIST_SLAVE_ENABLED.key}", "true")
+      
.set(s"spark.${CelebornConf.CLIENT_EXCLUDE_SLAVE_ON_FAILURE_ENABLED.key}", 
"true")
       .set(s"spark.${CelebornConf.CLIENT_PUSH_REPLICATE_ENABLED.key}", "true")
     val rssSparkSession = SparkSession.builder()
       .config(updateSparkConf(sparkConf, ShuffleMode.HASH))

Reply via email to