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/incubator-celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new 30608ea69 [CELEBORN-1225] Worker should build replicate factory to get
client for sending replicate data
30608ea69 is described below
commit 30608ea698de7c068a50df5f22dd027764ede810
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]>
---
.../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 a0b0dd243..7b2cfe473 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -1432,26 +1432,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)
@@ -1459,21 +1491,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)
@@ -1482,7 +1536,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)
@@ -1490,7 +1554,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)
@@ -1500,7 +1574,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)
@@ -1508,7 +1584,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")
@@ -1517,7 +1597,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)
@@ -1536,7 +1618,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")
@@ -1557,9 +1641,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")
@@ -1569,9 +1653,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)
@@ -1581,7 +1665,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")
@@ -1591,7 +1675,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)
@@ -1602,10 +1686,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 22290b0bc..71b0bffdb 100644
--- a/docs/configuration/network.md
+++ b/docs/configuration/network.md
@@ -19,26 +19,26 @@ license: |
<!--begin-include-->
| Key | Default | Description | Since |
| --- | ------- | ----------- | ----- |
-| celeborn.<module>.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.<module>.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.<module>.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.<module>.io.backLog | 0 | Requested maximum length of the
queue of incoming connections. Default 0 for no backlog. | |
-| celeborn.<module>.io.clientThreads | 0 | Number of threads used in the
client thread pool. Default to 0, which is 2x#cores. | |
-| celeborn.<module>.io.connectTimeout | <value of
celeborn.network.connect.timeout> | Socket connect timeout. | |
-| celeborn.<module>.io.connectionTimeout | <value of
celeborn.network.timeout> | Connection active timeout. | |
+| celeborn.<module>.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.<module>.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.<module>.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.<module>.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.<module>.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.<module>.io.connectTimeout | <value of
celeborn.network.connect.timeout> | 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.<module>.io.connectionTimeout | <value of
celeborn.network.timeout> | 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.<module>.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.<module>.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.<module>.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.<module>.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.<module>.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.<module>.io.mode | NIO | Netty EventLoopGroup backend,
available options: NIO, EPOLL. | |
-| celeborn.<module>.io.numConnectionsPerPeer | 1 | Number of concurrent
connections between two nodes. | |
-| celeborn.<module>.io.preferDirectBufs | true | If true, we will prefer
allocating off-heap byte buffers within Netty. | |
-| celeborn.<module>.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.<module>.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.<module>.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.<module>.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.<module>.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.<module>.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.<module>.io.saslTimeout | 30s | Timeout for a single round
trip of auth message exchange, in milliseconds. | 0.5.0 |
-| celeborn.<module>.io.sendBuffer | 0b | Send buffer size (SO_SNDBUF). |
0.2.0 |
-| celeborn.<module>.io.serverThreads | 0 | Number of threads used in the
server thread pool. Default to 0, which is 2x#cores. | |
-| celeborn.<module>.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.<module>.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.<module>.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.<module>.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.<module>.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.<module>.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.<role>.rpc.dispatcher.threads | <value of
celeborn.rpc.dispatcher.threads> | 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 edc3dc75f..37243671a 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 = _