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 88d59cd97 [CELEBORN-1247] Output config's alternatives to doc ### What
changes were proposed in this pull request? Add configs' alternatives to doc.
88d59cd97 is described below
commit 88d59cd971c283e83234a18af5af4bd7c754c54b
Author: mingji <[email protected]>
AuthorDate: Wed Jan 24 11:21:23 2024 +0800
[CELEBORN-1247] Output config's alternatives to doc
### What changes were proposed in this pull request?
Add configs' alternatives to doc.
### Why are the changes needed?
To help users use correct configs.
### Does this PR introduce _any_ user-facing change?
NO.
### How was this patch tested?
GA.
Closes #2253 from FMX/b1241.
Authored-by: mingji <[email protected]>
Signed-off-by: mingji <[email protected]>
---
.../org/apache/celeborn/common/CelebornConf.scala | 7 +-
.../org/apache/celeborn/ConfigurationSuite.scala | 7 +-
docs/configuration/client.md | 188 +++++++++---------
docs/configuration/columnar-shuffle.md | 16 +-
docs/configuration/ha.md | 16 +-
docs/configuration/master.md | 54 +++---
docs/configuration/metrics.md | 32 +--
docs/configuration/network.md | 72 +++----
docs/configuration/quota.md | 16 +-
docs/configuration/worker.md | 216 ++++++++++-----------
10 files changed, 316 insertions(+), 308 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 36da26a2c..b707f5b29 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -1691,7 +1691,12 @@ object CelebornConf extends Logging {
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.")
+ s"it works for replicate client of worker replicating data to peer
worker." +
+ "If you are using the \"celeborn.client.heartbeat.interval\", " +
+ "please use the new configs for each module according to your needs or
" +
+ "replace it with \"celeborn.rpc.heartbeat.interval\", " +
+ "\"celeborn.data.heartbeat.interval\" and" +
+ "\"celeborn.replicate.heartbeat.interval\". ")
.timeConf(TimeUnit.MILLISECONDS)
.createWithDefaultString("60s")
diff --git a/common/src/test/scala/org/apache/celeborn/ConfigurationSuite.scala
b/common/src/test/scala/org/apache/celeborn/ConfigurationSuite.scala
index 1b7240fbc..1a2599a66 100644
--- a/common/src/test/scala/org/apache/celeborn/ConfigurationSuite.scala
+++ b/common/src/test/scala/org/apache/celeborn/ConfigurationSuite.scala
@@ -104,7 +104,8 @@ class ConfigurationSuite extends AnyFunSuite {
s"${escape(entry.key)}",
s"${escape(entry.defaultValueString)}",
s"${entry.doc}",
- s"${entry.version}")
+ s"${entry.version}",
+ s"${escape(entry.alternatives.map(_._1).mkString(","))}")
output += seq.mkString("| ", " | ", " | ")
}
appendEndInclude(output)
@@ -138,8 +139,8 @@ class ConfigurationSuite extends AnyFunSuite {
}
def appendConfigurationTableHeader(output: ArrayBuffer[String]): Unit = {
- output += "| Key | Default | Description | Since |"
- output += "| --- | ------- | ----------- | ----- |"
+ output += "| Key | Default | Description | Since | Deprecated |"
+ output += "| --- | ------- | ----------- | ----- | ---------- |"
}
def appendBeginInclude(output: ArrayBuffer[String]): Unit = {
diff --git a/docs/configuration/client.md b/docs/configuration/client.md
index 76a3aefbb..284edf157 100644
--- a/docs/configuration/client.md
+++ b/docs/configuration/client.md
@@ -17,97 +17,99 @@ license: |
---
<!--begin-include-->
-| Key | Default | Description | Since |
-| --- | ------- | ----------- | ----- |
-| celeborn.client.application.heartbeatInterval | 10s | Interval for client to
send heartbeat message to master. | 0.3.0 |
-| celeborn.client.application.unregister.enabled | true | When true, Celeborn
client will inform celeborn master the application is already shutdown during
client exit, this allows the cluster to release resources immediately,
resulting in resource savings. | 0.3.2 |
-| celeborn.client.closeIdleConnections | true | Whether client will close idle
connections. | 0.3.0 |
-| celeborn.client.commitFiles.ignoreExcludedWorker | false | When true,
LifecycleManager will skip workers which are in the excluded list. | 0.3.0 |
-| celeborn.client.eagerlyCreateInputStream.threads | 32 | Threads count for
streamCreatorPool in CelebornShuffleReader. | 0.3.1 |
-| celeborn.client.excludePeerWorkerOnFailure.enabled | true | When true,
Celeborn will exclude partition's peer worker on failure when push data to
replica failed. | 0.3.0 |
-| celeborn.client.excludedWorker.expireTimeout | 180s | Timeout time for
LifecycleManager to clear reserved excluded worker. Default to be 1.5 *
`celeborn.master.heartbeat.worker.timeout`to cover worker heartbeat timeout
check period | 0.3.0 |
-| celeborn.client.fetch.buffer.size | 64k | Size of reducer partition buffer
memory for shuffle reader. The fetched data will be buffered in memory before
consuming. For performance consideration keep this buffer size not less than
`celeborn.client.push.buffer.max.size`. | 0.4.0 |
-| celeborn.client.fetch.dfsReadChunkSize | 8m | Max chunk size for
DfsPartitionReader. | 0.3.1 |
-| celeborn.client.fetch.excludeWorkerOnFailure.enabled | false | Whether to
enable shuffle client-side fetch exclude workers on failure. | 0.3.0 |
-| celeborn.client.fetch.excludedWorker.expireTimeout | <value of
celeborn.client.excludedWorker.expireTimeout> | ShuffleClient is a static
object, it will be used in the whole lifecycle of Executor,We give a expire
time for excluded workers to avoid a transient worker issues. | 0.3.0 |
-| celeborn.client.fetch.maxReqsInFlight | 3 | Amount of in-flight chunk fetch
request. | 0.3.0 |
-| celeborn.client.fetch.maxRetriesForEachReplica | 3 | Max retry times of
fetch chunk on each replica | 0.3.0 |
-| celeborn.client.fetch.timeout | 600s | Timeout for a task to open stream and
fetch chunk. | 0.3.0 |
-| celeborn.client.flink.compression.enabled | true | Whether to compress data
in Flink plugin. | 0.3.0 |
-| celeborn.client.flink.inputGate.concurrentReadings | 2147483647 | Max
concurrent reading channels for a input gate. | 0.3.0 |
-| celeborn.client.flink.inputGate.memory | 32m | Memory reserved for a input
gate. | 0.3.0 |
-| celeborn.client.flink.inputGate.minMemory | 8m | Min memory reserved for a
input gate. | 0.3.0 |
-| celeborn.client.flink.inputGate.supportFloatingBuffer | true | Whether to
support floating buffer in Flink input gates. | 0.3.0 |
-| celeborn.client.flink.resultPartition.memory | 64m | Memory reserved for a
result partition. | 0.3.0 |
-| celeborn.client.flink.resultPartition.minMemory | 8m | Min memory reserved
for a result partition. | 0.3.0 |
-| celeborn.client.flink.resultPartition.supportFloatingBuffer | true | Whether
to support floating buffer for result partitions. | 0.3.0 |
-| celeborn.client.mr.pushData.max | 32m | Max size for a push data sent from
mr client. | 0.4.0 |
-| celeborn.client.push.buffer.initial.size | 8k | | 0.3.0 |
-| celeborn.client.push.buffer.max.size | 64k | Max size of reducer partition
buffer memory for shuffle hash writer. The pushed data will be buffered in
memory before sending to Celeborn worker. For performance consideration keep
this buffer size higher than 32K. Example: If reducer amount is 2000, buffer
size is 64K, then each task will consume up to `64KiB * 2000 = 125MiB` heap
memory. | 0.3.0 |
-| celeborn.client.push.excludeWorkerOnFailure.enabled | false | Whether to
enable shuffle client-side push exclude workers on failures. | 0.3.0 |
-| celeborn.client.push.limit.inFlight.sleepInterval | 50ms | Sleep interval
when check netty in-flight requests to be done. | 0.3.0 |
-| celeborn.client.push.limit.inFlight.timeout | <undefined> | Timeout
for netty in-flight requests to be done.Default value should be
`celeborn.client.push.timeout * 2`. | 0.3.0 |
-| celeborn.client.push.limit.strategy | SIMPLE | The strategy used to control
the push speed. Valid strategies are SIMPLE and SLOWSTART. The SLOWSTART
strategy usually works with congestion control mechanism on the worker side. |
0.3.0 |
-| celeborn.client.push.maxReqsInFlight.perWorker | 32 | Amount of Netty
in-flight requests per worker. Default max memory of in flight requests per
worker is `celeborn.client.push.maxReqsInFlight.perWorker` *
`celeborn.client.push.buffer.max.size` * compression ratio(1 in worst case):
64KiB * 32 = 2MiB. The maximum memory will not exceed
`celeborn.client.push.maxReqsInFlight.total`. | 0.3.0 |
-| celeborn.client.push.maxReqsInFlight.total | 256 | Amount of total Netty
in-flight requests. The maximum memory is
`celeborn.client.push.maxReqsInFlight.total` *
`celeborn.client.push.buffer.max.size` * compression ratio(1 in worst case):
64KiB * 256 = 16MiB | 0.3.0 |
-| celeborn.client.push.queue.capacity | 512 | Push buffer queue size for a
task. The maximum memory is `celeborn.client.push.buffer.max.size` *
`celeborn.client.push.queue.capacity`, default: 64KiB * 512 = 32MiB | 0.3.0 |
-| celeborn.client.push.replicate.enabled | false | When true, Celeborn worker
will replicate shuffle data to another Celeborn worker asynchronously to ensure
the pushed shuffle data won't be lost after the node failure. It's recommended
to set `false` when `HDFS` is enabled in `celeborn.storage.activeTypes`. |
0.3.0 |
-| celeborn.client.push.retry.threads | 8 | Thread number to process shuffle
re-send push data requests. | 0.3.0 |
-| celeborn.client.push.revive.batchSize | 2048 | Max number of partitions in
one Revive request. | 0.3.0 |
-| celeborn.client.push.revive.interval | 100ms | Interval for client to
trigger Revive to LifecycleManager. The number of partitions in one Revive
request is `celeborn.client.push.revive.batchSize`. | 0.3.0 |
-| celeborn.client.push.revive.maxRetries | 5 | Max retry times for reviving
when celeborn push data failed. | 0.3.0 |
-| celeborn.client.push.sendBufferPool.checkExpireInterval | 30s | Interval to
check expire for send buffer pool. If the pool has been idle for more than
`celeborn.client.push.sendBufferPool.expireTimeout`, the pooled send buffers
and push tasks will be cleaned up. | 0.3.1 |
-| celeborn.client.push.sendBufferPool.expireTimeout | 60s | Timeout before
clean up SendBufferPool. If SendBufferPool is idle for more than this time, the
send buffers and push tasks will be cleaned up. | 0.3.1 |
-| celeborn.client.push.slowStart.initialSleepTime | 500ms | The initial sleep
time if the current max in flight requests is 0 | 0.3.0 |
-| celeborn.client.push.slowStart.maxSleepTime | 2s | If
celeborn.client.push.limit.strategy is set to SLOWSTART, push side will take a
sleep strategy for each batch of requests, this controls the max sleep time if
the max in flight requests limit is 1 for a long time | 0.3.0 |
-| celeborn.client.push.sort.randomizePartitionId.enabled | false | Whether to
randomize partitionId in push sorter. If true, partitionId will be randomized
when sort data to avoid skew when push to worker | 0.3.0 |
-| celeborn.client.push.splitPartition.threads | 8 | Thread number to process
shuffle split request in shuffle client. | 0.3.0 |
-| celeborn.client.push.stageEnd.timeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for waiting
StageEnd. During this process, there are
`celeborn.client.requestCommitFiles.maxRetries` times for retry opportunities
for committing files and 1 times for releasing slots request. User can
customize this value according to your setting. By default, the value is the
max timeout value `celeborn.<module>.io.connectionTimeout`. | 0.3.0 |
-| celeborn.client.push.takeTaskMaxWaitAttempts | 1 | Max wait times if no task
available to push to worker. | 0.3.0 |
-| celeborn.client.push.takeTaskWaitInterval | 50ms | Wait interval if no task
available to push to worker. | 0.3.0 |
-| celeborn.client.push.timeout | 120s | Timeout for a task to push data rpc
message. This value should better be more than twice of
`celeborn.<module>.push.timeoutCheck.interval` | 0.3.0 |
-| celeborn.client.readLocalShuffleFile.enabled | false | Enable read local
shuffle file for clusters that co-deployed with yarn node manager. | 0.3.1 |
-| celeborn.client.readLocalShuffleFile.threads | 4 | Threads count for read
local shuffle file. | 0.3.1 |
-| celeborn.client.registerShuffle.maxRetries | 3 | Max retry times for client
to register shuffle. | 0.3.0 |
-| celeborn.client.registerShuffle.retryWait | 3s | Wait time before next retry
if register shuffle failed. | 0.3.0 |
-| celeborn.client.requestCommitFiles.maxRetries | 4 | Max retry times for
requestCommitFiles RPC. | 0.3.0 |
-| celeborn.client.reserveSlots.maxRetries | 3 | Max retry times for client to
reserve slots. | 0.3.0 |
-| celeborn.client.reserveSlots.rackaware.enabled | false | Whether need to
place different replicates on different racks when allocating slots. | 0.3.1 |
-| celeborn.client.reserveSlots.retryWait | 3s | Wait time before next retry if
reserve slots failed. | 0.3.0 |
-| celeborn.client.rpc.cache.concurrencyLevel | 32 | The number of write locks
to update rpc cache. | 0.3.0 |
-| celeborn.client.rpc.cache.expireTime | 15s | The time before a cache item is
removed. | 0.3.0 |
-| celeborn.client.rpc.cache.size | 256 | The max cache items count for rpc
cache. | 0.3.0 |
-| celeborn.client.rpc.getReducerFileGroup.askTimeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for ask operations
during getting reducer file group information. During this process, there are
`celeborn.client.requestCommitFiles.maxRetries` times for retry opportunities
for committing files and 1 times for releasing slots request. User can
customize this value according to your setting. By default, the value is the
max timeout value `celeborn.<module>.io.co [...]
-| celeborn.client.rpc.maxRetries | 3 | Max RPC retry times in
LifecycleManager. | 0.3.2 |
-| celeborn.client.rpc.registerShuffle.askTimeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for ask operations
during register shuffle. During this process, there are two times for retry
opportunities for requesting slots, one request for establishing a connection
with Worker and `celeborn.client.reserveSlots.maxRetries` times for retry
opportunities for reserving slots. User can customize this value according to
your setting. By default, the value is the m [...]
-| celeborn.client.rpc.requestPartition.askTimeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for ask operations
during requesting change partition location, such as reviving or splitting
partition. During this process, there are
`celeborn.client.reserveSlots.maxRetries` times for retry opportunities for
reserving slots. User can customize this value according to your setting. By
default, the value is the max timeout value `celeborn.<module>.io.connectionTim
[...]
-| celeborn.client.rpc.reserveSlots.askTimeout | <value of
celeborn.rpc.askTimeout> | Timeout for LifecycleManager request reserve
slots. | 0.3.0 |
-| celeborn.client.rpc.shared.threads | 16 | Number of shared rpc threads in
LifecycleManager. | 0.3.2 |
-| celeborn.client.shuffle.batchHandleChangePartition.interval | 100ms |
Interval for LifecycleManager to schedule handling change partition requests in
batch. | 0.3.0 |
-| celeborn.client.shuffle.batchHandleChangePartition.threads | 8 | Threads
number for LifecycleManager to handle change partition request in batch. |
0.3.0 |
-| celeborn.client.shuffle.batchHandleCommitPartition.interval | 5s | Interval
for LifecycleManager to schedule handling commit partition requests in batch. |
0.3.0 |
-| celeborn.client.shuffle.batchHandleCommitPartition.threads | 8 | Threads
number for LifecycleManager to handle commit partition request in batch. |
0.3.0 |
-| celeborn.client.shuffle.batchHandleReleasePartition.interval | 5s | Interval
for LifecycleManager to schedule handling release partition requests in batch.
| 0.3.0 |
-| celeborn.client.shuffle.batchHandleReleasePartition.threads | 8 | Threads
number for LifecycleManager to handle release partition request in batch. |
0.3.0 |
-| celeborn.client.shuffle.compression.codec | LZ4 | The codec used to compress
shuffle data. By default, Celeborn provides three codecs: `lz4`, `zstd`,
`none`. | 0.3.0 |
-| celeborn.client.shuffle.compression.zstd.level | 1 | Compression level for
Zstd compression codec, its value should be an integer between -5 and 22.
Increasing the compression level will result in better compression at the
expense of more CPU and memory. | 0.3.0 |
-| celeborn.client.shuffle.decompression.lz4.xxhash.instance |
<undefined> | Decompression XXHash instance for Lz4. Available options:
JNI, JAVASAFE, JAVAUNSAFE. | 0.3.2 |
-| celeborn.client.shuffle.expired.checkInterval | 60s | Interval for client to
check expired shuffles. | 0.3.0 |
-| celeborn.client.shuffle.manager.port | 0 | Port used by the LifecycleManager
on the Driver. | 0.3.0 |
-| celeborn.client.shuffle.mapPartition.split.enabled | false | whether to
enable shuffle partition split. Currently, this only applies to MapPartition. |
0.3.1 |
-| celeborn.client.shuffle.partition.type | REDUCE | Type of shuffle's
partition. | 0.3.0 |
-| celeborn.client.shuffle.partitionSplit.mode | SOFT | soft: the shuffle file
size might be larger than split threshold. hard: the shuffle file size will be
limited to split threshold. | 0.3.0 |
-| celeborn.client.shuffle.partitionSplit.threshold | 1G | Shuffle file size
threshold, if file size exceeds this, trigger split. | 0.3.0 |
-| celeborn.client.shuffle.rangeReadFilter.enabled | false | If a spark
application have skewed partition, this value can set to true to improve
performance. | 0.2.0 |
-| celeborn.client.shuffle.register.filterExcludedWorker.enabled | false |
Whether to filter excluded worker when register shuffle. | 0.4.0 |
-| celeborn.client.slot.assign.maxWorkers | 10000 | Max workers that slots of
one shuffle can be allocated on. Will choose the smaller positive one from
Master side and Client side, see `celeborn.master.slot.assign.maxWorkers`. |
0.3.1 |
-| celeborn.client.spark.fetch.throwsFetchFailure | false | client throws
FetchFailedException instead of CelebornIOException | 0.4.0 |
-| celeborn.client.spark.push.sort.memory.threshold | 64m | When
SortBasedPusher use memory over the threshold, will trigger push data. | 0.3.0
|
-| celeborn.client.spark.push.unsafeRow.fastWrite.enabled | true | This is
Celeborn's optimization on UnsafeRow for Spark and it's true by default. If you
have changed UnsafeRow's memory layout set this to false. | 0.2.2 |
-| celeborn.client.spark.shuffle.forceFallback.enabled | false | Whether force
fallback shuffle to Spark's default. | 0.3.0 |
-| celeborn.client.spark.shuffle.forceFallback.numPartitionsThreshold |
2147483647 | Celeborn will only accept shuffle of partition number lower than
this configuration value. | 0.3.0 |
-| celeborn.client.spark.shuffle.writer | HASH | Celeborn supports the
following kind of shuffle writers. 1. hash: hash-based shuffle writer works
fine when shuffle partition count is normal; 2. sort: sort-based shuffle writer
works fine when memory pressure is high or shuffle partition count is huge. |
0.3.0 |
-| celeborn.master.endpoints | <localhost>:9097 | Endpoints of master
nodes for celeborn client to connect, allowed pattern is:
`<host1>:<port1>[,<host2>:<port2>]*`, e.g. `clb1:9097,clb2:9098,clb3:9099`. If
the port is omitted, 9097 will be used. | 0.2.0 |
-| celeborn.storage.availableTypes | HDD | Enabled storages. Available options:
MEMORY,HDD,SSD,HDFS. Note: HDD and SSD would be treated as identical. | 0.3.0 |
-| celeborn.storage.hdfs.dir | <undefined> | HDFS base directory for
Celeborn to store shuffle data. | 0.2.0 |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.client.application.heartbeatInterval | 10s | Interval for client to
send heartbeat message to master. | 0.3.0 |
celeborn.application.heartbeatInterval |
+| celeborn.client.application.unregister.enabled | true | When true, Celeborn
client will inform celeborn master the application is already shutdown during
client exit, this allows the cluster to release resources immediately,
resulting in resource savings. | 0.3.2 | |
+| celeborn.client.closeIdleConnections | true | Whether client will close idle
connections. | 0.3.0 | |
+| celeborn.client.commitFiles.ignoreExcludedWorker | false | When true,
LifecycleManager will skip workers which are in the excluded list. | 0.3.0 | |
+| celeborn.client.eagerlyCreateInputStream.threads | 32 | Threads count for
streamCreatorPool in CelebornShuffleReader. | 0.3.1 | |
+| celeborn.client.excludePeerWorkerOnFailure.enabled | true | When true,
Celeborn will exclude partition's peer worker on failure when push data to
replica failed. | 0.3.0 | |
+| celeborn.client.excludedWorker.expireTimeout | 180s | Timeout time for
LifecycleManager to clear reserved excluded worker. Default to be 1.5 *
`celeborn.master.heartbeat.worker.timeout`to cover worker heartbeat timeout
check period | 0.3.0 | celeborn.worker.excluded.expireTimeout |
+| celeborn.client.fetch.buffer.size | 64k | Size of reducer partition buffer
memory for shuffle reader. The fetched data will be buffered in memory before
consuming. For performance consideration keep this buffer size not less than
`celeborn.client.push.buffer.max.size`. | 0.4.0 | |
+| celeborn.client.fetch.dfsReadChunkSize | 8m | Max chunk size for
DfsPartitionReader. | 0.3.1 | |
+| celeborn.client.fetch.excludeWorkerOnFailure.enabled | false | Whether to
enable shuffle client-side fetch exclude workers on failure. | 0.3.0 | |
+| celeborn.client.fetch.excludedWorker.expireTimeout | <value of
celeborn.client.excludedWorker.expireTimeout> | ShuffleClient is a static
object, it will be used in the whole lifecycle of Executor,We give a expire
time for excluded workers to avoid a transient worker issues. | 0.3.0 | |
+| celeborn.client.fetch.maxReqsInFlight | 3 | Amount of in-flight chunk fetch
request. | 0.3.0 | celeborn.fetch.maxReqsInFlight |
+| celeborn.client.fetch.maxRetriesForEachReplica | 3 | Max retry times of
fetch chunk on each replica | 0.3.0 |
celeborn.fetch.maxRetriesForEachReplica,celeborn.fetch.maxRetries |
+| celeborn.client.fetch.timeout | 600s | Timeout for a task to open stream and
fetch chunk. | 0.3.0 | celeborn.fetch.timeout |
+| celeborn.client.flink.compression.enabled | true | Whether to compress data
in Flink plugin. | 0.3.0 | remote-shuffle.job.enable-data-compression |
+| celeborn.client.flink.inputGate.concurrentReadings | 2147483647 | Max
concurrent reading channels for a input gate. | 0.3.0 |
remote-shuffle.job.concurrent-readings-per-gate |
+| celeborn.client.flink.inputGate.memory | 32m | Memory reserved for a input
gate. | 0.3.0 | remote-shuffle.job.memory-per-gate |
+| celeborn.client.flink.inputGate.minMemory | 8m | Min memory reserved for a
input gate. | 0.3.0 | remote-shuffle.job.min.memory-per-gate |
+| celeborn.client.flink.inputGate.supportFloatingBuffer | true | Whether to
support floating buffer in Flink input gates. | 0.3.0 |
remote-shuffle.job.support-floating-buffer-per-input-gate |
+| celeborn.client.flink.resultPartition.memory | 64m | Memory reserved for a
result partition. | 0.3.0 | remote-shuffle.job.memory-per-partition |
+| celeborn.client.flink.resultPartition.minMemory | 8m | Min memory reserved
for a result partition. | 0.3.0 | remote-shuffle.job.min.memory-per-partition |
+| celeborn.client.flink.resultPartition.supportFloatingBuffer | true | Whether
to support floating buffer for result partitions. | 0.3.0 |
remote-shuffle.job.support-floating-buffer-per-output-gate |
+| celeborn.client.mr.pushData.max | 32m | Max size for a push data sent from
mr client. | 0.4.0 | |
+| celeborn.client.push.buffer.initial.size | 8k | | 0.3.0 |
celeborn.push.buffer.initial.size |
+| celeborn.client.push.buffer.max.size | 64k | Max size of reducer partition
buffer memory for shuffle hash writer. The pushed data will be buffered in
memory before sending to Celeborn worker. For performance consideration keep
this buffer size higher than 32K. Example: If reducer amount is 2000, buffer
size is 64K, then each task will consume up to `64KiB * 2000 = 125MiB` heap
memory. | 0.3.0 | celeborn.push.buffer.max.size |
+| celeborn.client.push.excludeWorkerOnFailure.enabled | false | Whether to
enable shuffle client-side push exclude workers on failures. | 0.3.0 | |
+| celeborn.client.push.limit.inFlight.sleepInterval | 50ms | Sleep interval
when check netty in-flight requests to be done. | 0.3.0 |
celeborn.push.limit.inFlight.sleepInterval |
+| celeborn.client.push.limit.inFlight.timeout | <undefined> | Timeout
for netty in-flight requests to be done.Default value should be
`celeborn.client.push.timeout * 2`. | 0.3.0 |
celeborn.push.limit.inFlight.timeout |
+| celeborn.client.push.limit.strategy | SIMPLE | The strategy used to control
the push speed. Valid strategies are SIMPLE and SLOWSTART. The SLOWSTART
strategy usually works with congestion control mechanism on the worker side. |
0.3.0 | |
+| celeborn.client.push.maxReqsInFlight.perWorker | 32 | Amount of Netty
in-flight requests per worker. Default max memory of in flight requests per
worker is `celeborn.client.push.maxReqsInFlight.perWorker` *
`celeborn.client.push.buffer.max.size` * compression ratio(1 in worst case):
64KiB * 32 = 2MiB. The maximum memory will not exceed
`celeborn.client.push.maxReqsInFlight.total`. | 0.3.0 | |
+| celeborn.client.push.maxReqsInFlight.total | 256 | Amount of total Netty
in-flight requests. The maximum memory is
`celeborn.client.push.maxReqsInFlight.total` *
`celeborn.client.push.buffer.max.size` * compression ratio(1 in worst case):
64KiB * 256 = 16MiB | 0.3.0 | celeborn.push.maxReqsInFlight |
+| celeborn.client.push.queue.capacity | 512 | Push buffer queue size for a
task. The maximum memory is `celeborn.client.push.buffer.max.size` *
`celeborn.client.push.queue.capacity`, default: 64KiB * 512 = 32MiB | 0.3.0 |
celeborn.push.queue.capacity |
+| celeborn.client.push.replicate.enabled | false | When true, Celeborn worker
will replicate shuffle data to another Celeborn worker asynchronously to ensure
the pushed shuffle data won't be lost after the node failure. It's recommended
to set `false` when `HDFS` is enabled in `celeborn.storage.activeTypes`. |
0.3.0 | celeborn.push.replicate.enabled |
+| celeborn.client.push.retry.threads | 8 | Thread number to process shuffle
re-send push data requests. | 0.3.0 | celeborn.push.retry.threads |
+| celeborn.client.push.revive.batchSize | 2048 | Max number of partitions in
one Revive request. | 0.3.0 | |
+| celeborn.client.push.revive.interval | 100ms | Interval for client to
trigger Revive to LifecycleManager. The number of partitions in one Revive
request is `celeborn.client.push.revive.batchSize`. | 0.3.0 | |
+| celeborn.client.push.revive.maxRetries | 5 | Max retry times for reviving
when celeborn push data failed. | 0.3.0 | |
+| celeborn.client.push.sendBufferPool.checkExpireInterval | 30s | Interval to
check expire for send buffer pool. If the pool has been idle for more than
`celeborn.client.push.sendBufferPool.expireTimeout`, the pooled send buffers
and push tasks will be cleaned up. | 0.3.1 | |
+| celeborn.client.push.sendBufferPool.expireTimeout | 60s | Timeout before
clean up SendBufferPool. If SendBufferPool is idle for more than this time, the
send buffers and push tasks will be cleaned up. | 0.3.1 | |
+| celeborn.client.push.slowStart.initialSleepTime | 500ms | The initial sleep
time if the current max in flight requests is 0 | 0.3.0 | |
+| celeborn.client.push.slowStart.maxSleepTime | 2s | If
celeborn.client.push.limit.strategy is set to SLOWSTART, push side will take a
sleep strategy for each batch of requests, this controls the max sleep time if
the max in flight requests limit is 1 for a long time | 0.3.0 | |
+| celeborn.client.push.sort.randomizePartitionId.enabled | false | Whether to
randomize partitionId in push sorter. If true, partitionId will be randomized
when sort data to avoid skew when push to worker | 0.3.0 |
celeborn.push.sort.randomizePartitionId.enabled |
+| celeborn.client.push.splitPartition.threads | 8 | Thread number to process
shuffle split request in shuffle client. | 0.3.0 |
celeborn.push.splitPartition.threads |
+| celeborn.client.push.stageEnd.timeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for waiting
StageEnd. During this process, there are
`celeborn.client.requestCommitFiles.maxRetries` times for retry opportunities
for committing files and 1 times for releasing slots request. User can
customize this value according to your setting. By default, the value is the
max timeout value `celeborn.<module>.io.connectionTimeout`. | 0.3.0 |
celeborn.push.stageEnd.timeout |
+| celeborn.client.push.takeTaskMaxWaitAttempts | 1 | Max wait times if no task
available to push to worker. | 0.3.0 | |
+| celeborn.client.push.takeTaskWaitInterval | 50ms | Wait interval if no task
available to push to worker. | 0.3.0 | |
+| celeborn.client.push.timeout | 120s | Timeout for a task to push data rpc
message. This value should better be more than twice of
`celeborn.<module>.push.timeoutCheck.interval` | 0.3.0 |
celeborn.push.data.timeout |
+| celeborn.client.readLocalShuffleFile.enabled | false | Enable read local
shuffle file for clusters that co-deployed with yarn node manager. | 0.3.1 | |
+| celeborn.client.readLocalShuffleFile.threads | 4 | Threads count for read
local shuffle file. | 0.3.1 | |
+| celeborn.client.registerShuffle.maxRetries | 3 | Max retry times for client
to register shuffle. | 0.3.0 | celeborn.shuffle.register.maxRetries |
+| celeborn.client.registerShuffle.retryWait | 3s | Wait time before next retry
if register shuffle failed. | 0.3.0 | celeborn.shuffle.register.retryWait |
+| celeborn.client.requestCommitFiles.maxRetries | 4 | Max retry times for
requestCommitFiles RPC. | 0.3.0 | |
+| celeborn.client.reserveSlots.maxRetries | 3 | Max retry times for client to
reserve slots. | 0.3.0 | celeborn.slots.reserve.maxRetries |
+| celeborn.client.reserveSlots.rackaware.enabled | false | Whether need to
place different replicates on different racks when allocating slots. | 0.3.1 |
celeborn.client.reserveSlots.rackware.enabled |
+| celeborn.client.reserveSlots.retryWait | 3s | Wait time before next retry if
reserve slots failed. | 0.3.0 | celeborn.slots.reserve.retryWait |
+| celeborn.client.rpc.cache.concurrencyLevel | 32 | The number of write locks
to update rpc cache. | 0.3.0 | celeborn.rpc.cache.concurrencyLevel |
+| celeborn.client.rpc.cache.expireTime | 15s | The time before a cache item is
removed. | 0.3.0 | celeborn.rpc.cache.expireTime |
+| celeborn.client.rpc.cache.size | 256 | The max cache items count for rpc
cache. | 0.3.0 | celeborn.rpc.cache.size |
+| celeborn.client.rpc.getReducerFileGroup.askTimeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for ask operations
during getting reducer file group information. During this process, there are
`celeborn.client.requestCommitFiles.maxRetries` times for retry opportunities
for committing files and 1 times for releasing slots request. User can
customize this value according to your setting. By default, the value is the
max timeout value `celeborn.<module>.io.co [...]
+| celeborn.client.rpc.maxRetries | 3 | Max RPC retry times in
LifecycleManager. | 0.3.2 | |
+| celeborn.client.rpc.registerShuffle.askTimeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for ask operations
during register shuffle. During this process, there are two times for retry
opportunities for requesting slots, one request for establishing a connection
with Worker and `celeborn.client.reserveSlots.maxRetries` times for retry
opportunities for reserving slots. User can customize this value according to
your setting. By default, the value is the m [...]
+| celeborn.client.rpc.requestPartition.askTimeout | <value of
celeborn.<module>.io.connectionTimeout> | Timeout for ask operations
during requesting change partition location, such as reviving or splitting
partition. During this process, there are
`celeborn.client.reserveSlots.maxRetries` times for retry opportunities for
reserving slots. User can customize this value according to your setting. By
default, the value is the max timeout value `celeborn.<module>.io.connectionTim
[...]
+| celeborn.client.rpc.reserveSlots.askTimeout | <value of
celeborn.rpc.askTimeout> | Timeout for LifecycleManager request reserve
slots. | 0.3.0 | |
+| celeborn.client.rpc.shared.threads | 16 | Number of shared rpc threads in
LifecycleManager. | 0.3.2 | |
+| celeborn.client.shuffle.batchHandleChangePartition.interval | 100ms |
Interval for LifecycleManager to schedule handling change partition requests in
batch. | 0.3.0 | celeborn.shuffle.batchHandleChangePartition.interval |
+| celeborn.client.shuffle.batchHandleChangePartition.threads | 8 | Threads
number for LifecycleManager to handle change partition request in batch. |
0.3.0 | celeborn.shuffle.batchHandleChangePartition.threads |
+| celeborn.client.shuffle.batchHandleCommitPartition.interval | 5s | Interval
for LifecycleManager to schedule handling commit partition requests in batch. |
0.3.0 | celeborn.shuffle.batchHandleCommitPartition.interval |
+| celeborn.client.shuffle.batchHandleCommitPartition.threads | 8 | Threads
number for LifecycleManager to handle commit partition request in batch. |
0.3.0 | celeborn.shuffle.batchHandleCommitPartition.threads |
+| celeborn.client.shuffle.batchHandleReleasePartition.interval | 5s | Interval
for LifecycleManager to schedule handling release partition requests in batch.
| 0.3.0 | |
+| celeborn.client.shuffle.batchHandleReleasePartition.threads | 8 | Threads
number for LifecycleManager to handle release partition request in batch. |
0.3.0 | |
+| celeborn.client.shuffle.compression.codec | LZ4 | The codec used to compress
shuffle data. By default, Celeborn provides three codecs: `lz4`, `zstd`,
`none`. | 0.3.0 |
celeborn.shuffle.compression.codec,remote-shuffle.job.compression.codec |
+| celeborn.client.shuffle.compression.zstd.level | 1 | Compression level for
Zstd compression codec, its value should be an integer between -5 and 22.
Increasing the compression level will result in better compression at the
expense of more CPU and memory. | 0.3.0 |
celeborn.shuffle.compression.zstd.level |
+| celeborn.client.shuffle.decompression.lz4.xxhash.instance |
<undefined> | Decompression XXHash instance for Lz4. Available options:
JNI, JAVASAFE, JAVAUNSAFE. | 0.3.2 | |
+| celeborn.client.shuffle.expired.checkInterval | 60s | Interval for client to
check expired shuffles. | 0.3.0 | celeborn.shuffle.expired.checkInterval |
+| celeborn.client.shuffle.manager.port | 0 | Port used by the LifecycleManager
on the Driver. | 0.3.0 | celeborn.shuffle.manager.port |
+| celeborn.client.shuffle.mapPartition.split.enabled | false | whether to
enable shuffle partition split. Currently, this only applies to MapPartition. |
0.3.1 | |
+| celeborn.client.shuffle.partition.type | REDUCE | Type of shuffle's
partition. | 0.3.0 | celeborn.shuffle.partition.type |
+| celeborn.client.shuffle.partitionSplit.mode | SOFT | soft: the shuffle file
size might be larger than split threshold. hard: the shuffle file size will be
limited to split threshold. | 0.3.0 | celeborn.shuffle.partitionSplit.mode |
+| celeborn.client.shuffle.partitionSplit.threshold | 1G | Shuffle file size
threshold, if file size exceeds this, trigger split. | 0.3.0 |
celeborn.shuffle.partitionSplit.threshold |
+| celeborn.client.shuffle.rangeReadFilter.enabled | false | If a spark
application have skewed partition, this value can set to true to improve
performance. | 0.2.0 | celeborn.shuffle.rangeReadFilter.enabled |
+| celeborn.client.shuffle.register.filterExcludedWorker.enabled | false |
Whether to filter excluded worker when register shuffle. | 0.4.0 | |
+| celeborn.client.slot.assign.maxWorkers | 10000 | Max workers that slots of
one shuffle can be allocated on. Will choose the smaller positive one from
Master side and Client side, see `celeborn.master.slot.assign.maxWorkers`. |
0.3.1 | |
+| celeborn.client.spark.fetch.throwsFetchFailure | false | client throws
FetchFailedException instead of CelebornIOException | 0.4.0 | |
+| celeborn.client.spark.push.dynamicWriteMode.enabled | false | Whether to
dynamically switch push write mode based on conditions.If true, shuffle mode
will be only determined by partition count | 0.5.0 | |
+| celeborn.client.spark.push.dynamicWriteMode.partitionNum.threshold | 2000 |
Threshold of shuffle partition number for dynamically switching push writer
mode. When the shuffle partition number is greater than this value, use the
sort-based shuffle writer for memory efficiency; otherwise use the hash-based
shuffle writer for speed. This configuration only takes effect when
celeborn.client.spark.push.dynamicWriteMode.enabled is true. | 0.5.0 | |
+| celeborn.client.spark.push.sort.memory.threshold | 64m | When
SortBasedPusher use memory over the threshold, will trigger push data. | 0.3.0
| celeborn.push.sortMemory.threshold |
+| celeborn.client.spark.push.unsafeRow.fastWrite.enabled | true | This is
Celeborn's optimization on UnsafeRow for Spark and it's true by default. If you
have changed UnsafeRow's memory layout set this to false. | 0.2.2 | |
+| celeborn.client.spark.shuffle.forceFallback.enabled | false | Whether force
fallback shuffle to Spark's default. | 0.3.0 |
celeborn.shuffle.forceFallback.enabled |
+| celeborn.client.spark.shuffle.forceFallback.numPartitionsThreshold |
2147483647 | Celeborn will only accept shuffle of partition number lower than
this configuration value. | 0.3.0 |
celeborn.shuffle.forceFallback.numPartitionsThreshold |
+| celeborn.client.spark.shuffle.writer | HASH | Celeborn supports the
following kind of shuffle writers. 1. hash: hash-based shuffle writer works
fine when shuffle partition count is normal; 2. sort: sort-based shuffle writer
works fine when memory pressure is high or shuffle partition count is huge.
This configuration only takes effect when
celeborn.client.spark.push.dynamicWriteMode.enabled is false. | 0.3.0 |
celeborn.shuffle.writer |
+| celeborn.master.endpoints | <localhost>:9097 | Endpoints of master
nodes for celeborn client to connect, allowed pattern is:
`<host1>:<port1>[,<host2>:<port2>]*`, e.g. `clb1:9097,clb2:9098,clb3:9099`. If
the port is omitted, 9097 will be used. | 0.2.0 | |
+| celeborn.storage.availableTypes | HDD | Enabled storages. Available options:
MEMORY,HDD,SSD,HDFS. Note: HDD and SSD would be treated as identical. | 0.3.0 |
celeborn.storage.activeTypes |
+| celeborn.storage.hdfs.dir | <undefined> | HDFS base directory for
Celeborn to store shuffle data. | 0.2.0 | |
<!--end-include-->
diff --git a/docs/configuration/columnar-shuffle.md
b/docs/configuration/columnar-shuffle.md
index 65d60f0a3..942793b5f 100644
--- a/docs/configuration/columnar-shuffle.md
+++ b/docs/configuration/columnar-shuffle.md
@@ -17,12 +17,12 @@ license: |
---
<!--begin-include-->
-| Key | Default | Description | Since |
-| --- | ------- | ----------- | ----- |
-| celeborn.columnarShuffle.batch.size | 10000 | Vector batch size for columnar
shuffle. | 0.3.0 |
-| celeborn.columnarShuffle.codegen.enabled | false | Whether to use codegen
for columnar-based shuffle. | 0.3.0 |
-| celeborn.columnarShuffle.enabled | false | Whether to enable columnar-based
shuffle. | 0.2.0 |
-| celeborn.columnarShuffle.encoding.dictionary.enabled | false | Whether to
use dictionary encoding for columnar-based shuffle data. | 0.3.0 |
-| celeborn.columnarShuffle.encoding.dictionary.maxFactor | 0.3 | Max factor
for dictionary size. The max dictionary size is `min(32.0 KiB,
celeborn.columnarShuffle.batch.size *
celeborn.columnar.shuffle.encoding.dictionary.maxFactor)`. | 0.3.0 |
-| celeborn.columnarShuffle.offHeap.enabled | false | Whether to use off heap
columnar vector. | 0.3.0 |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.columnarShuffle.batch.size | 10000 | Vector batch size for columnar
shuffle. | 0.3.0 | celeborn.columnar.shuffle.batch.size |
+| celeborn.columnarShuffle.codegen.enabled | false | Whether to use codegen
for columnar-based shuffle. | 0.3.0 | celeborn.columnar.shuffle.codegen.enabled
|
+| celeborn.columnarShuffle.enabled | false | Whether to enable columnar-based
shuffle. | 0.2.0 | celeborn.columnar.shuffle.enabled |
+| celeborn.columnarShuffle.encoding.dictionary.enabled | false | Whether to
use dictionary encoding for columnar-based shuffle data. | 0.3.0 |
celeborn.columnar.shuffle.encoding.dictionary.enabled |
+| celeborn.columnarShuffle.encoding.dictionary.maxFactor | 0.3 | Max factor
for dictionary size. The max dictionary size is `min(32.0 KiB,
celeborn.columnarShuffle.batch.size *
celeborn.columnar.shuffle.encoding.dictionary.maxFactor)`. | 0.3.0 |
celeborn.columnar.shuffle.encoding.dictionary.maxFactor |
+| celeborn.columnarShuffle.offHeap.enabled | false | Whether to use off heap
columnar vector. | 0.3.0 | celeborn.columnar.offHeap.enabled |
<!--end-include-->
diff --git a/docs/configuration/ha.md b/docs/configuration/ha.md
index 494537ded..a6b2ea53f 100644
--- a/docs/configuration/ha.md
+++ b/docs/configuration/ha.md
@@ -17,12 +17,12 @@ license: |
---
<!--begin-include-->
-| Key | Default | Description | Since |
-| --- | ------- | ----------- | ----- |
-| celeborn.master.ha.enabled | false | When true, master nodes run as Raft
cluster mode. | 0.3.0 |
-| celeborn.master.ha.node.<id>.host | <required> | Host to bind of
master node <id> in HA mode. | 0.3.0 |
-| celeborn.master.ha.node.<id>.port | 9097 | Port to bind of master node
<id> in HA mode. | 0.3.0 |
-| celeborn.master.ha.node.<id>.ratis.port | 9872 | Ratis port to bind of
master node <id> in HA mode. | 0.3.0 |
-| celeborn.master.ha.ratis.raft.rpc.type | netty | RPC type for Ratis,
available options: netty, grpc. | 0.3.0 |
-| celeborn.master.ha.ratis.raft.server.storage.dir | /tmp/ratis | | 0.3.0 |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.master.ha.enabled | false | When true, master nodes run as Raft
cluster mode. | 0.3.0 | celeborn.ha.enabled |
+| celeborn.master.ha.node.<id>.host | <required> | Host to bind of
master node <id> in HA mode. | 0.3.0 | celeborn.ha.master.node.<id>.host
|
+| celeborn.master.ha.node.<id>.port | 9097 | Port to bind of master node
<id> in HA mode. | 0.3.0 | celeborn.ha.master.node.<id>.port |
+| celeborn.master.ha.node.<id>.ratis.port | 9872 | Ratis port to bind of
master node <id> in HA mode. | 0.3.0 |
celeborn.ha.master.node.<id>.ratis.port |
+| celeborn.master.ha.ratis.raft.rpc.type | netty | RPC type for Ratis,
available options: netty, grpc. | 0.3.0 |
celeborn.ha.master.ratis.raft.rpc.type |
+| celeborn.master.ha.ratis.raft.server.storage.dir | /tmp/ratis | | 0.3.0 |
celeborn.ha.master.ratis.raft.server.storage.dir |
<!--end-include-->
diff --git a/docs/configuration/master.md b/docs/configuration/master.md
index 92ddae700..c80d324c3 100644
--- a/docs/configuration/master.md
+++ b/docs/configuration/master.md
@@ -17,31 +17,31 @@ license: |
---
<!--begin-include-->
-| Key | Default | Description | Since |
-| --- | ------- | ----------- | ----- |
-| celeborn.dynamicConfig.refresh.interval | 120s | Interval for refreshing the
corresponding dynamic config periodically. | 0.4.0 |
-| celeborn.dynamicConfig.store.backend | NONE | Store backend for dynamic
config. Available options: NONE, FS. Note: NONE means disabling dynamic config
store. | 0.4.0 |
-| celeborn.master.estimatedPartitionSize.initialSize | 64mb | Initial
partition size for estimation, it will change according to runtime stats. |
0.3.0 |
-| celeborn.master.estimatedPartitionSize.update.initialDelay | 5min | Initial
delay time before start updating partition size for estimation. | 0.3.0 |
-| celeborn.master.estimatedPartitionSize.update.interval | 10min | Interval of
updating partition size for estimation. | 0.3.0 |
-| celeborn.master.hdfs.expireDirs.timeout | 1h | The timeout for a expire dirs
to be deleted on HDFS. | 0.3.0 |
-| celeborn.master.heartbeat.application.timeout | 300s | Application heartbeat
timeout. | 0.3.0 |
-| celeborn.master.heartbeat.worker.timeout | 120s | Worker heartbeat timeout.
| 0.3.0 |
-| celeborn.master.host | <localhost> | Hostname for master to bind. |
0.2.0 |
-| celeborn.master.http.host | <localhost> | Master's http host. | 0.4.0
|
-| celeborn.master.http.port | 9098 | Master's http port. | 0.4.0 |
-| celeborn.master.port | 9097 | Port for master to bind. | 0.2.0 |
-| celeborn.master.slot.assign.extraSlots | 2 | Extra slots number when master
assign slots. | 0.3.0 |
-| celeborn.master.slot.assign.loadAware.diskGroupGradient | 0.1 | This value
means how many more workload will be placed into a faster disk group than a
slower group. | 0.3.0 |
-| celeborn.master.slot.assign.loadAware.fetchTimeWeight | 1.0 | Weight of
average fetch time when calculating ordering in load-aware assignment strategy
| 0.3.0 |
-| celeborn.master.slot.assign.loadAware.flushTimeWeight | 0.0 | Weight of
average flush time when calculating ordering in load-aware assignment strategy
| 0.3.0 |
-| celeborn.master.slot.assign.loadAware.numDiskGroups | 5 | This configuration
is a guidance for load-aware slot allocation algorithm. This value is control
how many disk groups will be created. | 0.3.0 |
-| celeborn.master.slot.assign.maxWorkers | 10000 | Max workers that slots of
one shuffle can be allocated on. Will choose the smaller positive one from
Master side and Client side, see `celeborn.client.slot.assign.maxWorkers`. |
0.3.1 |
-| celeborn.master.slot.assign.policy | ROUNDROBIN | Policy for master to
assign slots, Celeborn supports two types of policy: roundrobin and loadaware.
Loadaware policy will be ignored when `HDFS` is enabled in
`celeborn.storage.activeTypes` | 0.3.0 |
-| celeborn.master.userResourceConsumption.update.interval | 30s | Time length
for a window about compute user resource consumption. | 0.3.0 |
-| celeborn.master.workerUnavailableInfo.expireTimeout | 1800s | Worker
unavailable info would be cleared when the retention period is expired | 0.3.1
|
-| celeborn.storage.availableTypes | HDD | Enabled storages. Available options:
MEMORY,HDD,SSD,HDFS. Note: HDD and SSD would be treated as identical. | 0.3.0 |
-| celeborn.storage.hdfs.dir | <undefined> | HDFS base directory for
Celeborn to store shuffle data. | 0.2.0 |
-| celeborn.storage.hdfs.kerberos.keytab | <undefined> | Kerberos keytab
file path for HDFS storage connection. | 0.3.2 |
-| celeborn.storage.hdfs.kerberos.principal | <undefined> | Kerberos
principal for HDFS storage connection. | 0.3.2 |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.dynamicConfig.refresh.interval | 120s | Interval for refreshing the
corresponding dynamic config periodically. | 0.4.0 | |
+| celeborn.dynamicConfig.store.backend | NONE | Store backend for dynamic
config. Available options: NONE, FS. Note: NONE means disabling dynamic config
store. | 0.4.0 | |
+| celeborn.master.estimatedPartitionSize.initialSize | 64mb | Initial
partition size for estimation, it will change according to runtime stats. |
0.3.0 | celeborn.shuffle.initialEstimatedPartitionSize |
+| celeborn.master.estimatedPartitionSize.update.initialDelay | 5min | Initial
delay time before start updating partition size for estimation. | 0.3.0 |
celeborn.shuffle.estimatedPartitionSize.update.initialDelay |
+| celeborn.master.estimatedPartitionSize.update.interval | 10min | Interval of
updating partition size for estimation. | 0.3.0 |
celeborn.shuffle.estimatedPartitionSize.update.interval |
+| celeborn.master.hdfs.expireDirs.timeout | 1h | The timeout for a expire dirs
to be deleted on HDFS. | 0.3.0 | |
+| celeborn.master.heartbeat.application.timeout | 300s | Application heartbeat
timeout. | 0.3.0 | celeborn.application.heartbeat.timeout |
+| celeborn.master.heartbeat.worker.timeout | 120s | Worker heartbeat timeout.
| 0.3.0 | celeborn.worker.heartbeat.timeout |
+| celeborn.master.host | <localhost> | Hostname for master to bind. |
0.2.0 | |
+| celeborn.master.http.host | <localhost> | Master's http host. | 0.4.0
|
celeborn.metrics.master.prometheus.host,celeborn.master.metrics.prometheus.host
|
+| celeborn.master.http.port | 9098 | Master's http port. | 0.4.0 |
celeborn.metrics.master.prometheus.port,celeborn.master.metrics.prometheus.port
|
+| celeborn.master.port | 9097 | Port for master to bind. | 0.2.0 | |
+| celeborn.master.slot.assign.extraSlots | 2 | Extra slots number when master
assign slots. | 0.3.0 | celeborn.slots.assign.extraSlots |
+| celeborn.master.slot.assign.loadAware.diskGroupGradient | 0.1 | This value
means how many more workload will be placed into a faster disk group than a
slower group. | 0.3.0 | celeborn.slots.assign.loadAware.diskGroupGradient |
+| celeborn.master.slot.assign.loadAware.fetchTimeWeight | 1.0 | Weight of
average fetch time when calculating ordering in load-aware assignment strategy
| 0.3.0 | celeborn.slots.assign.loadAware.fetchTimeWeight |
+| celeborn.master.slot.assign.loadAware.flushTimeWeight | 0.0 | Weight of
average flush time when calculating ordering in load-aware assignment strategy
| 0.3.0 | celeborn.slots.assign.loadAware.flushTimeWeight |
+| celeborn.master.slot.assign.loadAware.numDiskGroups | 5 | This configuration
is a guidance for load-aware slot allocation algorithm. This value is control
how many disk groups will be created. | 0.3.0 |
celeborn.slots.assign.loadAware.numDiskGroups |
+| celeborn.master.slot.assign.maxWorkers | 10000 | Max workers that slots of
one shuffle can be allocated on. Will choose the smaller positive one from
Master side and Client side, see `celeborn.client.slot.assign.maxWorkers`. |
0.3.1 | |
+| celeborn.master.slot.assign.policy | ROUNDROBIN | Policy for master to
assign slots, Celeborn supports two types of policy: roundrobin and loadaware.
Loadaware policy will be ignored when `HDFS` is enabled in
`celeborn.storage.activeTypes` | 0.3.0 | celeborn.slots.assign.policy |
+| celeborn.master.userResourceConsumption.update.interval | 30s | Time length
for a window about compute user resource consumption. | 0.3.0 | |
+| celeborn.master.workerUnavailableInfo.expireTimeout | 1800s | Worker
unavailable info would be cleared when the retention period is expired | 0.3.1
| |
+| celeborn.storage.availableTypes | HDD | Enabled storages. Available options:
MEMORY,HDD,SSD,HDFS. Note: HDD and SSD would be treated as identical. | 0.3.0 |
celeborn.storage.activeTypes |
+| celeborn.storage.hdfs.dir | <undefined> | HDFS base directory for
Celeborn to store shuffle data. | 0.2.0 | |
+| celeborn.storage.hdfs.kerberos.keytab | <undefined> | Kerberos keytab
file path for HDFS storage connection. | 0.3.2 | |
+| celeborn.storage.hdfs.kerberos.principal | <undefined> | Kerberos
principal for HDFS storage connection. | 0.3.2 | |
<!--end-include-->
diff --git a/docs/configuration/metrics.md b/docs/configuration/metrics.md
index fd1beadfa..1ba236f33 100644
--- a/docs/configuration/metrics.md
+++ b/docs/configuration/metrics.md
@@ -17,20 +17,20 @@ license: |
---
<!--begin-include-->
-| Key | Default | Description | Since |
-| --- | ------- | ----------- | ----- |
-| celeborn.metrics.app.topDiskUsage.count | 50 | Size for top items about top
disk usage applications list. | 0.2.0 |
-| celeborn.metrics.app.topDiskUsage.interval | 10min | Time length for a
window about top disk usage application list. | 0.2.0 |
-| celeborn.metrics.app.topDiskUsage.windowSize | 24 | Window size about top
disk usage application list. | 0.2.0 |
-| celeborn.metrics.capacity | 4096 | The maximum number of metrics which a
source can use to generate output strings. | 0.2.0 |
-| celeborn.metrics.collectPerfCritical.enabled | false | It controls whether
to collect metrics which may affect performance. When enable, Celeborn collects
them. | 0.2.0 |
-| celeborn.metrics.conf | <undefined> | Custom metrics configuration
file path. Default use `metrics.properties` in classpath. | 0.3.0 |
-| celeborn.metrics.enabled | true | When true, enable metrics system. | 0.2.0
|
-| celeborn.metrics.extraLabels | | If default metric labels are not enough,
extra metric labels can be customized. Labels' pattern is:
`<label1_key>=<label1_value>[,<label2_key>=<label2_value>]*`; e.g.
`env=prod,version=1` | 0.3.0 |
-| celeborn.metrics.json.path | /metrics/json | URI context path of json
metrics HTTP server. | 0.4.0 |
-| celeborn.metrics.json.pretty.enabled | true | When true, view metrics in
json pretty format | 0.4.0 |
-| celeborn.metrics.prometheus.path | /metrics/prometheus | URI context path of
prometheus metrics HTTP server. | 0.4.0 |
-| celeborn.metrics.sample.rate | 1.0 | It controls if Celeborn collect timer
metrics for some operations. Its value should be in [0.0, 1.0]. | 0.2.0 |
-| celeborn.metrics.timer.slidingWindow.size | 4096 | The sliding window size
of timer metric. | 0.2.0 |
-| celeborn.metrics.worker.pauseSpentTime.forceAppend.threshold | 10 | Force
append worker pause spent time even if worker still in pause serving state.Help
user can find worker pause spent time increase, when worker always been pause
state. | |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.metrics.app.topDiskUsage.count | 50 | Size for top items about top
disk usage applications list. | 0.2.0 | |
+| celeborn.metrics.app.topDiskUsage.interval | 10min | Time length for a
window about top disk usage application list. | 0.2.0 | |
+| celeborn.metrics.app.topDiskUsage.windowSize | 24 | Window size about top
disk usage application list. | 0.2.0 | |
+| celeborn.metrics.capacity | 4096 | The maximum number of metrics which a
source can use to generate output strings. | 0.2.0 | |
+| celeborn.metrics.collectPerfCritical.enabled | false | It controls whether
to collect metrics which may affect performance. When enable, Celeborn collects
them. | 0.2.0 | |
+| celeborn.metrics.conf | <undefined> | Custom metrics configuration
file path. Default use `metrics.properties` in classpath. | 0.3.0 | |
+| celeborn.metrics.enabled | true | When true, enable metrics system. | 0.2.0
| |
+| celeborn.metrics.extraLabels | | If default metric labels are not enough,
extra metric labels can be customized. Labels' pattern is:
`<label1_key>=<label1_value>[,<label2_key>=<label2_value>]*`; e.g.
`env=prod,version=1` | 0.3.0 | |
+| celeborn.metrics.json.path | /metrics/json | URI context path of json
metrics HTTP server. | 0.4.0 | |
+| celeborn.metrics.json.pretty.enabled | true | When true, view metrics in
json pretty format | 0.4.0 | |
+| celeborn.metrics.prometheus.path | /metrics/prometheus | URI context path of
prometheus metrics HTTP server. | 0.4.0 | |
+| celeborn.metrics.sample.rate | 1.0 | It controls if Celeborn collect timer
metrics for some operations. Its value should be in [0.0, 1.0]. | 0.2.0 | |
+| celeborn.metrics.timer.slidingWindow.size | 4096 | The sliding window size
of timer metric. | 0.2.0 | |
+| celeborn.metrics.worker.pauseSpentTime.forceAppend.threshold | 10 | Force
append worker pause spent time even if worker still in pause serving state.Help
user can find worker pause spent time increase, when worker always been pause
state. | | |
<!--end-include-->
diff --git a/docs/configuration/network.md b/docs/configuration/network.md
index c256a2f9f..f13f29b2a 100644
--- a/docs/configuration/network.md
+++ b/docs/configuration/network.md
@@ -17,40 +17,40 @@ 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. | 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. 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. 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).
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. If setting <module> to `replicate`, it works for replicate
client of worker replicating data to peer worker. | 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. If setting <module> to `replicate`, it works for replicate
client of worker replicating data to peer worker. | 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 |
-| celeborn.network.connect.timeout | 10s | Default socket connect timeout. |
0.2.0 |
-| celeborn.network.memory.allocator.numArenas | <undefined> | Number of
arenas for pooled memory allocator. Default value is
Runtime.getRuntime.availableProcessors, min value is 2. | 0.3.0 |
-| celeborn.network.memory.allocator.verbose.metric | false | Weather to enable
verbose metric for pooled allocator. | 0.3.0 |
-| celeborn.network.timeout | 240s | Default timeout for network operations. |
0.2.0 |
-| celeborn.port.maxRetries | 1 | When port is occupied, we will retry for max
retry times. | 0.2.0 |
-| celeborn.rpc.askTimeout | 60s | Timeout for RPC ask operations. It's
recommended to set at least `240s` when `HDFS` is enabled in
`celeborn.storage.activeTypes` | 0.2.0 |
-| celeborn.rpc.connect.threads | 64 | | 0.2.0 |
-| celeborn.rpc.dispatcher.threads | 0 | Threads number of message dispatcher
event loop. Default to 0, which is availableCore. | 0.3.0 |
-| celeborn.rpc.io.threads | <undefined> | Netty IO thread number of
NettyRpcEnv to handle RPC request. The default threads number is the number of
runtime available processors. | 0.2.0 |
-| celeborn.rpc.lookupTimeout | 30s | Timeout for RPC lookup operations. |
0.2.0 |
-| celeborn.shuffle.io.maxChunksBeingTransferred | <undefined> | The max
number of chunks allowed to be transferred at the same time on shuffle service.
Note that new incoming connections will be closed when the max number is hit.
The client will retry according to the shuffle retry configs (see
`celeborn.<module>.io.maxRetries` and `celeborn.<module>.io.retryWait`), if
those limits are reached the task will fail with fetch failure. | 0.2.0 |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| 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.If you are using the
"celeborn.client.heartbeat.interval", please use the new configs for each
module according to your needs or replace it wi [...]
+| 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. 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. 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).
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. If setting <module> to `replicate`, it works for replicate
client of worker replicating data to peer worker. | 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. If setting <module> to `replicate`, it works for replicate
client of worker replicating data to peer worker. | 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 | |
+| celeborn.network.connect.timeout | 10s | Default socket connect timeout. |
0.2.0 | |
+| celeborn.network.memory.allocator.numArenas | <undefined> | Number of
arenas for pooled memory allocator. Default value is
Runtime.getRuntime.availableProcessors, min value is 2. | 0.3.0 | |
+| celeborn.network.memory.allocator.verbose.metric | false | Whether to enable
verbose metric for pooled allocator. | 0.3.0 | |
+| celeborn.network.timeout | 240s | Default timeout for network operations. |
0.2.0 | |
+| celeborn.port.maxRetries | 1 | When port is occupied, we will retry for max
retry times. | 0.2.0 | |
+| celeborn.rpc.askTimeout | 60s | Timeout for RPC ask operations. It's
recommended to set at least `240s` when `HDFS` is enabled in
`celeborn.storage.activeTypes` | 0.2.0 | |
+| celeborn.rpc.connect.threads | 64 | | 0.2.0 | |
+| celeborn.rpc.dispatcher.threads | 0 | Threads number of message dispatcher
event loop. Default to 0, which is availableCore. | 0.3.0 |
celeborn.rpc.dispatcher.numThreads |
+| celeborn.rpc.io.threads | <undefined> | Netty IO thread number of
NettyRpcEnv to handle RPC request. The default threads number is the number of
runtime available processors. | 0.2.0 | |
+| celeborn.rpc.lookupTimeout | 30s | Timeout for RPC lookup operations. |
0.2.0 | |
+| celeborn.shuffle.io.maxChunksBeingTransferred | <undefined> | The max
number of chunks allowed to be transferred at the same time on shuffle service.
Note that new incoming connections will be closed when the max number is hit.
The client will retry according to the shuffle retry configs (see
`celeborn.<module>.io.maxRetries` and `celeborn.<module>.io.retryWait`), if
those limits are reached the task will fail with fetch failure. | 0.2.0 | |
<!--end-include-->
diff --git a/docs/configuration/quota.md b/docs/configuration/quota.md
index 565dde953..33fb79d81 100644
--- a/docs/configuration/quota.md
+++ b/docs/configuration/quota.md
@@ -17,12 +17,12 @@ license: |
---
<!--begin-include-->
-| Key | Default | Description | Since |
-| --- | ------- | ----------- | ----- |
-| celeborn.quota.configuration.path | <undefined> | Quota configuration
file path. The file format should be yaml. Quota configuration file template
can be found under conf directory. | 0.2.0 |
-| celeborn.quota.enabled | true | When true, before registering shuffle,
LifecycleManager should check if current user have enough quota space, if
cluster don't have enough quota space for current user, fallback to Spark's
default shuffle | 0.2.0 |
-| celeborn.quota.identity.provider |
org.apache.celeborn.common.identity.DefaultIdentityProvider | IdentityProvider
class name. Default class is
`org.apache.celeborn.common.identity.DefaultIdentityProvider`. Optional values:
org.apache.celeborn.common.identity.HadoopBasedIdentityProvider user name will
be obtained by UserGroupInformation.getUserName;
org.apache.celeborn.common.identity.DefaultIdentityProvider user name and
tenant id are default values or user-specific values. | 0.2.0 |
-| celeborn.quota.identity.user-specific.tenant | default | Tenant id if
celeborn.quota.identity.provider is
org.apache.celeborn.common.identity.DefaultIdentityProvider. | 0.3.0 |
-| celeborn.quota.identity.user-specific.userName | default | User name if
celeborn.quota.identity.provider is
org.apache.celeborn.common.identity.DefaultIdentityProvider. | 0.3.0 |
-| celeborn.quota.manager |
org.apache.celeborn.common.quota.DefaultQuotaManager | QuotaManger class name.
Default class is `org.apache.celeborn.common.quota.DefaultQuotaManager`. |
0.2.0 |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.quota.configuration.path | <undefined> | Quota configuration
file path. The file format should be yaml. Quota configuration file template
can be found under conf directory. | 0.2.0 | |
+| celeborn.quota.enabled | true | When true, before registering shuffle,
LifecycleManager should check if current user have enough quota space, if
cluster don't have enough quota space for current user, fallback to Spark's
default shuffle | 0.2.0 | |
+| celeborn.quota.identity.provider |
org.apache.celeborn.common.identity.DefaultIdentityProvider | IdentityProvider
class name. Default class is
`org.apache.celeborn.common.identity.DefaultIdentityProvider`. Optional values:
org.apache.celeborn.common.identity.HadoopBasedIdentityProvider user name will
be obtained by UserGroupInformation.getUserName;
org.apache.celeborn.common.identity.DefaultIdentityProvider user name and
tenant id are default values or user-specific values. | 0.2.0 | |
+| celeborn.quota.identity.user-specific.tenant | default | Tenant id if
celeborn.quota.identity.provider is
org.apache.celeborn.common.identity.DefaultIdentityProvider. | 0.3.0 | |
+| celeborn.quota.identity.user-specific.userName | default | User name if
celeborn.quota.identity.provider is
org.apache.celeborn.common.identity.DefaultIdentityProvider. | 0.3.0 | |
+| celeborn.quota.manager |
org.apache.celeborn.common.quota.DefaultQuotaManager | QuotaManger class name.
Default class is `org.apache.celeborn.common.quota.DefaultQuotaManager`. |
0.2.0 | |
<!--end-include-->
diff --git a/docs/configuration/worker.md b/docs/configuration/worker.md
index c29c44305..24340be0b 100644
--- a/docs/configuration/worker.md
+++ b/docs/configuration/worker.md
@@ -17,112 +17,112 @@ license: |
---
<!--begin-include-->
-| Key | Default | Description | Since |
-| --- | ------- | ----------- | ----- |
-| celeborn.dynamicConfig.refresh.interval | 120s | Interval for refreshing the
corresponding dynamic config periodically. | 0.4.0 |
-| celeborn.dynamicConfig.store.backend | NONE | Store backend for dynamic
config. Available options: NONE, FS. Note: NONE means disabling dynamic config
store. | 0.4.0 |
-| celeborn.master.endpoints | <localhost>:9097 | Endpoints of master
nodes for celeborn client to connect, allowed pattern is:
`<host1>:<port1>[,<host2>:<port2>]*`, e.g. `clb1:9097,clb2:9098,clb3:9099`. If
the port is omitted, 9097 will be used. | 0.2.0 |
-| celeborn.master.estimatedPartitionSize.minSize | 8mb | Ignore partition size
smaller than this configuration of partition size for estimation. | 0.3.0 |
-| celeborn.shuffle.chunk.size | 8m | Max chunk size of reducer's merged
shuffle data. For example, if a reducer's shuffle data is 128M and the data
will need 16 fetch chunk requests to fetch. | 0.2.0 |
-| celeborn.storage.availableTypes | HDD | Enabled storages. Available options:
MEMORY,HDD,SSD,HDFS. Note: HDD and SSD would be treated as identical. | 0.3.0 |
-| celeborn.storage.hdfs.dir | <undefined> | HDFS base directory for
Celeborn to store shuffle data. | 0.2.0 |
-| celeborn.storage.hdfs.kerberos.keytab | <undefined> | Kerberos keytab
file path for HDFS storage connection. | 0.3.2 |
-| celeborn.storage.hdfs.kerberos.principal | <undefined> | Kerberos
principal for HDFS storage connection. | 0.3.2 |
-| celeborn.worker.activeConnection.max | <undefined> | If the number of
active connections on a worker exceeds this configuration value, the worker
will be marked as high-load in the heartbeat report, and the master will not
include that node in the response of RequestSlots. | 0.3.1 |
-| celeborn.worker.bufferStream.threadsPerMountpoint | 8 | Threads count for
read buffer per mount point. | 0.3.0 |
-| celeborn.worker.clean.threads | 64 | Thread number of worker to clean up
expired shuffle keys. | 0.3.2 |
-| celeborn.worker.closeIdleConnections | false | Whether worker will close
idle connections. | 0.2.0 |
-| celeborn.worker.commitFiles.threads | 32 | Thread number of worker to commit
shuffle data files asynchronously. It's recommended to set at least `128` when
`HDFS` is enabled in `celeborn.storage.activeTypes`. | 0.3.0 |
-| celeborn.worker.commitFiles.timeout | 120s | Timeout for a Celeborn worker
to commit files of a shuffle. It's recommended to set at least `240s` when
`HDFS` is enabled in `celeborn.storage.activeTypes`. | 0.3.0 |
-| celeborn.worker.congestionControl.check.interval | 10ms | Interval of worker
checks congestion if celeborn.worker.congestionControl.enabled is true. | 0.3.2
|
-| celeborn.worker.congestionControl.enabled | false | Whether to enable
congestion control or not. | 0.3.0 |
-| celeborn.worker.congestionControl.high.watermark | <undefined> | If
the total bytes in disk buffer exceeds this configure, will start to
congestusers whose produce rate is higher than the potential average consume
rate. The congestion will stop if the produce rate is lower or equal to the
average consume rate, or the total pending bytes lower than
celeborn.worker.congestionControl.low.watermark | 0.3.0 |
-| celeborn.worker.congestionControl.low.watermark | <undefined> | Will
stop congest users if the total pending bytes of disk buffer is lower than this
configuration | 0.3.0 |
-| celeborn.worker.congestionControl.sample.time.window | 10s | The worker
holds a time sliding list to calculate users' produce/consume rate | 0.3.0 |
-| celeborn.worker.congestionControl.user.inactive.interval | 10min | How long
will consider this user is inactive if it doesn't send data | 0.3.0 |
-| celeborn.worker.decommission.checkInterval | 30s | The wait interval of
checking whether all the shuffle expired during worker decommission | 0.4.0 |
-| celeborn.worker.decommission.forceExitTimeout | 6h | The wait time of
waiting for all the shuffle expire during worker decommission. | 0.4.0 |
-| celeborn.worker.directMemoryRatioForMemoryShuffleStorage | 0.0 | Max ratio
of direct memory to store shuffle data | 0.2.0 |
-| celeborn.worker.directMemoryRatioForReadBuffer | 0.1 | Max ratio of direct
memory for read buffer | 0.2.0 |
-| celeborn.worker.directMemoryRatioToPauseReceive | 0.85 | If direct memory
usage reaches this limit, the worker will stop to receive data from Celeborn
shuffle clients. | 0.2.0 |
-| celeborn.worker.directMemoryRatioToPauseReplicate | 0.95 | If direct memory
usage reaches this limit, the worker will stop to receive replication data from
other workers. This value should be higher than
celeborn.worker.directMemoryRatioToPauseReceive. | 0.2.0 |
-| celeborn.worker.directMemoryRatioToResume | 0.7 | If direct memory usage is
less than this limit, worker will resume. | 0.2.0 |
-| celeborn.worker.disk.clean.threads | 4 | Thread number of worker to clean up
directories of expired shuffle keys on disk. | 0.3.2 |
-| celeborn.worker.fetch.heartbeat.enabled | false | enable the heartbeat from
worker to client when fetching data | 0.3.0 |
-| celeborn.worker.fetch.io.threads | <undefined> | Netty IO thread
number of worker to handle client fetch data. The default threads number is the
number of flush thread. | 0.2.0 |
-| celeborn.worker.fetch.port | 0 | Server port for Worker to receive fetch
data request from ShuffleClient. | 0.2.0 |
-| celeborn.worker.flusher.buffer.size | 256k | Size of buffer used by a single
flusher. | 0.2.0 |
-| celeborn.worker.flusher.diskTime.slidingWindow.size | 20 | The size of
sliding windows used to calculate statistics about flushed time and count. |
0.3.0 |
-| celeborn.worker.flusher.hdd.threads | 1 | Flusher's thread count per disk
used for write data to HDD disks. | 0.2.0 |
-| celeborn.worker.flusher.hdfs.buffer.size | 4m | Size of buffer used by a
HDFS flusher. | 0.3.0 |
-| celeborn.worker.flusher.hdfs.threads | 8 | Flusher's thread count used for
write data to HDFS. | 0.2.0 |
-| celeborn.worker.flusher.shutdownTimeout | 3s | Timeout for a flusher to
shutdown. | 0.2.0 |
-| celeborn.worker.flusher.ssd.threads | 16 | Flusher's thread count per disk
used for write data to SSD disks. | 0.2.0 |
-| celeborn.worker.flusher.threads | 16 | Flusher's thread count per disk for
unknown-type disks. | 0.2.0 |
-| celeborn.worker.graceful.shutdown.checkSlotsFinished.interval | 1s | The
wait interval of checking whether all released slots to be committed or
destroyed during worker graceful shutdown | 0.2.0 |
-| celeborn.worker.graceful.shutdown.checkSlotsFinished.timeout | 480s | The
wait time of waiting for the released slots to be committed or destroyed during
worker graceful shutdown. | 0.2.0 |
-| celeborn.worker.graceful.shutdown.enabled | false | When true, during worker
shutdown, the worker will wait for all released slots to be committed or
destroyed. | 0.2.0 |
-| celeborn.worker.graceful.shutdown.partitionSorter.shutdownTimeout | 120s |
The wait time of waiting for sorting partition files during worker graceful
shutdown. | 0.2.0 |
-| celeborn.worker.graceful.shutdown.recoverDbBackend | LEVELDB | Specifies a
disk-based store used in local db. LEVELDB or ROCKSDB. | 0.4.0 |
-| celeborn.worker.graceful.shutdown.recoverPath | <tmp>/recover | The
path to store DB. | 0.2.0 |
-| celeborn.worker.graceful.shutdown.saveCommittedFileInfo.interval | 5s |
Interval for a Celeborn worker to flush committed file infos into Level DB. |
0.3.1 |
-| celeborn.worker.graceful.shutdown.saveCommittedFileInfo.sync | false |
Whether to call sync method to save committed file infos into Level DB to
handle OS crash. | 0.3.1 |
-| celeborn.worker.graceful.shutdown.timeout | 600s | The worker's graceful
shutdown timeout time. | 0.2.0 |
-| celeborn.worker.http.host | <localhost> | Worker's http host. | 0.4.0
|
-| celeborn.worker.http.port | 9096 | Worker's http port. | 0.4.0 |
-| celeborn.worker.jvmQuake.check.interval | 1s | Interval of gc behavior
checking for worker jvm quake. | 0.4.0 |
-| celeborn.worker.jvmQuake.dump.enabled | true | Whether to heap dump for the
maximum GC 'deficit' during worker jvm quake. | 0.4.0 |
-| celeborn.worker.jvmQuake.dump.path | <tmp>/jvm-quake/dump/<pid>
| The path of heap dump for the maximum GC 'deficit' during worker jvm quake. |
0.4.0 |
-| celeborn.worker.jvmQuake.dump.threshold | 30s | The threshold of heap dump
for the maximum GC 'deficit' which can be accumulated before jvmquake takes
action. Meanwhile, there is no heap dump generated when dump threshold is
greater than kill threshold. | 0.4.0 |
-| celeborn.worker.jvmQuake.enabled | false | When true, Celeborn worker will
start the jvm quake to monitor of gc behavior, which enables early detection of
memory management issues and facilitates fast failure. | 0.4.0 |
-| celeborn.worker.jvmQuake.exitCode | 502 | The exit code of system kill for
the maximum GC 'deficit' during worker jvm quake. | 0.4.0 |
-| celeborn.worker.jvmQuake.kill.threshold | 60s | The threshold of system kill
for the maximum GC 'deficit' which can be accumulated before jvmquake takes
action. | 0.4.0 |
-| celeborn.worker.jvmQuake.runtimeWeight | 5.0 | The factor by which to
multiply running JVM time, when weighing it against GCing time. 'Deficit' is
accumulated as `gc_time - runtime * runtime_weight`, and is compared against
threshold to determine whether to take action. | 0.4.0 |
-| celeborn.worker.monitor.disk.check.interval | 30s | Intervals between device
monitor to check disk. | 0.3.0 |
-| celeborn.worker.monitor.disk.check.timeout | 30s | Timeout time for worker
check device status. | 0.3.0 |
-| celeborn.worker.monitor.disk.checklist | readwrite,diskusage | Monitor type
for disk, available items are: iohang, readwrite and diskusage. | 0.2.0 |
-| celeborn.worker.monitor.disk.enabled | true | When true, worker will monitor
device and report to master. | 0.3.0 |
-| celeborn.worker.monitor.disk.notifyError.expireTimeout | 10m | The expire
timeout of non-critical device error. Only notify critical error when the
number of non-critical errors for a period of time exceeds threshold. | 0.3.0 |
-| celeborn.worker.monitor.disk.notifyError.threshold | 64 | Device monitor
will only notify critical error once the accumulated valid non-critical error
number exceeding this threshold. | 0.3.0 |
-| celeborn.worker.monitor.disk.sys.block.dir | /sys/block | The directory
where linux file block information is stored. | 0.2.0 |
-| celeborn.worker.monitor.memory.check.interval | 10ms | Interval of worker
direct memory checking. | 0.3.0 |
-| celeborn.worker.monitor.memory.report.interval | 10s | Interval of worker
direct memory tracker reporting to log. | 0.3.0 |
-| celeborn.worker.monitor.memory.trimChannelWaitInterval | 1s | Wait time
after worker trigger channel to trim cache. | 0.3.0 |
-| celeborn.worker.monitor.memory.trimFlushWaitInterval | 1s | Wait time after
worker trigger StorageManger to flush data. | 0.3.0 |
-| celeborn.worker.partition.initial.readBuffersMax | 1024 | Max number of
initial read buffers | 0.3.0 |
-| celeborn.worker.partition.initial.readBuffersMin | 1 | Min number of initial
read buffers | 0.3.0 |
-| celeborn.worker.partitionSorter.directMemoryRatioThreshold | 0.1 | Max ratio
of partition sorter's memory for sorting, when reserved memory is higher than
max partition sorter memory, partition sorter will stop sorting. | 0.2.0 |
-| celeborn.worker.push.heartbeat.enabled | false | enable the heartbeat from
worker to client when pushing data | 0.3.0 |
-| celeborn.worker.push.io.threads | <undefined> | Netty IO thread number
of worker to handle client push data. The default threads number is the number
of flush thread. | 0.2.0 |
-| celeborn.worker.push.port | 0 | Server port for Worker to receive push data
request from ShuffleClient. | 0.2.0 |
-| celeborn.worker.readBuffer.allocationWait | 50ms | The time to wait when
buffer dispatcher can not allocate a buffer. | 0.3.0 |
-| celeborn.worker.readBuffer.target.changeThreshold | 1mb | The target ratio
for pre read memory usage. | 0.3.0 |
-| celeborn.worker.readBuffer.target.ratio | 0.9 | The target ratio for read
ahead buffer's memory usage. | 0.3.0 |
-| celeborn.worker.readBuffer.target.updateInterval | 100ms | The interval for
memory manager to calculate new read buffer's target memory. | 0.3.0 |
-| celeborn.worker.readBuffer.toTriggerReadMin | 32 | Min buffers count for map
data partition to trigger read. | 0.3.0 |
-| celeborn.worker.register.timeout | 180s | Worker register timeout. | 0.2.0 |
-| celeborn.worker.replicate.fastFail.duration | 60s | If a replicate request
not replied during the duration, worker will mark the replicate data request as
failed.It's recommended to set at least `240s` when `HDFS` is enabled in
`celeborn.storage.activeTypes`. | 0.2.0 |
-| celeborn.worker.replicate.io.threads | <undefined> | Netty IO thread
number of worker to replicate shuffle data. The default threads number is the
number of flush thread. | 0.2.0 |
-| celeborn.worker.replicate.port | 0 | Server port for Worker to receive
replicate data request from other Workers. | 0.2.0 |
-| celeborn.worker.replicate.randomConnection.enabled | true | Whether worker
will create random connection to peer when replicate data. When false, worker
tend to reuse the same cached TransportClient to a specific replicate worker;
when true, worker tend to use different cached TransportClient. Netty will use
the same thread to serve the same connection, so with more connections
replicate server can leverage more netty threads | 0.2.1 |
-| celeborn.worker.replicate.threads | 64 | Thread number of worker to
replicate shuffle data. | 0.2.0 |
-| celeborn.worker.rpc.port | 0 | Server port for Worker to receive RPC
request. | 0.2.0 |
-| celeborn.worker.shuffle.partitionSplit.enabled | true | enable the partition
split on worker side | 0.3.0 |
-| celeborn.worker.shuffle.partitionSplit.max | 2g | Specify the maximum
partition size for splitting, and ensure that individual partition files are
always smaller than this limit. | 0.3.0 |
-| celeborn.worker.shuffle.partitionSplit.min | 1m | Min size for a partition
to split | 0.3.0 |
-| celeborn.worker.sortPartition.indexCache.expire | 180s | PartitionSorter's
cache item expire time. | 0.4.0 |
-| celeborn.worker.sortPartition.indexCache.maxWeight | 100000 |
PartitionSorter's cache max weight for index buffer. | 0.4.0 |
-| celeborn.worker.sortPartition.reservedMemoryPerPartition | 1mb | Reserved
memory when sorting a shuffle file off-heap. | 0.3.0 |
-| celeborn.worker.sortPartition.threads | <undefined> |
PartitionSorter's thread counts. It's recommended to set at least `64` when
`HDFS` is enabled in `celeborn.storage.activeTypes`. | 0.3.0 |
-| celeborn.worker.sortPartition.timeout | 220s | Timeout for a shuffle file to
sort. | 0.3.0 |
-| celeborn.worker.storage.checkDirsEmpty.maxRetries | 3 | The number of
retries for a worker to check if the working directory is cleaned up before
registering with the master. | 0.3.0 |
-| celeborn.worker.storage.checkDirsEmpty.timeout | 1000ms | The wait time per
retry for a worker to check if the working directory is cleaned up before
registering with the master. | 0.3.0 |
-| celeborn.worker.storage.dirs | <undefined> | Directory list to store
shuffle data. It's recommended to configure one directory on each disk. Storage
size limit can be set for each directory. For the sake of performance, there
should be no more than 2 flush threads on the same disk partition if you are
using HDD, and should be 8 or more flush threads on the same disk partition if
you are using SSD. For example:
`dir1[:capacity=][:disktype=][:flushthread=],dir2[:capacity=][:disktyp [...]
-| celeborn.worker.storage.disk.reserve.ratio | <undefined> | Celeborn
worker reserved ratio for each disk. The minimum usable size for each disk is
the max space between the reserved space and the space calculate via reserved
ratio. | 0.3.2 |
-| celeborn.worker.storage.disk.reserve.size | 5G | Celeborn worker reserved
space for each disk. | 0.3.0 |
-| celeborn.worker.storage.expireDirs.timeout | 1h | The timeout for a expire
dirs to be deleted on disk. | 0.3.2 |
-| celeborn.worker.storage.workingDir | celeborn-worker/shuffle_data | Worker's
working dir path name. | 0.3.0 |
-| celeborn.worker.userResourceConsumption.update.interval | 30s | Time length
for a window about compute user resource consumption. | 0.3.2 |
-| celeborn.worker.writer.close.timeout | 120s | Timeout for a file writer to
close | 0.2.0 |
-| celeborn.worker.writer.create.maxAttempts | 3 | Retry count for a file
writer to create if its creation was failed. | 0.2.0 |
+| Key | Default | Description | Since | Deprecated |
+| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.dynamicConfig.refresh.interval | 120s | Interval for refreshing the
corresponding dynamic config periodically. | 0.4.0 | |
+| celeborn.dynamicConfig.store.backend | NONE | Store backend for dynamic
config. Available options: NONE, FS. Note: NONE means disabling dynamic config
store. | 0.4.0 | |
+| celeborn.master.endpoints | <localhost>:9097 | Endpoints of master
nodes for celeborn client to connect, allowed pattern is:
`<host1>:<port1>[,<host2>:<port2>]*`, e.g. `clb1:9097,clb2:9098,clb3:9099`. If
the port is omitted, 9097 will be used. | 0.2.0 | |
+| celeborn.master.estimatedPartitionSize.minSize | 8mb | Ignore partition size
smaller than this configuration of partition size for estimation. | 0.3.0 |
celeborn.shuffle.minPartitionSizeToEstimate |
+| celeborn.shuffle.chunk.size | 8m | Max chunk size of reducer's merged
shuffle data. For example, if a reducer's shuffle data is 128M and the data
will need 16 fetch chunk requests to fetch. | 0.2.0 | |
+| celeborn.storage.availableTypes | HDD | Enabled storages. Available options:
MEMORY,HDD,SSD,HDFS. Note: HDD and SSD would be treated as identical. | 0.3.0 |
celeborn.storage.activeTypes |
+| celeborn.storage.hdfs.dir | <undefined> | HDFS base directory for
Celeborn to store shuffle data. | 0.2.0 | |
+| celeborn.storage.hdfs.kerberos.keytab | <undefined> | Kerberos keytab
file path for HDFS storage connection. | 0.3.2 | |
+| celeborn.storage.hdfs.kerberos.principal | <undefined> | Kerberos
principal for HDFS storage connection. | 0.3.2 | |
+| celeborn.worker.activeConnection.max | <undefined> | If the number of
active connections on a worker exceeds this configuration value, the worker
will be marked as high-load in the heartbeat report, and the master will not
include that node in the response of RequestSlots. | 0.3.1 | |
+| celeborn.worker.bufferStream.threadsPerMountpoint | 8 | Threads count for
read buffer per mount point. | 0.3.0 | |
+| celeborn.worker.clean.threads | 64 | Thread number of worker to clean up
expired shuffle keys. | 0.3.2 | |
+| celeborn.worker.closeIdleConnections | false | Whether worker will close
idle connections. | 0.2.0 | |
+| celeborn.worker.commitFiles.threads | 32 | Thread number of worker to commit
shuffle data files asynchronously. It's recommended to set at least `128` when
`HDFS` is enabled in `celeborn.storage.activeTypes`. | 0.3.0 |
celeborn.worker.commit.threads |
+| celeborn.worker.commitFiles.timeout | 120s | Timeout for a Celeborn worker
to commit files of a shuffle. It's recommended to set at least `240s` when
`HDFS` is enabled in `celeborn.storage.activeTypes`. | 0.3.0 |
celeborn.worker.shuffle.commit.timeout |
+| celeborn.worker.congestionControl.check.interval | 10ms | Interval of worker
checks congestion if celeborn.worker.congestionControl.enabled is true. | 0.3.2
| |
+| celeborn.worker.congestionControl.enabled | false | Whether to enable
congestion control or not. | 0.3.0 | |
+| celeborn.worker.congestionControl.high.watermark | <undefined> | If
the total bytes in disk buffer exceeds this configure, will start to
congestusers whose produce rate is higher than the potential average consume
rate. The congestion will stop if the produce rate is lower or equal to the
average consume rate, or the total pending bytes lower than
celeborn.worker.congestionControl.low.watermark | 0.3.0 | |
+| celeborn.worker.congestionControl.low.watermark | <undefined> | Will
stop congest users if the total pending bytes of disk buffer is lower than this
configuration | 0.3.0 | |
+| celeborn.worker.congestionControl.sample.time.window | 10s | The worker
holds a time sliding list to calculate users' produce/consume rate | 0.3.0 | |
+| celeborn.worker.congestionControl.user.inactive.interval | 10min | How long
will consider this user is inactive if it doesn't send data | 0.3.0 | |
+| celeborn.worker.decommission.checkInterval | 30s | The wait interval of
checking whether all the shuffle expired during worker decommission | 0.4.0 |
|
+| celeborn.worker.decommission.forceExitTimeout | 6h | The wait time of
waiting for all the shuffle expire during worker decommission. | 0.4.0 | |
+| celeborn.worker.directMemoryRatioForMemoryShuffleStorage | 0.0 | Max ratio
of direct memory to store shuffle data | 0.2.0 | |
+| celeborn.worker.directMemoryRatioForReadBuffer | 0.1 | Max ratio of direct
memory for read buffer | 0.2.0 | |
+| celeborn.worker.directMemoryRatioToPauseReceive | 0.85 | If direct memory
usage reaches this limit, the worker will stop to receive data from Celeborn
shuffle clients. | 0.2.0 | |
+| celeborn.worker.directMemoryRatioToPauseReplicate | 0.95 | If direct memory
usage reaches this limit, the worker will stop to receive replication data from
other workers. This value should be higher than
celeborn.worker.directMemoryRatioToPauseReceive. | 0.2.0 | |
+| celeborn.worker.directMemoryRatioToResume | 0.7 | If direct memory usage is
less than this limit, worker will resume. | 0.2.0 | |
+| celeborn.worker.disk.clean.threads | 4 | Thread number of worker to clean up
directories of expired shuffle keys on disk. | 0.3.2 | |
+| celeborn.worker.fetch.heartbeat.enabled | false | enable the heartbeat from
worker to client when fetching data | 0.3.0 | |
+| celeborn.worker.fetch.io.threads | <undefined> | Netty IO thread
number of worker to handle client fetch data. The default threads number is the
number of flush thread. | 0.2.0 | |
+| celeborn.worker.fetch.port | 0 | Server port for Worker to receive fetch
data request from ShuffleClient. | 0.2.0 | |
+| celeborn.worker.flusher.buffer.size | 256k | Size of buffer used by a single
flusher. | 0.2.0 | |
+| celeborn.worker.flusher.diskTime.slidingWindow.size | 20 | The size of
sliding windows used to calculate statistics about flushed time and count. |
0.3.0 | celeborn.worker.flusher.avgFlushTime.slidingWindow.size |
+| celeborn.worker.flusher.hdd.threads | 1 | Flusher's thread count per disk
used for write data to HDD disks. | 0.2.0 | |
+| celeborn.worker.flusher.hdfs.buffer.size | 4m | Size of buffer used by a
HDFS flusher. | 0.3.0 | |
+| celeborn.worker.flusher.hdfs.threads | 8 | Flusher's thread count used for
write data to HDFS. | 0.2.0 | |
+| celeborn.worker.flusher.shutdownTimeout | 3s | Timeout for a flusher to
shutdown. | 0.2.0 | |
+| celeborn.worker.flusher.ssd.threads | 16 | Flusher's thread count per disk
used for write data to SSD disks. | 0.2.0 | |
+| celeborn.worker.flusher.threads | 16 | Flusher's thread count per disk for
unknown-type disks. | 0.2.0 | |
+| celeborn.worker.graceful.shutdown.checkSlotsFinished.interval | 1s | The
wait interval of checking whether all released slots to be committed or
destroyed during worker graceful shutdown | 0.2.0 | |
+| celeborn.worker.graceful.shutdown.checkSlotsFinished.timeout | 480s | The
wait time of waiting for the released slots to be committed or destroyed during
worker graceful shutdown. | 0.2.0 | |
+| celeborn.worker.graceful.shutdown.enabled | false | When true, during worker
shutdown, the worker will wait for all released slots to be committed or
destroyed. | 0.2.0 | |
+| celeborn.worker.graceful.shutdown.partitionSorter.shutdownTimeout | 120s |
The wait time of waiting for sorting partition files during worker graceful
shutdown. | 0.2.0 | |
+| celeborn.worker.graceful.shutdown.recoverDbBackend | LEVELDB | Specifies a
disk-based store used in local db. LEVELDB or ROCKSDB. | 0.4.0 | |
+| celeborn.worker.graceful.shutdown.recoverPath | <tmp>/recover | The
path to store DB. | 0.2.0 | |
+| celeborn.worker.graceful.shutdown.saveCommittedFileInfo.interval | 5s |
Interval for a Celeborn worker to flush committed file infos into Level DB. |
0.3.1 | |
+| celeborn.worker.graceful.shutdown.saveCommittedFileInfo.sync | false |
Whether to call sync method to save committed file infos into Level DB to
handle OS crash. | 0.3.1 | |
+| celeborn.worker.graceful.shutdown.timeout | 600s | The worker's graceful
shutdown timeout time. | 0.2.0 | |
+| celeborn.worker.http.host | <localhost> | Worker's http host. | 0.4.0
|
celeborn.metrics.worker.prometheus.host,celeborn.worker.metrics.prometheus.host
|
+| celeborn.worker.http.port | 9096 | Worker's http port. | 0.4.0 |
celeborn.metrics.worker.prometheus.port,celeborn.worker.metrics.prometheus.port
|
+| celeborn.worker.jvmQuake.check.interval | 1s | Interval of gc behavior
checking for worker jvm quake. | 0.4.0 | |
+| celeborn.worker.jvmQuake.dump.enabled | true | Whether to heap dump for the
maximum GC 'deficit' during worker jvm quake. | 0.4.0 | |
+| celeborn.worker.jvmQuake.dump.path | <tmp>/jvm-quake/dump/<pid>
| The path of heap dump for the maximum GC 'deficit' during worker jvm quake. |
0.4.0 | |
+| celeborn.worker.jvmQuake.dump.threshold | 30s | The threshold of heap dump
for the maximum GC 'deficit' which can be accumulated before jvmquake takes
action. Meanwhile, there is no heap dump generated when dump threshold is
greater than kill threshold. | 0.4.0 | |
+| celeborn.worker.jvmQuake.enabled | false | When true, Celeborn worker will
start the jvm quake to monitor of gc behavior, which enables early detection of
memory management issues and facilitates fast failure. | 0.4.0 | |
+| celeborn.worker.jvmQuake.exitCode | 502 | The exit code of system kill for
the maximum GC 'deficit' during worker jvm quake. | 0.4.0 | |
+| celeborn.worker.jvmQuake.kill.threshold | 60s | The threshold of system kill
for the maximum GC 'deficit' which can be accumulated before jvmquake takes
action. | 0.4.0 | |
+| celeborn.worker.jvmQuake.runtimeWeight | 5.0 | The factor by which to
multiply running JVM time, when weighing it against GCing time. 'Deficit' is
accumulated as `gc_time - runtime * runtime_weight`, and is compared against
threshold to determine whether to take action. | 0.4.0 | |
+| celeborn.worker.monitor.disk.check.interval | 30s | Intervals between device
monitor to check disk. | 0.3.0 | celeborn.worker.monitor.disk.checkInterval |
+| celeborn.worker.monitor.disk.check.timeout | 30s | Timeout time for worker
check device status. | 0.3.0 | celeborn.worker.disk.check.timeout |
+| celeborn.worker.monitor.disk.checklist | readwrite,diskusage | Monitor type
for disk, available items are: iohang, readwrite and diskusage. | 0.2.0 | |
+| celeborn.worker.monitor.disk.enabled | true | When true, worker will monitor
device and report to master. | 0.3.0 | |
+| celeborn.worker.monitor.disk.notifyError.expireTimeout | 10m | The expire
timeout of non-critical device error. Only notify critical error when the
number of non-critical errors for a period of time exceeds threshold. | 0.3.0 |
|
+| celeborn.worker.monitor.disk.notifyError.threshold | 64 | Device monitor
will only notify critical error once the accumulated valid non-critical error
number exceeding this threshold. | 0.3.0 | |
+| celeborn.worker.monitor.disk.sys.block.dir | /sys/block | The directory
where linux file block information is stored. | 0.2.0 | |
+| celeborn.worker.monitor.memory.check.interval | 10ms | Interval of worker
direct memory checking. | 0.3.0 | celeborn.worker.memory.checkInterval |
+| celeborn.worker.monitor.memory.report.interval | 10s | Interval of worker
direct memory tracker reporting to log. | 0.3.0 |
celeborn.worker.memory.reportInterval |
+| celeborn.worker.monitor.memory.trimChannelWaitInterval | 1s | Wait time
after worker trigger channel to trim cache. | 0.3.0 | |
+| celeborn.worker.monitor.memory.trimFlushWaitInterval | 1s | Wait time after
worker trigger StorageManger to flush data. | 0.3.0 | |
+| celeborn.worker.partition.initial.readBuffersMax | 1024 | Max number of
initial read buffers | 0.3.0 | |
+| celeborn.worker.partition.initial.readBuffersMin | 1 | Min number of initial
read buffers | 0.3.0 | |
+| celeborn.worker.partitionSorter.directMemoryRatioThreshold | 0.1 | Max ratio
of partition sorter's memory for sorting, when reserved memory is higher than
max partition sorter memory, partition sorter will stop sorting. | 0.2.0 | |
+| celeborn.worker.push.heartbeat.enabled | false | enable the heartbeat from
worker to client when pushing data | 0.3.0 | |
+| celeborn.worker.push.io.threads | <undefined> | Netty IO thread number
of worker to handle client push data. The default threads number is the number
of flush thread. | 0.2.0 | |
+| celeborn.worker.push.port | 0 | Server port for Worker to receive push data
request from ShuffleClient. | 0.2.0 | |
+| celeborn.worker.readBuffer.allocationWait | 50ms | The time to wait when
buffer dispatcher can not allocate a buffer. | 0.3.0 | |
+| celeborn.worker.readBuffer.target.changeThreshold | 1mb | The target ratio
for pre read memory usage. | 0.3.0 | |
+| celeborn.worker.readBuffer.target.ratio | 0.9 | The target ratio for read
ahead buffer's memory usage. | 0.3.0 | |
+| celeborn.worker.readBuffer.target.updateInterval | 100ms | The interval for
memory manager to calculate new read buffer's target memory. | 0.3.0 | |
+| celeborn.worker.readBuffer.toTriggerReadMin | 32 | Min buffers count for map
data partition to trigger read. | 0.3.0 | |
+| celeborn.worker.register.timeout | 180s | Worker register timeout. | 0.2.0 |
|
+| celeborn.worker.replicate.fastFail.duration | 60s | If a replicate request
not replied during the duration, worker will mark the replicate data request as
failed.It's recommended to set at least `240s` when `HDFS` is enabled in
`celeborn.storage.activeTypes`. | 0.2.0 | |
+| celeborn.worker.replicate.io.threads | <undefined> | Netty IO thread
number of worker to replicate shuffle data. The default threads number is the
number of flush thread. | 0.2.0 | |
+| celeborn.worker.replicate.port | 0 | Server port for Worker to receive
replicate data request from other Workers. | 0.2.0 | |
+| celeborn.worker.replicate.randomConnection.enabled | true | Whether worker
will create random connection to peer when replicate data. When false, worker
tend to reuse the same cached TransportClient to a specific replicate worker;
when true, worker tend to use different cached TransportClient. Netty will use
the same thread to serve the same connection, so with more connections
replicate server can leverage more netty threads | 0.2.1 | |
+| celeborn.worker.replicate.threads | 64 | Thread number of worker to
replicate shuffle data. | 0.2.0 | |
+| celeborn.worker.rpc.port | 0 | Server port for Worker to receive RPC
request. | 0.2.0 | |
+| celeborn.worker.shuffle.partitionSplit.enabled | true | enable the partition
split on worker side | 0.3.0 | celeborn.worker.partition.split.enabled |
+| celeborn.worker.shuffle.partitionSplit.max | 2g | Specify the maximum
partition size for splitting, and ensure that individual partition files are
always smaller than this limit. | 0.3.0 | |
+| celeborn.worker.shuffle.partitionSplit.min | 1m | Min size for a partition
to split | 0.3.0 | celeborn.shuffle.partitionSplit.min |
+| celeborn.worker.sortPartition.indexCache.expire | 180s | PartitionSorter's
cache item expire time. | 0.4.0 | |
+| celeborn.worker.sortPartition.indexCache.maxWeight | 100000 |
PartitionSorter's cache max weight for index buffer. | 0.4.0 | |
+| celeborn.worker.sortPartition.reservedMemoryPerPartition | 1mb | Reserved
memory when sorting a shuffle file off-heap. | 0.3.0 |
celeborn.worker.partitionSorter.reservedMemoryPerPartition |
+| celeborn.worker.sortPartition.threads | <undefined> |
PartitionSorter's thread counts. It's recommended to set at least `64` when
`HDFS` is enabled in `celeborn.storage.activeTypes`. | 0.3.0 |
celeborn.worker.partitionSorter.threads |
+| celeborn.worker.sortPartition.timeout | 220s | Timeout for a shuffle file to
sort. | 0.3.0 | celeborn.worker.partitionSorter.sort.timeout |
+| celeborn.worker.storage.checkDirsEmpty.maxRetries | 3 | The number of
retries for a worker to check if the working directory is cleaned up before
registering with the master. | 0.3.0 |
celeborn.worker.disk.checkFileClean.maxRetries |
+| celeborn.worker.storage.checkDirsEmpty.timeout | 1000ms | The wait time per
retry for a worker to check if the working directory is cleaned up before
registering with the master. | 0.3.0 |
celeborn.worker.disk.checkFileClean.timeout |
+| celeborn.worker.storage.dirs | <undefined> | Directory list to store
shuffle data. It's recommended to configure one directory on each disk. Storage
size limit can be set for each directory. For the sake of performance, there
should be no more than 2 flush threads on the same disk partition if you are
using HDD, and should be 8 or more flush threads on the same disk partition if
you are using SSD. For example:
`dir1[:capacity=][:disktype=][:flushthread=],dir2[:capacity=][:disktyp [...]
+| celeborn.worker.storage.disk.reserve.ratio | <undefined> | Celeborn
worker reserved ratio for each disk. The minimum usable size for each disk is
the max space between the reserved space and the space calculate via reserved
ratio. | 0.3.2 | |
+| celeborn.worker.storage.disk.reserve.size | 5G | Celeborn worker reserved
space for each disk. | 0.3.0 | celeborn.worker.disk.reserve.size |
+| celeborn.worker.storage.expireDirs.timeout | 1h | The timeout for a expire
dirs to be deleted on disk. | 0.3.2 | |
+| celeborn.worker.storage.workingDir | celeborn-worker/shuffle_data | Worker's
working dir path name. | 0.3.0 | celeborn.worker.workingDir |
+| celeborn.worker.userResourceConsumption.update.interval | 30s | Time length
for a window about compute user resource consumption. | 0.3.2 | |
+| celeborn.worker.writer.close.timeout | 120s | Timeout for a file writer to
close | 0.2.0 | |
+| celeborn.worker.writer.create.maxAttempts | 3 | Retry count for a file
writer to create if its creation was failed. | 0.2.0 | |
<!--end-include-->