| celeborn.application.heartbeatInterval |
10s |
Interval for client to send heartbeat message to master. |
0.2.0 |
| celeborn.client.maxRetries |
15 |
Max retry times for client to connect master endpoint |
0.2.0 |
| celeborn.client.rpc.askTimeout |
<value of celeborn.network.timeout> |
Timeout for client RPC ask operations. |
0.2.0 |
| celeborn.fetch.maxReqsInFlight |
3 |
Amount of in-flight chunk fetch request. |
0.2.0 |
| celeborn.fetch.timeout |
120s |
Timeout for a task to fetch chunk. |
0.2.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.push.buffer.initial.size |
8k |
|
0.2.0 |
| celeborn.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.2.0 |
| celeborn.push.limit.inFlight.sleepInterval |
50ms |
Sleep interval when check netty in-flight requests to be done. |
0.2.0 |
| celeborn.push.limit.inFlight.timeout |
240s |
Timeout for netty in-flight requests to be done. |
0.2.0 |
| celeborn.push.maxReqsInFlight |
32 |
Amount of Netty in-flight requests. The maximum memory is celeborn.push.maxReqsInFlight * celeborn.push.buffer.max.size * compression ratio(1 in worst case), default: 64Kib * 32 = 2Mib |
0.2.0 |
| celeborn.push.queue.capacity |
512 |
Push buffer queue size for a task. The maximum memory is celeborn.push.buffer.max.size * celeborn.push.queue.capacity, default: 64KiB * 512 = 32MiB |
0.2.0 |
| celeborn.push.replicate.enabled |
true |
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. |
0.2.0 |
| celeborn.push.retry.threads |
8 |
Thread number to process shuffle re-send push data requests. |
0.2.0 |
| celeborn.push.sortMemory.threshold |
64m |
When SortBasedPusher use memory over the threshold, will trigger push data. |
0.2.0 |
| celeborn.push.splitPartition.threads |
8 |
Thread number to process shuffle split request in shuffle client. |
0.2.0 |
| celeborn.push.stageEnd.timeout |
240s |
Timeout for StageEnd. |
0.2.0 |
| celeborn.rpc.cache.concurrencyLevel |
32 |
The number of write locks to update rpc cache. |
0.2.0 |
| celeborn.rpc.cache.expireTime |
15s |
The time before a cache item is removed. |
0.2.0 |
| celeborn.rpc.cache.size |
256 |
The max cache items count for rpc cache. |
0.2.0 |
| celeborn.rpc.maxParallelism |
1024 |
Max parallelism of client on sending RPC requests. |
0.2.0 |
| celeborn.shuffle.batchHandleChangePartition.enabled |
false |
When true, LifecycleManager will handle change partition request in batch. Otherwise, LifecycleManager will process the requests one by one |
0.2.0 |
| celeborn.shuffle.batchHandleChangePartition.interval |
100ms |
Interval for LifecycleManager to schedule handling change partition requests in batch. |
0.2.0 |
| celeborn.shuffle.batchHandleChangePartition.threads |
8 |
Threads number for LifecycleManager to handle change partition request in batch. |
0.2.0 |
| celeborn.shuffle.chuck.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.shuffle.compression.codec |
LZ4 |
The codec used to compress shuffle data. By default, Celeborn provides two codecs: lz4 and zstd. |
0.2.0 |
| celeborn.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.2.0 |
| celeborn.shuffle.expired.checkInterval |
60s |
Interval for client to check expired shuffles. |
0.2.0 |
| celeborn.shuffle.forceFallback.enabled |
false |
Whether force fallback shuffle to Spark's default. |
0.2.0 |
| celeborn.shuffle.forceFallback.numPartitionsThreshold |
500000 |
Celeborn will only accept shuffle of partition number lower than this configuration value. |
0.2.0 |
| celeborn.shuffle.manager.port |
0 |
Port used by the LifecycleManager on the Driver. |
0.2.0 |
| celeborn.shuffle.partition.type |
REDUCE |
Type of shuffle's partition. |
0.2.0 |
| celeborn.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.2.0 |
| celeborn.shuffle.partitionSplit.threshold |
1G |
Shuffle file size threshold, if file size exceeds this, trigger split. |
0.2.0 |
| celeborn.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.register.maxRetries |
3 |
Max retry times for client to register shuffle. |
0.2.0 |
| celeborn.shuffle.register.retryWait |
3s |
Wait time before next retry if register shuffle failed. |
0.2.0 |
| celeborn.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 it huge. |
0.2.0 |
| celeborn.slots.reserve.maxRetries |
3 |
Max retry times for client to reserve slots. |
0.2.0 |
| celeborn.slots.reserve.retryWait |
3s |
Wait time before next retry if reserve slots failed. |
0.2.0 |
| celeborn.storage.hdfs.dir |
<undefined> |
HDFS dir configuration for Celeborn to access HDFS. |
0.2.0 |
| celeborn.worker.excluded.checkInterval |
30s |
Interval for client to refresh excluded worker list. |
0.2.0 |
| celeborn.worker.excluded.expireTimeout |
600s |
Timeout time for LifecycleManager to clear reserved excluded worker. |
0.2.0 |