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 afab4a0a3 [CELEBORN-696][FOLLOWUP] Remove new allocated peer workers
from pushExecludedWrkers
afab4a0a3 is described below
commit afab4a0a3bff8c6730570a77f025d4863fd3ef88
Author: Angerszhuuuu <[email protected]>
AuthorDate: Wed Jun 28 17:38:36 2023 +0800
[CELEBORN-696][FOLLOWUP] Remove new allocated peer workers from
pushExecludedWrkers
### What changes were proposed in this pull request?
Remove new allocated location's workers from pushExecludedWrkers should
also remove peers
### Why are the changes needed?
### Does this PR introduce _any_ user-facing change?
### How was this patch tested?
Closes #1636 from AngersZhuuuu/CELEBORN-696-FOLLOWUP.
Authored-by: Angerszhuuuu <[email protected]>
Signed-off-by: zky.zhoukeyong <[email protected]>
---
.../java/org/apache/celeborn/client/ShuffleClientImpl.java | 11 +++++++++--
1 file changed, 9 insertions(+), 2 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 0a18dc767..1324db481 100644
--- a/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
+++ b/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
@@ -530,7 +530,9 @@ public class ShuffleClientImpl extends ShuffleClient {
PartitionLocation partitionLoc =
PbSerDeUtils.fromPbPartitionLocation(response.getPartitionLocationsList().get(i));
pushExcludedWorkers.remove(partitionLoc.hostAndPushPort());
-
+ if (partitionLoc.hasPeer()) {
+
pushExcludedWorkers.remove(partitionLoc.getPeer().hostAndPushPort());
+ }
result.put(partitionLoc.getId(), partitionLoc);
}
return result;
@@ -706,7 +708,9 @@ public class ShuffleClientImpl extends ShuffleClient {
int partitionId = partitionInfo.getPartitionId();
int statusCode = partitionInfo.getStatus();
if (partitionInfo.getOldAvailable()) {
-
pushExcludedWorkers.remove(oldLocMap.get(partitionId).hostAndPushPort());
+ PartitionLocation oldLoc = oldLocMap.get(partitionId);
+ // Currently, revive only check if main location available, here
won't remove peer loc.
+ pushExcludedWorkers.remove(oldLoc.hostAndPushPort());
}
if (StatusCode.SUCCESS.getValue() == statusCode) {
@@ -714,6 +718,9 @@ public class ShuffleClientImpl extends ShuffleClient {
PbSerDeUtils.fromPbPartitionLocation(partitionInfo.getPartition());
partitionLocationMap.put(partitionId, loc);
pushExcludedWorkers.remove(loc.hostAndPushPort());
+ if (loc.hasPeer()) {
+ pushExcludedWorkers.remove(loc.getPeer().hostAndPushPort());
+ }
} else if (StatusCode.STAGE_ENDED.getValue() == statusCode) {
stageEnded(shuffleId);
return results;