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);
+ }
+ }
}