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 bd465aa7a [CELEBORN-936] Shuffle master urls to avoid always connect 
first mast…
bd465aa7a is described below

commit bd465aa7a6d8da07b14aaf1776e4e54ebf97820a
Author: hongzhaoyang <[email protected]>
AuthorDate: Wed Aug 30 17:33:38 2023 +0800

    [CELEBORN-936] Shuffle master urls to avoid always connect first mast…
    
    ### What changes were proposed in this pull request?
    Shuffle master urls to avoid always connect first master first time
    
    ### Why are the changes needed?
    
    ### Does this PR introduce _any_ user-facing change?
    
    ### How was this patch tested?
    
    Closes #1866 from zy-jordan/CELEBORN-936.
    
    Authored-by: hongzhaoyang <[email protected]>
    Signed-off-by: zky.zhoukeyong <[email protected]>
---
 .../org/apache/celeborn/common/client/MasterClient.java     | 13 +++++++------
 1 file changed, 7 insertions(+), 6 deletions(-)

diff --git 
a/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java 
b/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java
index 8ad3f267a..e63e3b87b 100644
--- a/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java
+++ b/common/src/main/java/org/apache/celeborn/common/client/MasterClient.java
@@ -18,7 +18,7 @@
 package org.apache.celeborn.common.client;
 
 import java.io.IOException;
-import java.util.UUID;
+import java.util.*;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
@@ -52,7 +52,7 @@ public class MasterClient {
   private static final Logger LOG = 
LoggerFactory.getLogger(MasterClient.class);
 
   private final RpcEnv rpcEnv;
-  private final String[] masterEndpoints;
+  private final List<String> masterEndpoints;
   private final int maxRetries;
 
   private final RpcTimeout rpcTimeout;
@@ -62,8 +62,9 @@ public class MasterClient {
 
   public MasterClient(RpcEnv rpcEnv, CelebornConf conf) {
     this.rpcEnv = rpcEnv;
-    this.masterEndpoints = conf.masterEndpoints();
-    this.maxRetries = Math.max(masterEndpoints.length, 
conf.masterClientMaxRetries());
+    this.masterEndpoints = Arrays.asList(conf.masterEndpoints());
+    Collections.shuffle(this.masterEndpoints);
+    this.maxRetries = Math.max(masterEndpoints.size(), 
conf.masterClientMaxRetries());
     this.rpcTimeout = conf.masterClientRpcAskTimeout();
     this.rpcEndpointRef = new AtomicReference<>();
     this.oneWayMessageSender = 
ThreadUtils.newDaemonSingleThreadExecutor("One-Way-Message-Sender");
@@ -224,9 +225,9 @@ public class MasterClient {
     if (endpointRef == null) {
       int index = currentIndex.get();
       do {
-        RpcEndpointRef tempEndpointRef = 
setupEndpointRef(masterEndpoints[index]);
+        RpcEndpointRef tempEndpointRef = 
setupEndpointRef(masterEndpoints.get(index));
         if (rpcEndpointRef.compareAndSet(null, tempEndpointRef)) {
-          index = (index + 1) % masterEndpoints.length;
+          index = (index + 1) % masterEndpoints.size();
         }
         endpointRef = rpcEndpointRef.get();
       } while (endpointRef == null && index != currentIndex.get());

Reply via email to