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

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


The following commit(s) were added to refs/heads/branch-0.4 by this push:
     new e968cf47e [CELEBORN-1225] Worker should build replicate factory to get 
client for sending replicate data
e968cf47e is described below

commit e968cf47e46b677fa61498baef37f4cebe929e96
Author: SteNicholas <[email protected]>
AuthorDate: Wed Jan 17 16:40:46 2024 +0800

    [CELEBORN-1225] Worker should build replicate factory to get client for 
sending replicate data
    
    ### What changes were proposed in this pull request?
    
    `PushDataHandler` should build replicate factory to get client for sending 
replicate data instead of push client factory.
    
    ### Why are the changes needed?
    
    `PushDataHandler` uses push client factory to create client for 
replicating, which should use replicate factory, otherwise replicate module 
configuration does not take effect for replicating of worker server.
    
    ### Does this PR introduce _any_ user-facing change?
    
    No.
    
    ### How was this patch tested?
    
    GA and cluster.
    
    Closes #2232 from SteNicholas/CELEBORN-1225.
    
    Authored-by: SteNicholas <[email protected]>
    Signed-off-by: mingji <[email protected]>
    (cherry picked from commit 30608ea698de7c068a50df5f22dd027764ede810)
    Signed-off-by: mingji <[email protected]>
---
 .../org/apache/celeborn/common/CelebornConf.scala  | 132 +++++++++++++++++----
 docs/configuration/network.md                      |  34 +++---
 .../service/deploy/worker/PushDataHandler.scala    |  14 +--
 .../celeborn/service/deploy/worker/Worker.scala    |  12 +-
 4 files changed, 139 insertions(+), 53 deletions(-)

diff --git 
a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala 
b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
index e2d34c18d..ced84bfde 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -1428,26 +1428,58 @@ object CelebornConf extends Logging {
   val NETWORK_IO_PREFER_DIRECT_BUFS: ConfigEntry[Boolean] =
     buildConf("celeborn.<module>.io.preferDirectBufs")
       .categories("network")
-      .doc("If true, we will prefer allocating off-heap byte buffers within 
Netty.")
+      .doc("If true, we will prefer allocating off-heap byte buffers within 
Netty. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client, master or worker. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for worker receiving push data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for replicate server or client of worker replicating data 
to peer worker. " +
+        s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, " 
+
+        s"it works for worker fetch server.")
       .booleanConf
       .createWithDefault(true)
 
   val NETWORK_IO_CONNECT_TIMEOUT: ConfigEntry[Long] =
     buildConf("celeborn.<module>.io.connectTimeout")
       .categories("network")
-      .doc("Socket connect timeout.")
+      .doc("Socket connect timeout. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for the replicate client of worker replicating data to peer 
worker.")
       .fallbackConf(NETWORK_CONNECT_TIMEOUT)
 
   val NETWORK_IO_CONNECTION_TIMEOUT: ConfigEntry[Long] =
     buildConf("celeborn.<module>.io.connectionTimeout")
       .categories("network")
-      .doc("Connection active timeout.")
+      .doc("Connection active timeout. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client, master or worker. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for worker receiving push data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for replicate server or client of worker replicating data 
to peer worker. " +
+        s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, " 
+
+        s"it works for worker fetch server.")
       .fallbackConf(NETWORK_TIMEOUT)
 
   val NETWORK_IO_NUM_CONNECTIONS_PER_PEER: ConfigEntry[Int] =
     buildConf("celeborn.<module>.io.numConnectionsPerPeer")
       .categories("network")
-      .doc("Number of concurrent connections between two nodes.")
+      .doc("Number of concurrent connections between two nodes. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for replicate client of worker replicating data to peer 
worker.")
       .intConf
       .createWithDefault(1)
 
@@ -1455,21 +1487,43 @@ object CelebornConf extends Logging {
     buildConf("celeborn.<module>.io.backLog")
       .categories("network")
       .doc(
-        "Requested maximum length of the queue of incoming connections. 
Default 0 for no backlog.")
+        "Requested maximum length of the queue of incoming connections. 
Default 0 for no backlog. " +
+          s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " 
+
+          s"it works for master or worker. " +
+          s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, 
" +
+          s"it works for worker receiving push data. " +
+          s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+          s"it works for replicate server of worker replicating data to peer 
worker. " +
+          s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, 
" +
+          s"it works for worker fetch server.")
       .intConf
       .createWithDefault(0)
 
   val NETWORK_IO_SERVER_THREADS: ConfigEntry[Int] =
     buildConf("celeborn.<module>.io.serverThreads")
       .categories("network")
-      .doc("Number of threads used in the server thread pool. Default to 0, 
which is 2x#cores.")
+      .doc("Number of threads used in the server thread pool. Default to 0, 
which is 2x#cores. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for master or worker. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for worker receiving push data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for replicate server of worker replicating data to peer 
worker. " +
+        s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, " 
+
+        s"it works for worker fetch server.")
       .intConf
       .createWithDefault(0)
 
   val NETWORK_IO_CLIENT_THREADS: ConfigEntry[Int] =
     buildConf("celeborn.<module>.io.clientThreads")
       .categories("network")
-      .doc("Number of threads used in the client thread pool. Default to 0, 
which is 2x#cores.")
+      .doc("Number of threads used in the client thread pool. Default to 0, 
which is 2x#cores. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for replicate client of worker replicating data to peer 
worker.")
       .intConf
       .createWithDefault(0)
 
@@ -1478,7 +1532,17 @@ object CelebornConf extends Logging {
       .categories("network")
       .doc("Receive buffer size (SO_RCVBUF). Note: the optimal size for 
receive buffer and send buffer " +
         "should be latency * network_bandwidth. Assuming latency = 1ms, 
network_bandwidth = 10Gbps " +
-        "buffer size should be ~ 1.25MB.")
+        "buffer size should be ~ 1.25MB. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client, master or worker. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for worker receiving push data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for replicate server or client of worker replicating data 
to peer worker. " +
+        s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, " 
+
+        s"it works for worker fetch server.")
       .version("0.2.0")
       .bytesConf(ByteUnit.BYTE)
       .createWithDefault(0)
@@ -1486,7 +1550,17 @@ object CelebornConf extends Logging {
   val NETWORK_IO_SEND_BUFFER: ConfigEntry[Long] =
     buildConf("celeborn.<module>.io.sendBuffer")
       .categories("network")
-      .doc("Send buffer size (SO_SNDBUF).")
+      .doc("Send buffer size (SO_SNDBUF). " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client, master or worker. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for worker receiving push data. " +
+        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
+        s"it works for replicate server or client of worker replicating data 
to peer worker. " +
+        s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, " 
+
+        s"it works for worker fetch server.")
       .version("0.2.0")
       .bytesConf(ByteUnit.BYTE)
       .createWithDefault(0)
@@ -1496,7 +1570,9 @@ object CelebornConf extends Logging {
       .categories("network")
       .doc(
         "Max number of times we will try IO exceptions (such as connection 
timeouts) per request. " +
-          "If set to 0, we will not do any retries.")
+          "If set to 0, we will not do any retries. " +
+          s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, 
" +
+          s"it works for Flink shuffle client push data.")
       .intConf
       .createWithDefault(3)
 
@@ -1504,7 +1580,11 @@ object CelebornConf extends Logging {
     buildConf("celeborn.<module>.io.retryWait")
       .categories("network")
       .doc("Time that we will wait in order to perform a retry after an 
IOException. " +
-        "Only relevant if maxIORetries > 0.")
+        "Only relevant if maxIORetries > 0. " +
+        s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
+        s"it works for shuffle client push and fetch data. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for Flink shuffle client push data.")
       .version("0.2.0")
       .timeConf(TimeUnit.MILLISECONDS)
       .createWithDefaultString("5s")
@@ -1513,7 +1593,9 @@ object CelebornConf extends Logging {
     buildConf("celeborn.<module>.io.lazyFD")
       .categories("network")
       .doc("Whether to initialize FileDescriptor lazily or not. If true, file 
descriptors are created only " +
-        "when data is going to be transferred. This can reduce the number of 
open files.")
+        "when data is going to be transferred. This can reduce the number of 
open files. " +
+        s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, " 
+
+        s"it works for worker fetch server.")
       .booleanConf
       .createWithDefault(true)
 
@@ -1532,7 +1614,9 @@ object CelebornConf extends Logging {
       .internal
       .doc("Minimum size of a block that we should start using memory map 
rather than reading in through " +
         "normal IO operations. This prevents Celeborn from memory mapping very 
small blocks. In general, " +
-        "memory mapping has high overhead for blocks close to or below the 
page size of the OS.")
+        "memory mapping has high overhead for blocks close to or below the 
page size of the OS. " +
+        s"If setting <module> to `${TransportModuleConstants.FETCH_MODULE}`, " 
+
+        s"it works for worker fetch server.")
       .version("0.3.0")
       .bytesConf(ByteUnit.BYTE)
       .createWithDefaultString("2m")
@@ -1553,9 +1637,9 @@ object CelebornConf extends Logging {
       .categories("network")
       .doc("Interval for checking push data timeout. " +
         s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
-        s"it works for shuffle client push data and should be configured on 
client side. " +
-        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
-        s"it works for worker replicate data to peer worker and should be 
configured on worker side.")
+        s"it works for shuffle client push data. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for Flink shuffle client push data.")
       .version("0.3.0")
       .timeConf(TimeUnit.MILLISECONDS)
       .createWithDefaultString("5s")
@@ -1565,9 +1649,9 @@ object CelebornConf extends Logging {
       .categories("network")
       .doc("Threads num for checking push data timeout. " +
         s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
-        s"it works for shuffle client push data and should be configured on 
client side. " +
-        s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
-        s"it works for worker replicate data to peer worker and should be 
configured on worker side.")
+        s"it works for shuffle client push data. " +
+        s"If setting <module> to `${TransportModuleConstants.PUSH_MODULE}`, " +
+        s"it works for Flink shuffle client push data.")
       .version("0.3.0")
       .intConf
       .createWithDefault(4)
@@ -1577,7 +1661,7 @@ object CelebornConf extends Logging {
       .categories("network")
       .doc("Interval for checking fetch data timeout. " +
         s"It only support setting <module> to 
`${TransportModuleConstants.DATA_MODULE}` " +
-        s"since it works for shuffle client fetch data and should be 
configured on client side.")
+        s"since it works for shuffle client fetch data.")
       .version("0.3.0")
       .timeConf(TimeUnit.MILLISECONDS)
       .createWithDefaultString("5s")
@@ -1587,7 +1671,7 @@ object CelebornConf extends Logging {
       .categories("network")
       .doc("Threads num for checking fetch data timeout. " +
         s"It only support setting <module> to 
`${TransportModuleConstants.DATA_MODULE}` " +
-        s"since it works for shuffle client fetch data and should be 
configured on client side.")
+        s"since it works for shuffle client fetch data.")
       .version("0.3.0")
       .intConf
       .createWithDefault(4)
@@ -1598,10 +1682,12 @@ object CelebornConf extends Logging {
       .categories("network")
       .version("0.3.0")
       .doc("The heartbeat interval between worker and client. " +
+        s"If setting <module> to `${TransportModuleConstants.RPC_MODULE}`, " +
+        s"it works for shuffle client. " +
         s"If setting <module> to `${TransportModuleConstants.DATA_MODULE}`, " +
-        s"it works for shuffle client push and fetch data and should be 
configured on client side. " +
+        s"it works for shuffle client push and fetch data. " +
         s"If setting <module> to 
`${TransportModuleConstants.REPLICATE_MODULE}`, " +
-        s"it works for worker replicate data to peer worker and should be 
configured on worker side.")
+        s"it works for replicate client of worker replicating data to peer 
worker.")
       .timeConf(TimeUnit.MILLISECONDS)
       .createWithDefaultString("60s")
 
diff --git a/docs/configuration/network.md b/docs/configuration/network.md
index fb30a705b..7cc84e35a 100644
--- a/docs/configuration/network.md
+++ b/docs/configuration/network.md
@@ -19,26 +19,26 @@ license: |
 <!--begin-include-->
 | Key | Default | Description | Since |
 | --- | ------- | ----------- | ----- |
-| celeborn.&lt;module&gt;.fetch.timeoutCheck.interval | 5s | Interval for 
checking fetch data timeout. It only support setting <module> to `data` since 
it works for shuffle client fetch data and should be configured on client side. 
| 0.3.0 | 
-| celeborn.&lt;module&gt;.fetch.timeoutCheck.threads | 4 | Threads num for 
checking fetch data timeout. It only support setting <module> to `data` since 
it works for shuffle client fetch data and should be configured on client side. 
| 0.3.0 | 
-| celeborn.&lt;module&gt;.heartbeat.interval | 60s | The heartbeat interval 
between worker and client. If setting <module> to `data`, it works for shuffle 
client push and fetch data and should be configured on client side. If setting 
<module> to `replicate`, it works for worker replicate data to peer worker and 
should be configured on worker side. | 0.3.0 | 
-| celeborn.&lt;module&gt;.io.backLog | 0 | Requested maximum length of the 
queue of incoming connections. Default 0 for no backlog. |  | 
-| celeborn.&lt;module&gt;.io.clientThreads | 0 | Number of threads used in the 
client thread pool. Default to 0, which is 2x#cores. |  | 
-| celeborn.&lt;module&gt;.io.connectTimeout | &lt;value of 
celeborn.network.connect.timeout&gt; | Socket connect timeout. |  | 
-| celeborn.&lt;module&gt;.io.connectionTimeout | &lt;value of 
celeborn.network.timeout&gt; | Connection active timeout. |  | 
+| celeborn.&lt;module&gt;.fetch.timeoutCheck.interval | 5s | Interval for 
checking fetch data timeout. It only support setting <module> to `data` since 
it works for shuffle client fetch data. | 0.3.0 | 
+| celeborn.&lt;module&gt;.fetch.timeoutCheck.threads | 4 | Threads num for 
checking fetch data timeout. It only support setting <module> to `data` since 
it works for shuffle client fetch data. | 0.3.0 | 
+| celeborn.&lt;module&gt;.heartbeat.interval | 60s | The heartbeat interval 
between worker and client. If setting <module> to `rpc`, it works for shuffle 
client. If setting <module> to `data`, it works for shuffle client push and 
fetch data. If setting <module> to `replicate`, it works for replicate client 
of worker replicating data to peer worker. | 0.3.0 | 
+| celeborn.&lt;module&gt;.io.backLog | 0 | Requested maximum length of the 
queue of incoming connections. Default 0 for no backlog. If setting <module> to 
`rpc`, it works for master or worker. If setting <module> to `push`, it works 
for worker receiving push data. If setting <module> to `replicate`, it works 
for replicate server of worker replicating data to peer worker. If setting 
<module> to `fetch`, it works for worker fetch server. |  | 
+| celeborn.&lt;module&gt;.io.clientThreads | 0 | Number of threads used in the 
client thread pool. Default to 0, which is 2x#cores. If setting <module> to 
`rpc`, it works for shuffle client. If setting <module> to `data`, it works for 
shuffle client push and fetch data. If setting <module> to `replicate`, it 
works for replicate client of worker replicating data to peer worker. |  | 
+| celeborn.&lt;module&gt;.io.connectTimeout | &lt;value of 
celeborn.network.connect.timeout&gt; | Socket connect timeout. If setting 
<module> to `rpc`, it works for shuffle client. If setting <module> to `data`, 
it works for shuffle client push and fetch data. If setting <module> to 
`replicate`, it works for the replicate client of worker replicating data to 
peer worker. |  | 
+| celeborn.&lt;module&gt;.io.connectionTimeout | &lt;value of 
celeborn.network.timeout&gt; | Connection active timeout. If setting <module> 
to `rpc`, it works for shuffle client, master or worker. If setting <module> to 
`data`, it works for shuffle client push and fetch data. If setting <module> to 
`push`, it works for worker receiving push data. If setting <module> to 
`replicate`, it works for replicate server or client of worker replicating data 
to peer worker. If setting <module> to ` [...]
 | celeborn.&lt;module&gt;.io.enableVerboseMetrics | false | Whether to track 
Netty memory detailed metrics. If true, the detailed metrics of Netty 
PoolByteBufAllocator will be gotten, otherwise only general memory usage will 
be tracked. |  | 
-| celeborn.&lt;module&gt;.io.lazyFD | true | Whether to initialize 
FileDescriptor lazily or not. If true, file descriptors are created only when 
data is going to be transferred. This can reduce the number of open files. |  | 
-| celeborn.&lt;module&gt;.io.maxRetries | 3 | Max number of times we will try 
IO exceptions (such as connection timeouts) per request. If set to 0, we will 
not do any retries. |  | 
+| celeborn.&lt;module&gt;.io.lazyFD | true | Whether to initialize 
FileDescriptor lazily or not. If true, file descriptors are created only when 
data is going to be transferred. This can reduce the number of open files. If 
setting <module> to `fetch`, it works for worker fetch server. |  | 
+| celeborn.&lt;module&gt;.io.maxRetries | 3 | Max number of times we will try 
IO exceptions (such as connection timeouts) per request. If set to 0, we will 
not do any retries. If setting <module> to `push`, it works for Flink shuffle 
client push data. |  | 
 | celeborn.&lt;module&gt;.io.mode | NIO | Netty EventLoopGroup backend, 
available options: NIO, EPOLL. |  | 
-| celeborn.&lt;module&gt;.io.numConnectionsPerPeer | 1 | Number of concurrent 
connections between two nodes. |  | 
-| celeborn.&lt;module&gt;.io.preferDirectBufs | true | If true, we will prefer 
allocating off-heap byte buffers within Netty. |  | 
-| celeborn.&lt;module&gt;.io.receiveBuffer | 0b | Receive buffer size 
(SO_RCVBUF). Note: the optimal size for receive buffer and send buffer should 
be latency * network_bandwidth. Assuming latency = 1ms, network_bandwidth = 
10Gbps buffer size should be ~ 1.25MB. | 0.2.0 | 
-| celeborn.&lt;module&gt;.io.retryWait | 5s | Time that we will wait in order 
to perform a retry after an IOException. Only relevant if maxIORetries > 0. | 
0.2.0 | 
+| celeborn.&lt;module&gt;.io.numConnectionsPerPeer | 1 | Number of concurrent 
connections between two nodes. If setting <module> to `rpc`, it works for 
shuffle client. If setting <module> to `data`, it works for shuffle client push 
and fetch data. If setting <module> to `replicate`, it works for replicate 
client of worker replicating data to peer worker. |  | 
+| celeborn.&lt;module&gt;.io.preferDirectBufs | true | If true, we will prefer 
allocating off-heap byte buffers within Netty. If setting <module> to `rpc`, it 
works for shuffle client, master or worker. If setting <module> to `data`, it 
works for shuffle client push and fetch data. If setting <module> to `push`, it 
works for worker receiving push data. If setting <module> to `replicate`, it 
works for replicate server or client of worker replicating data to peer worker. 
If setting <module [...]
+| celeborn.&lt;module&gt;.io.receiveBuffer | 0b | Receive buffer size 
(SO_RCVBUF). Note: the optimal size for receive buffer and send buffer should 
be latency * network_bandwidth. Assuming latency = 1ms, network_bandwidth = 
10Gbps buffer size should be ~ 1.25MB. If setting <module> to `rpc`, it works 
for shuffle client, master or worker. If setting <module> to `data`, it works 
for shuffle client push and fetch data. If setting <module> to `push`, it works 
for worker receiving push data.  [...]
+| celeborn.&lt;module&gt;.io.retryWait | 5s | Time that we will wait in order 
to perform a retry after an IOException. Only relevant if maxIORetries > 0. If 
setting <module> to `data`, it works for shuffle client push and fetch data. If 
setting <module> to `push`, it works for Flink shuffle client push data. | 
0.2.0 | 
 | celeborn.&lt;module&gt;.io.saslTimeout | 30s | Timeout for a single round 
trip of auth message exchange, in milliseconds. | 0.5.0 | 
-| celeborn.&lt;module&gt;.io.sendBuffer | 0b | Send buffer size (SO_SNDBUF). | 
0.2.0 | 
-| celeborn.&lt;module&gt;.io.serverThreads | 0 | Number of threads used in the 
server thread pool. Default to 0, which is 2x#cores. |  | 
-| celeborn.&lt;module&gt;.push.timeoutCheck.interval | 5s | Interval for 
checking push data timeout. If setting <module> to `data`, it works for shuffle 
client push data and should be configured on client side. If setting <module> 
to `replicate`, it works for worker replicate data to peer worker and should be 
configured on worker side. | 0.3.0 | 
-| celeborn.&lt;module&gt;.push.timeoutCheck.threads | 4 | Threads num for 
checking push data timeout. If setting <module> to `data`, it works for shuffle 
client push data and should be configured on client side. If setting <module> 
to `replicate`, it works for worker replicate data to peer worker and should be 
configured on worker side. | 0.3.0 | 
+| celeborn.&lt;module&gt;.io.sendBuffer | 0b | Send buffer size (SO_SNDBUF). 
If setting <module> to `rpc`, it works for shuffle client, master or worker. If 
setting <module> to `data`, it works for shuffle client push and fetch data. If 
setting <module> to `push`, it works for worker receiving push data. If setting 
<module> to `replicate`, it works for replicate server or client of worker 
replicating data to peer worker. If setting <module> to `fetch`, it works for 
worker fetch server. | [...]
+| celeborn.&lt;module&gt;.io.serverThreads | 0 | Number of threads used in the 
server thread pool. Default to 0, which is 2x#cores. If setting <module> to 
`rpc`, it works for master or worker. If setting <module> to `push`, it works 
for worker receiving push data. If setting <module> to `replicate`, it works 
for replicate server of worker replicating data to peer worker. If setting 
<module> to `fetch`, it works for worker fetch server. |  | 
+| celeborn.&lt;module&gt;.push.timeoutCheck.interval | 5s | Interval for 
checking push data timeout. If setting <module> to `data`, it works for shuffle 
client push data. If setting <module> to `push`, it works for Flink shuffle 
client push data. | 0.3.0 | 
+| celeborn.&lt;module&gt;.push.timeoutCheck.threads | 4 | Threads num for 
checking push data timeout. If setting <module> to `data`, it works for shuffle 
client push data. If setting <module> to `push`, it works for Flink shuffle 
client push data. | 0.3.0 | 
 | celeborn.&lt;role&gt;.rpc.dispatcher.threads | &lt;value of 
celeborn.rpc.dispatcher.threads&gt; | Threads number of message dispatcher 
event loop for roles |  | 
 | celeborn.io.maxDefaultNettyThreads | 64 | Max default netty threads | 0.3.2 
| 
 | celeborn.network.bind.preferIpAddress | true | When `ture`, prefer to use IP 
address, otherwise FQDN. This configuration only takes effects when the bind 
hostname is not set explicitly, in such case, Celeborn will find the first 
non-loopback address to bind. | 0.3.0 | 
diff --git 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/PushDataHandler.scala
 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/PushDataHandler.scala
index 6705325e0..509a33209 100644
--- 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/PushDataHandler.scala
+++ 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/PushDataHandler.scala
@@ -55,7 +55,7 @@ class PushDataHandler(val workerSource: WorkerSource) extends 
BaseMessageHandler
   private var shufflePushDataTimeout: ConcurrentHashMap[String, Long] = _
   private var replicateThreadPool: ThreadPoolExecutor = _
   private var unavailablePeers: ConcurrentHashMap[WorkerInfo, Long] = _
-  private var pushClientFactory: TransportClientFactory = _
+  private var replicateClientFactory: TransportClientFactory = _
   private var registered: AtomicBoolean = _
   private var workerInfo: WorkerInfo = _
   private var diskReserveSize: Long = _
@@ -78,7 +78,7 @@ class PushDataHandler(val workerSource: WorkerSource) extends 
BaseMessageHandler
     shuffleMapperAttempts = worker.shuffleMapperAttempts
     replicateThreadPool = worker.replicateThreadPool
     unavailablePeers = worker.unavailablePeers
-    pushClientFactory = worker.pushClientFactory
+    replicateClientFactory = worker.replicateClientFactory
     registered = worker.registered
     workerInfo = worker.workerInfo
     diskReserveSize = worker.conf.workerDiskReserveSize
@@ -347,7 +347,7 @@ class PushDataHandler(val workerSource: WorkerSource) 
extends BaseMessageHandler
             }
           }
           try {
-            val client = getClient(peer.getHost, peer.getReplicatePort, 
location.getId)
+            val client = getReplicateClient(peer.getHost, 
peer.getReplicatePort, location.getId)
             val newPushData = new PushData(
               PartitionLocation.Mode.REPLICA.mode(),
               shuffleKey,
@@ -610,7 +610,7 @@ class PushDataHandler(val workerSource: WorkerSource) 
extends BaseMessageHandler
           }
 
           try {
-            val client = getClient(peer.getHost, peer.getReplicatePort, 
location.getId)
+            val client = getReplicateClient(peer.getHost, 
peer.getReplicatePort, location.getId)
             val newPushMergedData = new PushMergedData(
               PartitionLocation.Mode.REPLICA.mode(),
               shuffleKey,
@@ -1216,11 +1216,11 @@ class PushDataHandler(val workerSource: WorkerSource) 
extends BaseMessageHandler
     false
   }
 
-  private def getClient(host: String, port: Int, partitionId: Int): 
TransportClient = {
+  private def getReplicateClient(host: String, port: Int, partitionId: Int): 
TransportClient = {
     if (workerReplicateRandomConnectionEnabled) {
-      pushClientFactory.createClient(host, port)
+      replicateClientFactory.createClient(host, port)
     } else {
-      pushClientFactory.createClient(host, port, partitionId)
+      replicateClientFactory.createClient(host, port, partitionId)
     }
   }
 
diff --git 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
index dc4152c86..643a077c8 100644
--- 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
+++ 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
@@ -141,7 +141,7 @@ private[celeborn] class Worker(
   rpcEnv.setupEndpoint(RpcNameConstants.WORKER_EP, controller)
 
   val pushDataHandler = new PushDataHandler(workerSource)
-  val (pushServer, pushClientFactory) = {
+  private val pushServer = {
     val closeIdleConnections = conf.workerCloseIdleConnections
     val numThreads = 
conf.workerPushIoThreads.getOrElse(storageManager.totalFlusherThread)
     val transportConf =
@@ -155,13 +155,11 @@ private[celeborn] class Worker(
         pushServerLimiter,
         conf.workerPushHeartbeatEnabled,
         workerSource)
-    (
-      transportContext.createServer(conf.workerPushPort),
-      transportContext.createClientFactory())
+    transportContext.createServer(conf.workerPushPort)
   }
 
   val replicateHandler = new PushDataHandler(workerSource)
-  private val replicateServer = {
+  val (replicateServer, replicateClientFactory) = {
     val closeIdleConnections = conf.workerCloseIdleConnections
     val numThreads =
       
conf.workerReplicateIoThreads.getOrElse(storageManager.totalFlusherThread)
@@ -176,7 +174,9 @@ private[celeborn] class Worker(
         replicateLimiter,
         false,
         workerSource)
-    transportContext.createServer(conf.workerReplicatePort)
+    (
+      transportContext.createServer(conf.workerReplicatePort),
+      transportContext.createClientFactory())
   }
 
   var fetchHandler: FetchHandler = _

Reply via email to