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(

Reply via email to