This is an automated email from the ASF dual-hosted git repository.
chengpan 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 590198ece [CELEBORN-666][FOLLOWUP] Rename all RPC blacklist fields
590198ece is described below
commit 590198eceab12e7e07d2b541a8308ee17805077a
Author: Angerszhuuuu <[email protected]>
AuthorDate: Wed Jun 28 19:49:44 2023 +0800
[CELEBORN-666][FOLLOWUP] Rename all RPC blacklist fields
### What changes were proposed in this pull request?
In this pr, we rename all RPC blacklist fields, it won't have have
compatibility issues.
For RPC `GetBlacklist` and `GetBlacklistResponse` we won't change it, since
it won't be used in next release, so we can remove these two RPC in next
release.
### Why are the changes needed?
### Does this PR introduce _any_ user-facing change?
### How was this patch tested?
Closes #1643 from AngersZhuuuu/CELEBORN-666-RPC.
Authored-by: Angerszhuuuu <[email protected]>
Signed-off-by: Cheng Pan <[email protected]>
---
common/src/main/proto/TransportMessages.proto | 8 ++++----
.../common/protocol/message/ControlMessages.scala | 23 +++++++++++++---------
.../apache/celeborn/common/util/PbSerDeUtils.scala | 2 +-
.../master/clustermeta/AbstractMetaManager.java | 2 +-
4 files changed, 20 insertions(+), 15 deletions(-)
diff --git a/common/src/main/proto/TransportMessages.proto
b/common/src/main/proto/TransportMessages.proto
index d1430d72b..29c73595a 100644
--- a/common/src/main/proto/TransportMessages.proto
+++ b/common/src/main/proto/TransportMessages.proto
@@ -295,18 +295,18 @@ message PbHeartbeatFromApplication {
message PbHeartbeatFromApplicationResponse {
int32 status = 1;
- repeated PbWorkerInfo blacklist = 2;
+ repeated PbWorkerInfo excludedWorkers = 2;
repeated PbWorkerInfo unknownWorkers = 3;
repeated PbWorkerInfo shuttingWorkers = 4;
}
message PbGetBlacklist {
- repeated PbWorkerInfo localBlackList = 1;
+ repeated PbWorkerInfo localExcludedWorkers = 1;
}
message PbGetBlacklistResponse {
int32 status = 1;
- repeated PbWorkerInfo blacklist = 2;
+ repeated PbWorkerInfo excludedWorkers = 2;
repeated PbWorkerInfo unknownWorkers = 3;
}
@@ -465,7 +465,7 @@ message PbSnapshotMetaInfo {
int64 estimatedPartitionSize = 1;
repeated string registeredShuffle = 2;
repeated string hostnameSet = 3;
- repeated PbWorkerInfo blacklist = 4;
+ repeated PbWorkerInfo excludedWorkers = 4;
repeated PbWorkerInfo workerLostEvents = 5;
map<string, int64> appHeartbeatTime = 6;
repeated PbWorkerInfo workers = 7;
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 bafbb2bc0..8c66a0a15 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
@@ -627,11 +627,15 @@ object ControlMessages extends Logging {
.build().toByteArray
new TransportMessage(MessageType.HEARTBEAT_FROM_APPLICATION, payload)
- case HeartbeatFromApplicationResponse(statusCode, blacklist,
unknownWorkers, shuttingWorkers) =>
+ case HeartbeatFromApplicationResponse(
+ statusCode,
+ excludedWorkers,
+ unknownWorkers,
+ shuttingWorkers) =>
val payload = PbHeartbeatFromApplicationResponse.newBuilder()
.setStatus(statusCode.getValue)
- .addAllBlacklist(
- blacklist.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
+ .addAllExcludedWorkers(
+ excludedWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
.addAllUnknownWorkers(
unknownWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
.addAllShuttingWorkers(
@@ -641,7 +645,7 @@ object ControlMessages extends Logging {
case GetBlacklist(localExcludedWorkers) =>
val payload = PbGetBlacklist.newBuilder()
- .addAllLocalBlackList(localExcludedWorkers.asScala.map { workerInfo =>
+ .addAllLocalExcludedWorkers(localExcludedWorkers.asScala.map {
workerInfo =>
PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
}.toList.asJava)
.build().toByteArray
@@ -650,7 +654,7 @@ object ControlMessages extends Logging {
case GetBlacklistResponse(statusCode, excludedWorkers, unknownWorkers) =>
val builder = PbGetBlacklistResponse.newBuilder()
.setStatus(statusCode.getValue)
- builder.addAllBlacklist(
+ builder.addAllExcludedWorkers(
excludedWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
builder.addAllUnknownWorkers(
unknownWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
@@ -961,7 +965,7 @@ object ControlMessages extends Logging {
PbHeartbeatFromApplicationResponse.parseFrom(message.getPayload)
HeartbeatFromApplicationResponse(
Utils.toStatusCode(pbHeartbeatFromApplicationResponse.getStatus),
- pbHeartbeatFromApplicationResponse.getBlacklistList.asScala
+ pbHeartbeatFromApplicationResponse.getExcludedWorkersList.asScala
.map(PbSerDeUtils.fromPbWorkerInfo).toList.asJava,
pbHeartbeatFromApplicationResponse.getUnknownWorkersList.asScala
.map(PbSerDeUtils.fromPbWorkerInfo).toList.asJava,
@@ -970,14 +974,15 @@ object ControlMessages extends Logging {
case GET_BLACKLIST =>
val pbGetBlacklist = PbGetBlacklist.parseFrom(message.getPayload)
- GetBlacklist(new
util.ArrayList[WorkerInfo](pbGetBlacklist.getLocalBlackListList.asScala
- .map(PbSerDeUtils.fromPbWorkerInfo).toList.asJava))
+ GetBlacklist(
+ new
util.ArrayList[WorkerInfo](pbGetBlacklist.getLocalExcludedWorkersList.asScala
+ .map(PbSerDeUtils.fromPbWorkerInfo).toList.asJava))
case GET_BLACKLIST_RESPONSE =>
val pbGetBlacklistResponse =
PbGetBlacklistResponse.parseFrom(message.getPayload)
GetBlacklistResponse(
Utils.toStatusCode(pbGetBlacklistResponse.getStatus),
- pbGetBlacklistResponse.getBlacklistList.asScala
+ pbGetBlacklistResponse.getExcludedWorkersList.asScala
.map(PbSerDeUtils.fromPbWorkerInfo).toList.asJava,
pbGetBlacklistResponse.getUnknownWorkersList.asScala
.map(PbSerDeUtils.fromPbWorkerInfo).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 232cc1c30..992697656 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
@@ -378,7 +378,7 @@ object PbSerDeUtils {
.setEstimatedPartitionSize(estimatedPartitionSize)
.addAllRegisteredShuffle(registeredShuffle)
.addAllHostnameSet(hostnameSet)
- .addAllBlacklist(excludedWorkers.asScala.map(toPbWorkerInfo(_,
true)).asJava)
+ .addAllExcludedWorkers(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/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 2ab16409a..e2a2f1e40 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
@@ -270,7 +270,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
registeredShuffle.addAll(snapshotMetaInfo.getRegisteredShuffleList());
hostnameSet.addAll(snapshotMetaInfo.getHostnameSetList());
excludedWorkers.addAll(
- snapshotMetaInfo.getBlacklistList().stream()
+ snapshotMetaInfo.getExcludedWorkersList().stream()
.map(PbSerDeUtils::fromPbWorkerInfo)
.collect(Collectors.toSet()));
workerLostEvents.addAll(