This is an automated email from the ASF dual-hosted git repository.

ethanfeng pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/celeborn.git


The following commit(s) were added to refs/heads/main by this push:
     new 8e4ddaa46 [CELEBORN-1392] TransportClientFactory should regard as zero 
for negative celeborn.<module>.io.connectTimeout/connectionTimeout
8e4ddaa46 is described below

commit 8e4ddaa467387c21de62f87a018a3fcbf95bb5c0
Author: SteNicholas <[email protected]>
AuthorDate: Fri Apr 19 19:28:02 2024 +0800

    [CELEBORN-1392] TransportClientFactory should regard as zero for negative 
celeborn.<module>.io.connectTimeout/connectionTimeout
    
    ### What changes were proposed in this pull request?
    
    `TransportClientFactory` should regard as zero for negative 
`celeborn.<module>.io.connectTimeout` and 
`celeborn.<module>.io.connectionTimeout`.
    
    ### Why are the changes needed?
    
    When `celeborn.<module>.io.connectionTimeout` is 0 that means unlimited to 
netty, `ChannelFuture.await(0)` fails directly and inappropriately. Meanwhile, 
whhen `celeborn.<module>.io.connectionTimeout` is less than 0 that causes 
meaningless transport client reconnections and endless reconstructions.
    
    Backport:
    
    - https://github.com/apache/spark/pull/41785
    - https://github.com/apache/spark/pull/42619
    
    ### Does this PR introduce _any_ user-facing change?
    
    No.
    
    ### How was this patch tested?
    
    `TransportClientFactorySuiteJ#unlimitedConnectAndConnectionTimeouts`
    
    Closes #2467 from SteNicholas/CELEBORN-1392.
    
    Authored-by: SteNicholas <[email protected]>
    Signed-off-by: mingji <[email protected]>
---
 .../network/client/TransportClientFactory.java     | 10 +++++++-
 .../network/TransportClientFactorySuiteJ.java      | 27 ++++++++++++++++++++++
 2 files changed, 36 insertions(+), 1 deletion(-)

diff --git 
a/common/src/main/java/org/apache/celeborn/common/network/client/TransportClientFactory.java
 
b/common/src/main/java/org/apache/celeborn/common/network/client/TransportClientFactory.java
index 26191be26..a712965a4 100644
--- 
a/common/src/main/java/org/apache/celeborn/common/network/client/TransportClientFactory.java
+++ 
b/common/src/main/java/org/apache/celeborn/common/network/client/TransportClientFactory.java
@@ -251,7 +251,15 @@ public class TransportClientFactory implements Closeable {
     // Connect to the remote server
     long preConnect = System.nanoTime();
     ChannelFuture cf = bootstrap.connect(address);
-    if (!cf.await(connectTimeoutMs)) {
+    if (connectTimeoutMs <= 0) {
+      cf.await();
+      assert cf.isDone();
+      if (cf.isCancelled()) {
+        throw new IOException(String.format("Connecting to %s cancelled", 
address));
+      } else if (!cf.isSuccess()) {
+        throw new IOException(String.format("Failed to connect to %s", 
address), cf.cause());
+      }
+    } else if (!cf.await(connectTimeoutMs)) {
       throw new CelebornIOException(
           String.format("Connecting to %s timed out (%s ms)", address, 
connectTimeoutMs));
     } else if (cf.cause() != null) {
diff --git 
a/common/src/test/java/org/apache/celeborn/common/network/TransportClientFactorySuiteJ.java
 
b/common/src/test/java/org/apache/celeborn/common/network/TransportClientFactorySuiteJ.java
index 26a9b4885..b77a9c7d0 100644
--- 
a/common/src/test/java/org/apache/celeborn/common/network/TransportClientFactorySuiteJ.java
+++ 
b/common/src/test/java/org/apache/celeborn/common/network/TransportClientFactorySuiteJ.java
@@ -211,4 +211,31 @@ public class TransportClientFactorySuiteJ {
     factory.close();
     factory.createClient(getLocalHost(), server1.getPort());
   }
+
+  @Test
+  public void unlimitedConnectionAndCreationTimeouts() throws IOException, 
InterruptedException {
+    CelebornConf _conf = new CelebornConf();
+    _conf.set("celeborn.shuffle.io.connectTimeout", "-1");
+    _conf.set("celeborn.shuffle.io.connectionTimeout", "-1");
+    TransportConf conf = new TransportConf(TEST_MODULE, _conf);
+    try (TransportContext ctx = new TransportContext(conf, new 
BaseMessageHandler(), true);
+        TransportClientFactory factory = ctx.createClientFactory()) {
+      TransportClient c1 = factory.createClient(getLocalHost(), 
server1.getPort());
+      assertTrue(c1.isActive());
+      long expiredTime = System.currentTimeMillis() + 5000;
+      while (c1.isActive() && System.currentTimeMillis() < expiredTime) {
+        Thread.sleep(10);
+      }
+      assertTrue(c1.isActive());
+      // When connectionTimeout is unlimited, the connection shall be able to 
fail when the server
+      // is not reachable.
+      TransportServer server = ctx.createServer();
+      int unreachablePort = server.getPort();
+      JavaUtils.closeQuietly(server);
+      IOException exception =
+          assertThrows(
+              IOException.class, () -> factory.createClient(getLocalHost(), 
unreachablePort));
+      assertNotEquals(exception.getCause(), null);
+    }
+  }
 }

Reply via email to