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 | <value of
celeborn.client.excludedWorker.expireTimeout> | 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 | <value of
celeborn.client.excludedWorker.expireTimeout> | 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 | <undefined> | 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))