This is an automated email from the ASF dual-hosted git repository.
zhouky 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 3985a5cbd [CELEBORN-666][FOLLOWUP] Unify all blacklist related code
and comment
3985a5cbd is described below
commit 3985a5cbd7d51503bcd570c7f6af2bd323d125a4
Author: Angerszhuuuu <[email protected]>
AuthorDate: Wed Jun 28 16:28:03 2023 +0800
[CELEBORN-666][FOLLOWUP] Unify all blacklist related code and comment
### What changes were proposed in this pull request?
Unify all blacklist related code and comment
### Why are the changes needed?
### Does this PR introduce _any_ user-facing change?
### How was this patch tested?
Closes #1638 from AngersZhuuuu/CELEBORN-666-FOLLOWUP.
Authored-by: Angerszhuuuu <[email protected]>
Signed-off-by: zky.zhoukeyong <[email protected]>
---
METRICS.md | 2 +-
assets/grafana/rss-dashboard.json | 6 +-
.../apache/celeborn/client/ShuffleClientImpl.java | 25 ++++----
.../celeborn/client/read/RssInputStream.java | 4 +-
.../celeborn/client/ChangePartitionManager.scala | 6 +-
.../apache/celeborn/client/LifecycleManager.scala | 11 ++--
.../celeborn/client/WorkerStatusTracker.scala | 74 +++++++++++-----------
.../celeborn/client/commit/CommitHandler.scala | 8 +--
.../celeborn/client/WorkerStatusTrackerSuite.scala | 42 ++++++------
.../common/protocol/message/StatusCode.java | 8 +--
.../common/protocol/message/ControlMessages.scala | 17 +++--
.../apache/celeborn/common/util/PbSerDeUtils.scala | 4 +-
.../org/apache/celeborn/common/util/Utils.scala | 8 +--
docs/migration.md | 4 ++
docs/monitoring.md | 34 +++++-----
.../master/clustermeta/AbstractMetaManager.java | 32 ++++------
.../celeborn/service/deploy/master/Master.scala | 33 +++++-----
.../service/deploy/master/MasterSource.scala | 2 +-
.../clustermeta/DefaultMetaSystemSuiteJ.java | 8 +--
.../clustermeta/ha/MasterStateMachineSuiteJ.java | 10 +--
.../ha/RatisMasterStatusSystemSuiteJ.java | 34 +++++-----
.../celeborn/server/common/HttpService.scala | 2 +-
.../server/common/http/HttpRequestHandler.scala | 4 +-
.../celeborn/tests/spark/PushDataTimeoutTest.scala | 12 ++--
.../celeborn/service/deploy/worker/Worker.scala | 4 +-
25 files changed, 197 insertions(+), 197 deletions(-)
diff --git a/METRICS.md b/METRICS.md
index 8aa51f070..3cd32313d 100644
--- a/METRICS.md
+++ b/METRICS.md
@@ -64,7 +64,7 @@ Here is an example of grafana dashboard importing.
| MetricName | Role |
Description
|
|:--------------------------------------:|:-----------------:|:---------------------------------------------------------------------------------------------------------------:|
| WorkerCount | master |
The count of active workers.
|
-| BlacklistedWorkerCount | master |
The count of workers in blacklist.
|
+| ExcludedWorkerCount | master |
The count of workers in excluded list.
|
| OfferSlotsTime | master |
The time of offer slots.
|
| PartitionSize | master | The
estimated partition size of last 20 flush window whose length is 15 seconds by
defaults. |
| RegisteredShuffleCount | master and worker |
The value means count of registered shuffle.
|
diff --git a/assets/grafana/rss-dashboard.json
b/assets/grafana/rss-dashboard.json
index 17da13655..2a424a984 100644
--- a/assets/grafana/rss-dashboard.json
+++ b/assets/grafana/rss-dashboard.json
@@ -627,7 +627,7 @@
"type": "prometheus",
"uid": "${DS_PROMETHEUS}"
},
- "description": "The count of workers in blacklist. ",
+ "description": "The count of workers in excluded list. ",
"fieldConfig": {
"defaults": {
"color": {
@@ -704,11 +704,11 @@
"type": "prometheus",
"uid": "${DS_PROMETHEUS}"
},
- "expr": "metrics_BlacklistedWorkerCount_Value",
+ "expr": "metrics_ExcludedWorkerCount_Value",
"refId": "A"
}
],
- "title": "metrics_BlacklistedWorkerCount_Value",
+ "title": "metrics_ExcludedWorkerCount_Value",
"type": "timeseries"
},
{
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 94d3df4a2..0a18dc767 100644
--- a/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
+++ b/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
@@ -194,15 +194,17 @@ public class ShuffleClientImpl extends ShuffleClient {
logger.info("Created ShuffleClientImpl, appUniqueId: {}", appUniqueId);
}
- private boolean checkPushBlacklisted(
+ private boolean isPushTargetWorkerExcluded(
PartitionLocation location, RpcResponseCallback wrappedCallback) {
- // If shuffleClientBlacklistEnabled = false, blacklist should be empty.
+ // If pushExcludeWorkerOnFailureEnabled = false, pushExcludedWorkers
should be empty.
if (pushExcludedWorkers.contains(location.hostAndPushPort())) {
- wrappedCallback.onFailure(new
CelebornIOException(StatusCode.PUSH_DATA_MASTER_BLACKLISTED));
+ wrappedCallback.onFailure(
+ new
CelebornIOException(StatusCode.PUSH_DATA_MASTER_WORKER_EXCLUDED));
return true;
} else if (location.hasPeer()
&& pushExcludedWorkers.contains(location.getPeer().hostAndPushPort()))
{
- wrappedCallback.onFailure(new
CelebornIOException(StatusCode.PUSH_DATA_SLAVE_BLACKLISTED));
+ wrappedCallback.onFailure(
+ new CelebornIOException(StatusCode.PUSH_DATA_SLAVE_WORKER_EXCLUDED));
return true;
} else {
return false;
@@ -270,7 +272,7 @@ public class ShuffleClientImpl extends ShuffleClient {
batchId,
newLoc);
try {
- if (!checkPushBlacklisted(newLoc, wrappedCallback)) {
+ if (!isPushTargetWorkerExcluded(newLoc, wrappedCallback)) {
if (!testRetryRevive || remainReviveTimes < 1) {
TransportClient client =
dataClientFactory.createClient(newLoc.getHost(),
newLoc.getPushPort(), partitionId);
@@ -621,7 +623,6 @@ public class ShuffleClientImpl extends ShuffleClient {
}
void excludeWorkerByCause(StatusCode cause, PartitionLocation oldLocation) {
- // Add ShuffleClient side blacklist
if (pushExcludeWorkerOnFailureEnabled && oldLocation != null) {
if (cause == StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_MASTER) {
pushExcludedWorkers.add(oldLocation.hostAndPushPort());
@@ -1034,7 +1035,7 @@ public class ShuffleClientImpl extends ShuffleClient {
// do push data
try {
- if (!checkPushBlacklisted(loc, wrappedCallback)) {
+ if (!isPushTargetWorkerExcluded(loc, wrappedCallback)) {
if (!testRetryRevive) {
TransportClient client =
dataClientFactory.createClient(loc.getHost(),
loc.getPushPort(), partitionId);
@@ -1411,7 +1412,7 @@ public class ShuffleClientImpl extends ShuffleClient {
// do push merged data
try {
- if (!checkPushBlacklisted(batches.get(0).loc, wrappedCallback)) {
+ if (!isPushTargetWorkerExcluded(batches.get(0).loc, wrappedCallback)) {
if (!testRetryRevive || remainReviveTimes < 1) {
TransportClient client = dataClientFactory.createClient(host, port);
client.pushMergedData(mergedData, pushDataTimeout, wrappedCallback);
@@ -1668,10 +1669,10 @@ public class ShuffleClientImpl extends ShuffleClient {
cause = StatusCode.PUSH_DATA_TIMEOUT_SLAVE;
} else if (message.startsWith(StatusCode.REPLICATE_DATA_FAILED.name())) {
cause = StatusCode.REPLICATE_DATA_FAILED;
- } else if
(message.startsWith(StatusCode.PUSH_DATA_MASTER_BLACKLISTED.name())) {
- cause = StatusCode.PUSH_DATA_MASTER_BLACKLISTED;
- } else if
(message.startsWith(StatusCode.PUSH_DATA_SLAVE_BLACKLISTED.name())) {
- cause = StatusCode.PUSH_DATA_SLAVE_BLACKLISTED;
+ } else if
(message.startsWith(StatusCode.PUSH_DATA_MASTER_WORKER_EXCLUDED.name())) {
+ cause = StatusCode.PUSH_DATA_MASTER_WORKER_EXCLUDED;
+ } else if
(message.startsWith(StatusCode.PUSH_DATA_SLAVE_WORKER_EXCLUDED.name())) {
+ cause = StatusCode.PUSH_DATA_SLAVE_WORKER_EXCLUDED;
} else if (connectFail(message)) {
// Throw when push to master worker connection causeException.
cause = StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_MASTER;
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 cf95cbd2d..e6d3b1df5 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
@@ -288,7 +288,7 @@ public abstract class RssInputStream extends InputStream {
while (fetchChunkRetryCnt < fetchChunkMaxRetry) {
try {
if (isExcluded(location)) {
- throw new CelebornIOException("Fetch data from blacklisted
location! " + location);
+ throw new CelebornIOException("Fetch data from excluded worker! "
+ location);
}
return createReader(location, fetchChunkRetryCnt,
fetchChunkMaxRetry);
} catch (Exception e) {
@@ -326,7 +326,7 @@ public abstract class RssInputStream extends InputStream {
try {
if (isExcluded(currentReader.getLocation())) {
throw new CelebornIOException(
- "Fetch data from blacklisted location! " +
currentReader.getLocation());
+ "Fetch data from excluded worker! " +
currentReader.getLocation());
}
return currentReader.next();
} catch (Exception e) {
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 164d041d9..ac2c5f690 100644
---
a/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
+++
b/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
@@ -203,10 +203,10 @@ class ChangePartitionManager(
}.mkString("[", ",", "]")
logWarning(s"Batch handle change partition for $changes")
- // Blacklist all failed workers
+ // Exclude all failed workers
if (changePartitions.exists(_.causes.isDefined)) {
changePartitions.filter(_.causes.isDefined).foreach { changePartition =>
- lifecycleManager.workerStatusTracker.blacklistWorkerFromPartition(
+ lifecycleManager.workerStatusTracker.excludeWorkerFromPartition(
shuffleId,
changePartition.oldPartition,
changePartition.causes.get)
@@ -253,7 +253,7 @@ class ChangePartitionManager(
}
}
- // Get candidate worker that not in blacklist of shuffleId
+ // Get candidate worker that not in excluded worker list of shuffleId
val candidates =
lifecycleManager
.workerSnapshots(shuffleId)
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 fece54cf2..4b4a01535 100644
--- a/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
+++ b/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
@@ -438,15 +438,15 @@ class LifecycleManager(val appUniqueId: String, val conf:
CelebornConf) extends
logError(s"Init rpc client failed for $shuffleId on $workerInfo
during reserve slots.", t)
connectFailedWorkers.put(
workerInfo,
- (StatusCode.UNKNOWN_WORKER, System.currentTimeMillis()))
+ (StatusCode.WORKER_UNKNOWN, System.currentTimeMillis()))
}
}
candidatesWorkers.removeAll(connectFailedWorkers.asScala.keys.toList.asJava)
workerStatusTracker.recordWorkerFailure(connectFailedWorkers)
// If newly allocated from master and can setup endpoint success,
LifecycleManager should remove worker from
- // the blacklist to improve the accuracy of the blacklist
- workerStatusTracker.removeFromBlacklist(candidatesWorkers)
+ // the excluded worker list to improve the accuracy of the list.
+ workerStatusTracker.removeFromExcludedWorkers(candidatesWorkers)
// Third, for each slot, LifecycleManager should ask Worker to reserve the
slot
// and prepare the pushing data env.
@@ -861,8 +861,9 @@ class LifecycleManager(val appUniqueId: String, val conf:
CelebornConf) extends
val retryCandidates = new util.HashSet(slots.keySet())
// add candidates to avoid revive action passed in slots only 2
worker
retryCandidates.addAll(candidates)
- // remove blacklist from retryCandidates
-
retryCandidates.removeAll(workerStatusTracker.blacklist.keys().asScala.toList.asJava)
+ // remove excluded workers from retryCandidates
+ retryCandidates.removeAll(
+ workerStatusTracker.excludedWorkers.keys().asScala.toList.asJava)
retryCandidates.removeAll(workerStatusTracker.shuttingWorkers.asScala.toList.asJava)
if (retryCandidates.size < 1 || (pushReplicateEnabled &&
retryCandidates.size < 2)) {
logError(s"Retry reserve slots for $shuffleId failed caused by not
enough slots.")
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 5f671080a..1c71c2a45 100644
--- a/client/src/main/scala/org/apache/celeborn/client/WorkerStatusTracker.scala
+++ b/client/src/main/scala/org/apache/celeborn/client/WorkerStatusTracker.scala
@@ -37,8 +37,7 @@ class WorkerStatusTracker(
private val excludedWorkerExpireTimeout =
conf.clientExcludedWorkerExpireTimeout
private val workerStatusListeners =
ConcurrentHashMap.newKeySet[WorkerStatusListener]()
- // blacklist
- val blacklist = new ShuffleFailedWorkers()
+ val excludedWorkers = new ShuffleFailedWorkers()
val shuttingWorkers: JSet[WorkerInfo] = new JHashSet[WorkerInfo]()
def registerWorkerStatusListener(workerStatusListener:
WorkerStatusListener): Unit = {
@@ -49,12 +48,12 @@ class WorkerStatusTracker(
if (conf.clientCheckedUseAllocatedWorkers) {
lifecycleManager.getAllocatedWorkers()
} else {
- blacklist.asScala.keys.toSet ++ shuttingWorkers.asScala.toSet
+ excludedWorkers.asScala.keys.toSet ++ shuttingWorkers.asScala.toSet
}
}
def workerAvailable(worker: WorkerInfo): Boolean = {
- !blacklist.containsKey(worker) && !shuttingWorkers.contains(worker)
+ !excludedWorkers.containsKey(worker) && !shuttingWorkers.contains(worker)
}
def workerAvailable(loc: PartitionLocation): Boolean = {
@@ -65,13 +64,13 @@ class WorkerStatusTracker(
}
}
- def blacklistWorkerFromPartition(
+ def excludeWorkerFromPartition(
shuffleId: Int,
oldPartition: PartitionLocation,
cause: StatusCode): Unit = {
val failedWorker = new ShuffleFailedWorkers()
- def blacklistWorker(partition: PartitionLocation, statusCode: StatusCode):
Unit = {
+ def excludeWorker(partition: PartitionLocation, statusCode: StatusCode):
Unit = {
val tmpWorker = partition.getWorker
val worker =
lifecycleManager.workerSnapshots(shuffleId).keySet().asScala.find(_.equals(tmpWorker))
@@ -83,29 +82,29 @@ class WorkerStatusTracker(
if (oldPartition != null) {
cause match {
case StatusCode.PUSH_DATA_WRITE_FAIL_MASTER =>
- blacklistWorker(oldPartition, StatusCode.PUSH_DATA_WRITE_FAIL_MASTER)
+ excludeWorker(oldPartition, StatusCode.PUSH_DATA_WRITE_FAIL_MASTER)
case StatusCode.PUSH_DATA_WRITE_FAIL_SLAVE
if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled
=>
- blacklistWorker(oldPartition.getPeer,
StatusCode.PUSH_DATA_WRITE_FAIL_SLAVE)
+ excludeWorker(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)
+ excludeWorker(oldPartition,
StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_MASTER)
case StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_SLAVE
if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled
=>
- blacklistWorker(
+ excludeWorker(
oldPartition.getPeer,
StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_SLAVE)
case StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_MASTER =>
- blacklistWorker(oldPartition,
StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_MASTER)
+ excludeWorker(oldPartition,
StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_MASTER)
case StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_SLAVE
if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled
=>
- blacklistWorker(
+ excludeWorker(
oldPartition.getPeer,
StatusCode.PUSH_DATA_CONNECTION_EXCEPTION_SLAVE)
case StatusCode.PUSH_DATA_TIMEOUT_MASTER =>
- blacklistWorker(oldPartition, StatusCode.PUSH_DATA_TIMEOUT_MASTER)
+ excludeWorker(oldPartition, StatusCode.PUSH_DATA_TIMEOUT_MASTER)
case StatusCode.PUSH_DATA_TIMEOUT_SLAVE
if oldPartition.hasPeer && conf.clientExcludeSlaveOnFailureEnabled
=>
- blacklistWorker(
+ excludeWorker(
oldPartition.getPeer,
StatusCode.PUSH_DATA_TIMEOUT_SLAVE)
case _ =>
@@ -120,47 +119,47 @@ class WorkerStatusTracker(
val failedWorkerMsg = failedWorker.asScala.map { case (worker, (status,
time)) =>
s"${worker.readableAddress()} ${status.name()} $time"
}.mkString("\n")
- val blacklistMsg = blacklist.asScala.map { case (worker, (status, time))
=>
+ val excludedWorkerMsg = excludedWorkers.asScala.map { case (worker,
(status, time)) =>
s"${worker.readableAddress()} ${status.name()} $time"
}.mkString("\n")
val shuttingDownMsg =
shuttingWorkers.asScala.map(_.readableAddress()).mkString("\n")
logInfo(
s"""
- |Reporting Worker Failure:
+ |Reporting failed worker:
|$failedWorkerMsg
- |Current blacklist:
- |$blacklistMsg
- |Current shutting down:
+ |Current excluded worker:
+ |$excludedWorkerMsg
+ |Current shutting down worker:
|$shuttingDownMsg""".stripMargin)
failedWorker.asScala.foreach {
case (worker, (StatusCode.WORKER_SHUTDOWN, _)) =>
shuttingWorkers.add(worker)
- case (worker, (statusCode, registerTime)) if
!blacklist.containsKey(worker) =>
- blacklist.put(worker, (statusCode, registerTime))
+ case (worker, (statusCode, registerTime)) if
!excludedWorkers.containsKey(worker) =>
+ excludedWorkers.put(worker, (statusCode, registerTime))
case (worker, (statusCode, _))
if statusCode == StatusCode.NO_AVAILABLE_WORKING_DIR ||
statusCode == StatusCode.RESERVE_SLOTS_FAILED ||
- statusCode == StatusCode.UNKNOWN_WORKER =>
- blacklist.put(worker, (statusCode, blacklist.get(worker)._2))
+ statusCode == StatusCode.WORKER_UNKNOWN =>
+ excludedWorkers.put(worker, (statusCode,
excludedWorkers.get(worker)._2))
case _ => // Not cover
}
}
}
- def removeFromBlacklist(workers: JHashSet[WorkerInfo]): Unit = {
- blacklist.keySet.removeAll(workers)
+ def removeFromExcludedWorkers(workers: JHashSet[WorkerInfo]): Unit = {
+ excludedWorkers.keySet.removeAll(workers)
}
def handleHeartbeatResponse(res: HeartbeatFromApplicationResponse): Unit = {
if (res.statusCode == StatusCode.SUCCESS) {
- logInfo(s"Received Blacklist from Master, blacklist: ${res.blacklist} " +
+ logInfo(s"Received Worker status from Master, excluded workers:
${res.excludedWorkers} " +
s"unknown workers: ${res.unknownWorkers}, shutdown workers:
${res.shuttingWorkers}")
val current = System.currentTimeMillis()
- blacklist.asScala.foreach {
+ excludedWorkers.asScala.foreach {
case (workerInfo: WorkerInfo, (statusCode, registerTime)) =>
statusCode match {
- case StatusCode.UNKNOWN_WORKER |
+ case StatusCode.WORKER_UNKNOWN |
StatusCode.NO_AVAILABLE_WORKING_DIR |
StatusCode.RESERVE_SLOTS_FAILED |
StatusCode.PUSH_DATA_CREATE_CONNECTION_FAIL_MASTER |
@@ -171,24 +170,24 @@ class WorkerStatusTracker(
StatusCode.PUSH_DATA_TIMEOUT_SLAVE
if current - registerTime < excludedWorkerExpireTimeout => //
reserve
case _ =>
- if (!res.blacklist.contains(workerInfo) &&
+ if (!res.excludedWorkers.contains(workerInfo) &&
!res.shuttingWorkers.contains(workerInfo) &&
!res.unknownWorkers.contains(workerInfo)) {
- blacklist.remove(workerInfo)
+ excludedWorkers.remove(workerInfo)
}
}
}
- if (!res.blacklist.isEmpty) {
- blacklist.putAll(res.blacklist.asScala.filterNot(blacklist.containsKey)
- .map(_ -> (StatusCode.WORKER_IN_BLACKLIST -> current)).toMap.asJava)
+ if (!res.excludedWorkers.isEmpty) {
+
excludedWorkers.putAll(res.excludedWorkers.asScala.filterNot(excludedWorkers.containsKey)
+ .map(_ -> (StatusCode.WORKER_EXCLUDED -> current)).toMap.asJava)
}
shuttingWorkers.retainAll(res.shuttingWorkers)
shuttingWorkers.addAll(res.shuttingWorkers)
if (!res.unknownWorkers.isEmpty || !res.shuttingWorkers.isEmpty) {
-
blacklist.putAll(res.unknownWorkers.asScala.filterNot(blacklist.containsKey)
- .map(_ -> (StatusCode.UNKNOWN_WORKER -> current)).toMap.asJava)
+
excludedWorkers.putAll(res.unknownWorkers.asScala.filterNot(excludedWorkers.containsKey)
+ .map(_ -> (StatusCode.WORKER_UNKNOWN -> current)).toMap.asJava)
val workerStatus = new WorkersStatus(res.unknownWorkers,
res.shuttingWorkers)
workerStatusListeners.asScala.foreach { listener =>
try {
@@ -200,8 +199,9 @@ class WorkerStatusTracker(
}
}
- logInfo(s"Current blacklist $blacklist, Current shuttingDown
${shuttingWorkers.asScala.map(
- _.readableAddress()).mkString("\n")}")
+ logInfo(
+ s"Current excluded workers $excludedWorkers, Current shuttingDown
${shuttingWorkers.asScala.map(
+ _.readableAddress()).mkString("\n")}")
}
}
}
diff --git
a/client/src/main/scala/org/apache/celeborn/client/commit/CommitHandler.scala
b/client/src/main/scala/org/apache/celeborn/client/commit/CommitHandler.scala
index d11c8eff0..28a8ac7f9 100644
---
a/client/src/main/scala/org/apache/celeborn/client/commit/CommitHandler.scala
+++
b/client/src/main/scala/org/apache/celeborn/client/commit/CommitHandler.scala
@@ -268,9 +268,9 @@ abstract class CommitHandler(
commitEpoch.incrementAndGet())
val res =
if (conf.clientCommitFilesIgnoreExcludedWorkers &&
- workerStatusTracker.blacklist.containsKey(worker)) {
+ workerStatusTracker.excludedWorkers.containsKey(worker)) {
CommitFilesResponse(
- StatusCode.WORKER_IN_BLACKLIST,
+ StatusCode.WORKER_EXCLUDED,
List.empty.asJava,
List.empty.asJava,
masterIds,
@@ -281,10 +281,10 @@ abstract class CommitHandler(
res.status match {
case StatusCode.SUCCESS => // do nothing
- case StatusCode.PARTIAL_SUCCESS | StatusCode.SHUFFLE_NOT_REGISTERED
| StatusCode.REQUEST_FAILED | StatusCode.WORKER_IN_BLACKLIST =>
+ case StatusCode.PARTIAL_SUCCESS | StatusCode.SHUFFLE_NOT_REGISTERED
| StatusCode.REQUEST_FAILED | StatusCode.WORKER_EXCLUDED =>
logInfo(s"Request $commitFiles return ${res.status} for " +
s"${Utils.makeShuffleKey(applicationId, shuffleId)}")
- if (res.status != StatusCode.WORKER_IN_BLACKLIST) {
+ if (res.status != StatusCode.WORKER_EXCLUDED) {
commitFilesFailedWorkers.put(worker, (res.status,
System.currentTimeMillis()))
}
case _ =>
diff --git
a/client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala
b/client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala
index 2be8bbacb..d060d7108 100644
---
a/client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala
+++
b/client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala
@@ -36,8 +36,8 @@ class WorkerStatusTrackerSuite extends CelebornFunSuite {
val statusTracker = new WorkerStatusTracker(celebornConf, null)
val registerTime = System.currentTimeMillis()
- statusTracker.blacklist.put(mock("host1"), (StatusCode.UNKNOWN_WORKER,
registerTime));
- statusTracker.blacklist.put(mock("host2"), (StatusCode.WORKER_SHUTDOWN,
registerTime));
+ statusTracker.excludedWorkers.put(mock("host1"),
(StatusCode.WORKER_UNKNOWN, registerTime));
+ statusTracker.excludedWorkers.put(mock("host2"),
(StatusCode.WORKER_SHUTDOWN, registerTime));
// test reserve (only statusCode list in handleHeartbeatResponse)
val empty = buildResponse(Array.empty, Array.empty, Array.empty)
@@ -45,53 +45,57 @@ class WorkerStatusTrackerSuite extends CelebornFunSuite {
// only reserve host1
Assert.assertEquals(
- statusTracker.blacklist.get(mock("host1")),
- (StatusCode.UNKNOWN_WORKER, registerTime))
- Assert.assertFalse(statusTracker.blacklist.containsKey(mock("host2")))
+ statusTracker.excludedWorkers.get(mock("host1")),
+ (StatusCode.WORKER_UNKNOWN, registerTime))
+
Assert.assertFalse(statusTracker.excludedWorkers.containsKey(mock("host2")))
- // add shutdown/blacklist
+ // add shutdown/excluded worker
val response1 = buildResponse(Array("host0"), Array("host1", "host3"),
Array("host4"))
statusTracker.handleHeartbeatResponse(response1)
// test keep Unknown register time
Assert.assertEquals(
- statusTracker.blacklist.get(mock("host1")),
- (StatusCode.UNKNOWN_WORKER, registerTime))
+ statusTracker.excludedWorkers.get(mock("host1")),
+ (StatusCode.WORKER_UNKNOWN, registerTime))
// test new added workers
- Assert.assertTrue(statusTracker.blacklist.containsKey(mock("host0")))
- Assert.assertTrue(statusTracker.blacklist.containsKey(mock("host3")))
- Assert.assertTrue(!statusTracker.blacklist.contains(mock("host4")))
+ Assert.assertTrue(statusTracker.excludedWorkers.containsKey(mock("host0")))
+ Assert.assertTrue(statusTracker.excludedWorkers.containsKey(mock("host3")))
+ Assert.assertTrue(!statusTracker.excludedWorkers.contains(mock("host4")))
Assert.assertTrue(statusTracker.shuttingWorkers.contains(mock("host4")))
// test re heartbeat with shutdown workers
val response3 = buildResponse(Array.empty, Array.empty, Array("host4"))
statusTracker.handleHeartbeatResponse(response3)
- Assert.assertTrue(!statusTracker.blacklist.contains(mock("host4")))
+ Assert.assertTrue(!statusTracker.excludedWorkers.contains(mock("host4")))
Assert.assertTrue(statusTracker.shuttingWorkers.contains(mock("host4")))
// test remove
val workers = new util.HashSet[WorkerInfo]
workers.add(mock("host3"))
- statusTracker.removeFromBlacklist(workers)
- Assert.assertFalse(statusTracker.blacklist.containsKey(mock("host3")))
+ statusTracker.removeFromExcludedWorkers(workers)
+
Assert.assertFalse(statusTracker.excludedWorkers.containsKey(mock("host3")))
// test register time elapsed
Thread.sleep(3000)
val response2 = buildResponse(Array.empty, Array("host5", "host6"),
Array.empty)
statusTracker.handleHeartbeatResponse(response2)
- Assert.assertEquals(statusTracker.blacklist.size(), 2)
- Assert.assertFalse(statusTracker.blacklist.containsKey(mock("host1")))
+ Assert.assertEquals(statusTracker.excludedWorkers.size(), 2)
+
Assert.assertFalse(statusTracker.excludedWorkers.containsKey(mock("host1")))
}
private def buildResponse(
- blackWorkerHosts: Array[String],
+ excludedWorkerHosts: Array[String],
unknownWorkerHosts: Array[String],
shuttingWorkerHosts: Array[String]): HeartbeatFromApplicationResponse = {
- val blacklist = mockWorkers(blackWorkerHosts)
+ val excludedWorkers = mockWorkers(excludedWorkerHosts)
val unknownWorkers = mockWorkers(unknownWorkerHosts)
val shuttingWorkers = mockWorkers(shuttingWorkerHosts)
- HeartbeatFromApplicationResponse(StatusCode.SUCCESS, blacklist,
unknownWorkers, shuttingWorkers)
+ HeartbeatFromApplicationResponse(
+ StatusCode.SUCCESS,
+ excludedWorkers,
+ unknownWorkers,
+ shuttingWorkers)
}
private def mockWorkers(workerHosts: Array[String]):
util.ArrayList[WorkerInfo] = {
diff --git
a/common/src/main/java/org/apache/celeborn/common/protocol/message/StatusCode.java
b/common/src/main/java/org/apache/celeborn/common/protocol/message/StatusCode.java
index 5ad47bed8..14ff78074 100644
---
a/common/src/main/java/org/apache/celeborn/common/protocol/message/StatusCode.java
+++
b/common/src/main/java/org/apache/celeborn/common/protocol/message/StatusCode.java
@@ -52,8 +52,8 @@ public enum StatusCode {
SHUFFLE_DATA_LOST(24),
WORKER_SHUTDOWN(25),
NO_AVAILABLE_WORKING_DIR(26),
- WORKER_IN_BLACKLIST(27),
- UNKNOWN_WORKER(28),
+ WORKER_EXCLUDED(27),
+ WORKER_UNKNOWN(28),
COMMIT_FILE_EXCEPTION(29),
@@ -74,8 +74,8 @@ public enum StatusCode {
PUSH_DATA_CONNECTION_EXCEPTION_SLAVE(41),
PUSH_DATA_TIMEOUT_MASTER(42),
PUSH_DATA_TIMEOUT_SLAVE(43),
- PUSH_DATA_MASTER_BLACKLISTED(44),
- PUSH_DATA_SLAVE_BLACKLISTED(45),
+ PUSH_DATA_MASTER_WORKER_EXCLUDED(44),
+ PUSH_DATA_SLAVE_WORKER_EXCLUDED(45),
FETCH_DATA_TIMEOUT(46),
REVIVE_INITIALIZED(47);
diff --git
a/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
b/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
index 9354bee62..8afbe07b8 100644
---
a/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
+++
b/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
@@ -328,15 +328,15 @@ object ControlMessages extends Logging {
case class HeartbeatFromApplicationResponse(
statusCode: StatusCode,
- blacklist: util.List[WorkerInfo],
+ excludedWorkers: util.List[WorkerInfo],
unknownWorkers: util.List[WorkerInfo],
shuttingWorkers: util.List[WorkerInfo]) extends Message
- case class GetBlacklist(localBlacklist: util.List[WorkerInfo]) extends
MasterMessage
+ case class GetBlacklist(localExcludedWorkers: util.List[WorkerInfo]) extends
MasterMessage
case class GetBlacklistResponse(
statusCode: StatusCode,
- blacklist: util.List[WorkerInfo],
+ excludedWorkers: util.List[WorkerInfo],
unknownWorkers: util.List[WorkerInfo]) extends Message
case class CheckQuota(userIdentifier: UserIdentifier) extends Message
@@ -645,20 +645,19 @@ object ControlMessages extends Logging {
.build().toByteArray
new TransportMessage(MessageType.HEARTBEAT_FROM_APPLICATION_RESPONSE,
payload)
- case GetBlacklist(localBlacklist) =>
+ case GetBlacklist(localExcludedWorkers) =>
val payload = PbGetBlacklist.newBuilder()
- .addAllLocalBlackList(localBlacklist.asScala.map { workerInfo =>
+ .addAllLocalBlackList(localExcludedWorkers.asScala.map { workerInfo =>
PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
- }
- .toList.asJava)
+ }.toList.asJava)
.build().toByteArray
new TransportMessage(MessageType.GET_BLACKLIST, payload)
- case GetBlacklistResponse(statusCode, blacklist, unknownWorkers) =>
+ case GetBlacklistResponse(statusCode, excludedWorkers, unknownWorkers) =>
val builder = PbGetBlacklistResponse.newBuilder()
.setStatus(statusCode.getValue)
builder.addAllBlacklist(
- blacklist.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
+ excludedWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
builder.addAllUnknownWorkers(
unknownWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
diff --git
a/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
b/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
index 9841a56ec..232cc1c30 100644
--- a/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
@@ -364,7 +364,7 @@ object PbSerDeUtils {
estimatedPartitionSize: java.lang.Long,
registeredShuffle: java.util.Set[String],
hostnameSet: java.util.Set[String],
- blacklist: java.util.Set[WorkerInfo],
+ excludedWorkers: java.util.Set[WorkerInfo],
workerLostEvent: java.util.Set[WorkerInfo],
appHeartbeatTime: java.util.Map[String, java.lang.Long],
workers: java.util.List[WorkerInfo],
@@ -378,7 +378,7 @@ object PbSerDeUtils {
.setEstimatedPartitionSize(estimatedPartitionSize)
.addAllRegisteredShuffle(registeredShuffle)
.addAllHostnameSet(hostnameSet)
- .addAllBlacklist(blacklist.asScala.map(toPbWorkerInfo(_, true)).asJava)
+ .addAllBlacklist(excludedWorkers.asScala.map(toPbWorkerInfo(_,
true)).asJava)
.addAllWorkerLostEvents(workerLostEvent.asScala.map(toPbWorkerInfo(_,
true)).asJava)
.putAllAppHeartbeatTime(appHeartbeatTime)
.addAllWorkers(workers.asScala.map(toPbWorkerInfo(_, true)).asJava)
diff --git a/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
b/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
index 7b1163e54..221dba251 100644
--- a/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
@@ -858,9 +858,9 @@ object Utils extends Logging {
case 26 =>
StatusCode.NO_AVAILABLE_WORKING_DIR
case 27 =>
- StatusCode.WORKER_IN_BLACKLIST
+ StatusCode.WORKER_EXCLUDED
case 28 =>
- StatusCode.UNKNOWN_WORKER
+ StatusCode.WORKER_UNKNOWN
case 30 =>
StatusCode.PUSH_DATA_SUCCESS_MASTER_CONGESTED
case 31 =>
@@ -878,9 +878,9 @@ object Utils extends Logging {
case 43 =>
StatusCode.PUSH_DATA_TIMEOUT_SLAVE
case 44 =>
- StatusCode.PUSH_DATA_MASTER_BLACKLISTED
+ StatusCode.PUSH_DATA_MASTER_WORKER_EXCLUDED
case 45 =>
- StatusCode.PUSH_DATA_SLAVE_BLACKLISTED
+ StatusCode.PUSH_DATA_SLAVE_WORKER_EXCLUDED
case 46 =>
StatusCode.FETCH_DATA_TIMEOUT
case 47 =>
diff --git a/docs/migration.md b/docs/migration.md
index 727817a32..af76683ee 100644
--- a/docs/migration.md
+++ b/docs/migration.md
@@ -43,3 +43,7 @@ license: |
- Since 0.3.0, Celeborn supports overriding Hadoop
configuration(`core-site.xml`, `hdfs-site.xml`, etc.) from Celeborn
configuration with the additional prefix `celeborn.hadoop.`.
On Spark client side, user should set Hadoop configuration like
`spark.celeborn.hadoop.foo=bar`, note that `spark.hadoop.foo=bar` does not take
effect;
on Flink client and Celeborn Master/Worker side, user should set like
`celeborn.hadoop.foo=bar`.
+
+ - Since 0.3.0, Celeborn master metrics `BlacklistedWorkerCount` is renamed as
`ExcludedWorkerCount`.
+
+ - Since 0.3.0, Celeborn master http request url `/blacklistedWorkers` is
renamed as `/excludedWorkers`.
diff --git a/docs/monitoring.md b/docs/monitoring.md
index d8b3c34ef..efe93b4d5 100644
--- a/docs/monitoring.md
+++ b/docs/monitoring.md
@@ -92,7 +92,7 @@ These metrics are exposed by Celeborn master.
- namespace=master
- WorkerCount
- LostWorkers
- - BlacklistedWorkerCount
+ - ExcludedWorkerCount
- RegisteredShuffleCount
- IsActiveMaster
- PartitionSize
@@ -282,19 +282,19 @@ The configuration of `<master-prometheus-host>`,
`<master-prometheus-port>`, `<w
API path listed as below:
-| Path | Service | Meaning
|
-|----------------------------|-----------------|--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
-| /metrics/prometheus | master, worker | List service metrics data in
prometheus format.
|
-| /conf | master, worker | List the conf setting of the
service.
|
-| /workerInfo | master, worker | List worker information of
the service. For the master, it will list all registered workers 's
information.
|
-| /lostWorkers | master | List all lost workers of the
master.
|
-| /blacklistedWorkers | master | List all blacklisted workers
of the master.
|
-| /threadDump | master, worker | List the current thread dump
of the service.
|
-| /hostnames | master | List all running
application's LifecycleManager's hostnames of the cluster.
|
-| /applications | master | List all running
application's ids of the cluster.
|
-| /shuffles | master, worker | List all running shuffle keys
of the service. For master, will return all running shuffle's key of the
cluster, for worker, only return keys of shuffles running in that worker. |
-| /listTopDiskUsedApps | master, worker | List the top disk usage
application ids. For master, will return the top disk usage application ids for
the cluster, for worker, only return application ids running in that worker. |
-| /listPartitionLocationInfo | worker | List all living
PartitionLocation information in that worker.
|
-| /unavailablePeers | worker | List the unavailable peers of
the worker, this always means the worker connect to the peer failed.
|
-| /isShutdown | worker | Show if the worker is during
the process of shutdown.
|
-| /isRegistered | worker | Show if the worker is
registered to the master success.
|
\ No newline at end of file
+| Path | Service | Meaning
|
+|----------------------------|-----------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
+| /metrics/prometheus | master, worker | List service metrics data in
prometheus format.
|
+| /conf | master, worker | List the conf setting of the
service.
|
+| /workerInfo | master, worker | List worker information of
the service. For the master, it will list all registered workers 's
information.
|
+| /lostWorkers | master | List all lost workers of the
master.
|
+| /excludedWorkers | master | List all excluded workers of
the master.
|
+| /threadDump | master, worker | List the current thread dump
of the service.
|
+| /hostnames | master | List all running
application's LifecycleManager's hostnames of the cluster.
|
+| /applications | master | List all running
application's ids of the cluster.
|
+| /shuffles | master, worker | List all running shuffle keys
of the service. For master, will return all running shuffle's key of the
cluster, for worker, only return keys of shuffles running in that worker. |
+| /listTopDiskUsedApps | master, worker | List the top disk usage
application ids. For master, will return the top disk usage application ids for
the cluster, for worker, only return application ids running in that worker. |
+| /listPartitionLocationInfo | worker | List all living
PartitionLocation information in that worker.
|
+| /unavailablePeers | worker | List the unavailable peers of
the worker, this always means the worker connect to the peer failed.
|
+| /isShutdown | worker | Show if the worker is during
the process of shutdown.
|
+| /isRegistered | worker | Show if the worker is
registered to the master success.
|
\ No newline at end of file
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
index d51079c99..2ab16409a 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
@@ -61,12 +61,8 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
public final ArrayList<WorkerInfo> workers = new ArrayList<>();
public final ConcurrentHashMap<WorkerInfo, Long> lostWorkers =
JavaUtils.newConcurrentHashMap();
public final ConcurrentHashMap<String, Long> appHeartbeatTime =
JavaUtils.newConcurrentHashMap();
- // blacklist
- public final Set<WorkerInfo> blacklist = ConcurrentHashMap.newKeySet();
- // shutdown workers
+ public final Set<WorkerInfo> excludedWorkers = ConcurrentHashMap.newKeySet();
public final Set<WorkerInfo> shutdownWorkers = ConcurrentHashMap.newKeySet();
-
- // workerLost events
public final Set<WorkerInfo> workerLostEvents =
ConcurrentHashMap.newKeySet();
protected RpcEnv rpcEnv;
@@ -159,8 +155,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
workers.remove(worker);
lostWorkers.put(worker, System.currentTimeMillis());
}
- // delete from blacklist
- blacklist.remove(worker);
+ excludedWorkers.remove(worker);
workerLostEvents.remove(worker);
}
@@ -172,8 +167,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
workers.remove(worker);
lostWorkers.put(worker, System.currentTimeMillis());
}
- // delete from blacklist
- blacklist.remove(worker);
+ excludedWorkers.remove(worker);
}
public void updateWorkerHeartbeatMeta(
@@ -202,13 +196,13 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
});
}
appDiskUsageMetric.update(estimatedAppDiskUsage);
- // If using HDFSONLY mode, workers with empty disks should not be put into
blacklist.
- if (!blacklist.contains(worker) && (disks.isEmpty() &&
!conf.hasHDFSStorage())) {
- LOG.debug("Worker: {} num total slots is 0, add to blacklist", worker);
- blacklist.add(worker);
+ // If using HDFSONLY mode, workers with empty disks should not be put into
excluded worker list.
+ if (!excludedWorkers.contains(worker) && (disks.isEmpty() &&
!conf.hasHDFSStorage())) {
+ LOG.debug("Worker: {} num total slots is 0, add to excluded list",
worker);
+ excludedWorkers.add(worker);
} else if (availableSlots.get() > 0) {
// only unblack if numSlots larger than 0
- blacklist.remove(worker);
+ excludedWorkers.remove(worker);
}
}
@@ -247,7 +241,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
estimatedPartitionSize,
registeredShuffle,
hostnameSet,
- blacklist,
+ excludedWorkers,
workerLostEvents,
appHeartbeatTime,
workers,
@@ -275,7 +269,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
registeredShuffle.addAll(snapshotMetaInfo.getRegisteredShuffleList());
hostnameSet.addAll(snapshotMetaInfo.getHostnameSetList());
- blacklist.addAll(
+ excludedWorkers.addAll(
snapshotMetaInfo.getBlacklistList().stream()
.map(PbSerDeUtils::fromPbWorkerInfo)
.collect(Collectors.toSet()));
@@ -335,10 +329,10 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
}
LOG.info("Successfully restore meta info from snapshot {}",
file.getAbsolutePath());
LOG.info(
- "Worker size: {}, Registered shuffle size: {}, Worker blacklist size:
{}.",
+ "Worker size: {}, Registered shuffle size: {}. Worker excluded list
size: {}.",
workers.size(),
registeredShuffle.size(),
- blacklist.size());
+ excludedWorkers.size());
workers.forEach(workerInfo -> LOG.info(workerInfo.toString()));
registeredShuffle.forEach(shuffle -> LOG.info("RegisteredShuffle {}",
shuffle));
}
@@ -367,7 +361,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
Utils.bytesToString(oldEstimatedPartitionSize),
Utils.bytesToString(estimatedPartitionSize));
workers.stream()
- .filter(worker -> !blacklist.contains(worker))
+ .filter(worker -> !excludedWorkers.contains(worker))
.forEach(workerInfo ->
workerInfo.updateDiskMaxSlots(estimatedPartitionSize));
}
}
diff --git
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
index e001a6f1c..ee1dc1488 100644
---
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
+++
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
@@ -144,13 +144,10 @@ private[celeborn] class Master(
masterSource.addGauge(
MasterSource.RegisteredShuffleCount,
_ => statusSystem.registeredShuffle.size())
- // blacklist worker count
- masterSource.addGauge(MasterSource.BlacklistedWorkerCount, _ =>
statusSystem.blacklist.size())
- // worker count
+ masterSource.addGauge(MasterSource.ExcludedWorkerCount, _ =>
statusSystem.excludedWorkers.size())
masterSource.addGauge(MasterSource.WorkerCount, _ =>
statusSystem.workers.size())
masterSource.addGauge(MasterSource.LostWorkerCount, _ =>
statusSystem.lostWorkers.size())
masterSource.addGauge(MasterSource.PartitionSize, _ =>
statusSystem.estimatedPartitionSize)
- // is master active under HA mode
masterSource.addGauge(MasterSource.IsActiveMaster, _ => isMasterActive)
metricsSystem.registerSource(resourceConsumptionSource)
@@ -234,7 +231,7 @@ private[celeborn] class Master(
appId,
totalWritten,
fileCount,
- localBlacklist,
+ needCheckedWorkerList,
requestId,
shouldResponse) =>
logDebug(s"Received heartbeat from app $appId")
@@ -245,7 +242,7 @@ private[celeborn] class Master(
appId,
totalWritten,
fileCount,
- localBlacklist,
+ needCheckedWorkerList,
requestId,
shouldResponse))
@@ -634,12 +631,12 @@ private[celeborn] class Master(
}
def handleGetBlacklist(context: RpcCallContext, msg: GetBlacklist): Unit = {
- msg.localBlacklist.removeAll(workersSnapShot)
+ msg.localExcludedWorkers.removeAll(workersSnapShot)
context.reply(
GetBlacklistResponse(
StatusCode.SUCCESS,
- new util.ArrayList(statusSystem.blacklist),
- msg.localBlacklist))
+ new util.ArrayList(statusSystem.excludedWorkers),
+ msg.localExcludedWorkers))
}
private def handleGetWorkerInfos(context: RpcCallContext): Unit = {
@@ -650,8 +647,8 @@ private[celeborn] class Master(
context: RpcCallContext,
failedWorkers: util.List[WorkerInfo],
requestId: String): Unit = {
- logInfo(s"Receive ReportNodeFailure $failedWorkers, current blacklist" +
- s"${statusSystem.blacklist}")
+ logInfo(s"Receive ReportNodeFailure $failedWorkers, current excluded
workers" +
+ s"${statusSystem.excludedWorkers}")
statusSystem.handleReportWorkerUnavailable(failedWorkers, requestId)
context.reply(OneWayMessageResponse)
}
@@ -685,7 +682,7 @@ private[celeborn] class Master(
if (shouldResponse) {
context.reply(HeartbeatFromApplicationResponse(
StatusCode.SUCCESS,
- new util.ArrayList(statusSystem.blacklist),
+ new util.ArrayList(statusSystem.excludedWorkers),
needCheckedWorkerList,
shutdownWorkerSnapshot))
} else {
@@ -745,10 +742,10 @@ private[celeborn] class Master(
}
private def workersAvailable(
- tmpBlacklist: Set[WorkerInfo] = Set.empty): util.List[WorkerInfo] = {
+ tmpExcludedWorkerList: Set[WorkerInfo] = Set.empty):
util.List[WorkerInfo] = {
workersSnapShot.asScala.filter { w =>
- !statusSystem.blacklist.contains(w) &&
!statusSystem.shutdownWorkers.contains(
- w) && !tmpBlacklist.contains(w)
+ !statusSystem.excludedWorkers.contains(w) &&
!statusSystem.shutdownWorkers.contains(
+ w) && !tmpExcludedWorkerList.contains(w)
}.asJava
}
@@ -779,10 +776,10 @@ private[celeborn] class Master(
sb.toString()
}
- override def getBlacklistedWorkers: String = {
+ override def getExcludedWorkers: String = {
val sb = new StringBuilder
- sb.append("==================== Blacklisted Workers in Master
=====================\n")
- statusSystem.blacklist.asScala.foreach { worker =>
+ sb.append("===================== Excluded Workers in Master
======================\n")
+ statusSystem.excludedWorkers.asScala.foreach { worker =>
sb.append(s"${worker.toUniqueId()}\n")
}
sb.toString()
diff --git
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/MasterSource.scala
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/MasterSource.scala
index 86d64edd9..5c3862cce 100644
---
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/MasterSource.scala
+++
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/MasterSource.scala
@@ -37,7 +37,7 @@ object MasterSource {
val LostWorkerCount = "LostWorkers"
- val BlacklistedWorkerCount = "BlacklistedWorkerCount"
+ val ExcludedWorkerCount = "ExcludedWorkerCount"
val RegisteredShuffleCount = "RegisteredShuffleCount"
diff --git
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
index 115426c8b..1b6df8205 100644
---
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
+++
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
@@ -511,7 +511,7 @@ public class DefaultMetaSystemSuiteJ {
1,
getNewReqeustId());
- Assert.assertEquals(statusSystem.blacklist.size(), 1);
+ Assert.assertEquals(statusSystem.excludedWorkers.size(), 1);
statusSystem.handleWorkerHeartbeat(
HOSTNAME2,
@@ -525,7 +525,7 @@ public class DefaultMetaSystemSuiteJ {
1,
getNewReqeustId());
- Assert.assertEquals(statusSystem.blacklist.size(), 2);
+ Assert.assertEquals(statusSystem.excludedWorkers.size(), 2);
statusSystem.handleWorkerHeartbeat(
HOSTNAME1,
@@ -539,7 +539,7 @@ public class DefaultMetaSystemSuiteJ {
1,
getNewReqeustId());
- Assert.assertEquals(statusSystem.blacklist.size(), 2);
+ Assert.assertEquals(statusSystem.excludedWorkers.size(), 2);
}
@Test
@@ -596,7 +596,7 @@ public class DefaultMetaSystemSuiteJ {
statusSystem.handleReportWorkerUnavailable(failedWorkers,
getNewReqeustId());
assert 1 == statusSystem.shutdownWorkers.size();
- assert 0 == statusSystem.blacklist.size();
+ assert 0 == statusSystem.excludedWorkers.size();
}
@Test
diff --git
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java
index c092fbf7a..bcf067826 100644
---
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java
+++
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java
@@ -254,9 +254,9 @@ public class MasterStateMachineSuiteJ extends
RatisBaseSuiteJ {
String host2 = "host2";
String host3 = "host3";
- masterStatusSystem.blacklist.add(info1);
- masterStatusSystem.blacklist.add(info2);
- masterStatusSystem.blacklist.add(info3);
+ masterStatusSystem.excludedWorkers.add(info1);
+ masterStatusSystem.excludedWorkers.add(info2);
+ masterStatusSystem.excludedWorkers.add(info3);
masterStatusSystem.hostnameSet.add(host1);
masterStatusSystem.hostnameSet.add(host2);
@@ -281,11 +281,11 @@ public class MasterStateMachineSuiteJ extends
RatisBaseSuiteJ {
masterStatusSystem.writeMetaInfoToFile(tmpFile);
masterStatusSystem.hostnameSet.clear();
- masterStatusSystem.blacklist.clear();
+ masterStatusSystem.excludedWorkers.clear();
masterStatusSystem.restoreMetaFromFile(tmpFile);
- Assert.assertEquals(3, masterStatusSystem.blacklist.size());
+ Assert.assertEquals(3, masterStatusSystem.excludedWorkers.size());
Assert.assertEquals(3, masterStatusSystem.hostnameSet.size());
Assert.assertEquals(
conf.metricsAppTopDiskUsageWindowSize(),
diff --git
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
index 9b9bea917..f2ab80996 100644
---
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
+++
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
@@ -736,9 +736,9 @@ public class RatisMasterStatusSystemSuiteJ {
getNewReqeustId());
Thread.sleep(3000L);
- Assert.assertEquals(1, STATUSSYSTEM1.blacklist.size());
- Assert.assertEquals(1, STATUSSYSTEM2.blacklist.size());
- Assert.assertEquals(1, STATUSSYSTEM3.blacklist.size());
+ Assert.assertEquals(1, STATUSSYSTEM1.excludedWorkers.size());
+ Assert.assertEquals(1, STATUSSYSTEM2.excludedWorkers.size());
+ Assert.assertEquals(1, STATUSSYSTEM3.excludedWorkers.size());
statusSystem.handleWorkerHeartbeat(
HOSTNAME2,
@@ -753,10 +753,10 @@ public class RatisMasterStatusSystemSuiteJ {
getNewReqeustId());
Thread.sleep(3000L);
- Assert.assertEquals(2, statusSystem.blacklist.size());
- Assert.assertEquals(2, STATUSSYSTEM1.blacklist.size());
- Assert.assertEquals(2, STATUSSYSTEM2.blacklist.size());
- Assert.assertEquals(2, STATUSSYSTEM3.blacklist.size());
+ Assert.assertEquals(2, statusSystem.excludedWorkers.size());
+ Assert.assertEquals(2, STATUSSYSTEM1.excludedWorkers.size());
+ Assert.assertEquals(2, STATUSSYSTEM2.excludedWorkers.size());
+ Assert.assertEquals(2, STATUSSYSTEM3.excludedWorkers.size());
statusSystem.handleWorkerHeartbeat(
HOSTNAME1,
@@ -771,10 +771,10 @@ public class RatisMasterStatusSystemSuiteJ {
getNewReqeustId());
Thread.sleep(3000L);
- Assert.assertEquals(1, statusSystem.blacklist.size());
- Assert.assertEquals(1, STATUSSYSTEM1.blacklist.size());
- Assert.assertEquals(1, STATUSSYSTEM2.blacklist.size());
- Assert.assertEquals(1, STATUSSYSTEM3.blacklist.size());
+ Assert.assertEquals(1, statusSystem.excludedWorkers.size());
+ Assert.assertEquals(1, STATUSSYSTEM1.excludedWorkers.size());
+ Assert.assertEquals(1, STATUSSYSTEM2.excludedWorkers.size());
+ Assert.assertEquals(1, STATUSSYSTEM3.excludedWorkers.size());
}
@Before
@@ -783,21 +783,21 @@ public class RatisMasterStatusSystemSuiteJ {
STATUSSYSTEM1.hostnameSet.clear();
STATUSSYSTEM1.workers.clear();
STATUSSYSTEM1.appHeartbeatTime.clear();
- STATUSSYSTEM1.blacklist.clear();
+ STATUSSYSTEM1.excludedWorkers.clear();
STATUSSYSTEM1.workerLostEvents.clear();
STATUSSYSTEM2.registeredShuffle.clear();
STATUSSYSTEM2.hostnameSet.clear();
STATUSSYSTEM2.workers.clear();
STATUSSYSTEM2.appHeartbeatTime.clear();
- STATUSSYSTEM2.blacklist.clear();
+ STATUSSYSTEM2.excludedWorkers.clear();
STATUSSYSTEM2.workerLostEvents.clear();
STATUSSYSTEM3.registeredShuffle.clear();
STATUSSYSTEM3.hostnameSet.clear();
STATUSSYSTEM3.workers.clear();
STATUSSYSTEM3.appHeartbeatTime.clear();
- STATUSSYSTEM3.blacklist.clear();
+ STATUSSYSTEM3.excludedWorkers.clear();
STATUSSYSTEM3.workerLostEvents.clear();
disks1.clear();
@@ -879,9 +879,9 @@ public class RatisMasterStatusSystemSuiteJ {
Assert.assertEquals(1, STATUSSYSTEM1.shutdownWorkers.size());
Assert.assertEquals(1, STATUSSYSTEM2.shutdownWorkers.size());
Assert.assertEquals(1, STATUSSYSTEM3.shutdownWorkers.size());
- Assert.assertEquals(0, STATUSSYSTEM1.blacklist.size());
- Assert.assertEquals(0, STATUSSYSTEM2.blacklist.size());
- Assert.assertEquals(0, STATUSSYSTEM3.blacklist.size());
+ Assert.assertEquals(0, STATUSSYSTEM1.excludedWorkers.size());
+ Assert.assertEquals(0, STATUSSYSTEM2.excludedWorkers.size());
+ Assert.assertEquals(0, STATUSSYSTEM3.excludedWorkers.size());
}
@Test
diff --git
a/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
b/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
index 47e2833cf..3cfb34f3c 100644
--- a/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
+++ b/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
@@ -47,7 +47,7 @@ abstract class HttpService extends Service with Logging {
def getShutdownWorkers: String = throw new UnsupportedOperationException()
- def getBlacklistedWorkers: String = throw new UnsupportedOperationException()
+ def getExcludedWorkers: String = throw new UnsupportedOperationException()
def getThreadDump: String
diff --git
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
index 53e4ba5f5..2747704dd 100644
---
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
+++
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpRequestHandler.scala
@@ -68,8 +68,8 @@ class HttpRequestHandler(
service.getWorkerInfo
case "/lostWorkers" if service.serviceName == Service.MASTER =>
service.getLostWorkers
- case "/blacklistedWorkers" if service.serviceName == Service.MASTER =>
- service.getBlacklistedWorkers
+ case "/excludedWorkers" if service.serviceName == Service.MASTER =>
+ service.getExcludedWorkers
case "/shutdownWorkers" if service.serviceName == Service.MASTER =>
service.getShutdownWorkers
case "/threadDump" =>
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 68e381fc6..2b31d7de1 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
@@ -41,7 +41,7 @@ class PushDataTimeoutTest extends AnyFunSuite
CelebornConf.TEST_CLIENT_PUSH_MASTER_DATA_TIMEOUT.key -> "true",
CelebornConf.TEST_WORKER_PUSH_SLAVE_DATA_TIMEOUT.key -> "true")
// required at least 4 workers, the reason behind this requirement is that
when replication is
- // enabled, there is a possibility that two workers might be added to the
blacklist due to
+ // enabled, there is a possibility that two workers might be added to the
excluded list due to
// master/slave timeout issues, then there are not enough workers to do
replication if available
// workers number = 1
setUpMiniCluster(masterConfs = null, workerConfs = workerConf, workerNum =
4)
@@ -120,7 +120,7 @@ class PushDataTimeoutTest extends AnyFunSuite
}
}
- test("celeborn spark integration test - pushdata timeout will add to
blacklist") {
+ test("celeborn spark integration test - pushdata timeout will add to
pushExcludedWorkers") {
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_EXCLUDE_SLAVE_ON_FAILURE_ENABLED.key}",
"true")
@@ -138,15 +138,15 @@ class PushDataTimeoutTest extends AnyFunSuite
assert(PushDataHandler.pushMasterMergeDataTimeoutTested.get())
assert(PushDataHandler.pushSlaveMergeDataTimeoutTested.get())
- val blacklist = SparkContextHelper.env
+ val excludedWorkers = SparkContextHelper.env
.shuffleManager
.asInstanceOf[RssShuffleManager]
.getLifecycleManager
.workerStatusTracker
- .blacklist
+ .excludedWorkers
- assert(blacklist.size() > 0)
- blacklist.asScala.foreach { case (_, (code, _)) =>
+ assert(excludedWorkers.size() > 0)
+ excludedWorkers.asScala.foreach { case (_, (code, _)) =>
assert(code == StatusCode.PUSH_DATA_TIMEOUT_MASTER ||
code == StatusCode.PUSH_DATA_TIMEOUT_SLAVE)
}
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 303c30342..f6775de28 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
@@ -528,8 +528,8 @@ private[celeborn] class Worker(
override def run(): Unit = {
logInfo("Shutdown hook called.")
// During shutdown, to avoid allocate slots in this worker,
- // add this worker to master's blacklist. When restart, register
worker will
- // make master remove this worker from blacklist.
+ // add this worker to master's excluded list. When restart, register
worker will
+ // make master remove this worker from excluded list.
try {
if (gracefulShutdown) {
rssHARetryClient.askSync(