This is an automated email from the ASF dual-hosted git repository.
zhouky pushed a commit to branch branch-0.3
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/branch-0.3 by this push:
new 44678e264 [CELEBORN-568] Support storage type selection
44678e264 is described below
commit 44678e26459012bca687f710df5e96605eb4b021
Author: mingji <[email protected]>
AuthorDate: Tue Jun 27 18:07:08 2023 +0800
[CELEBORN-568] Support storage type selection
### What changes were proposed in this pull request?
1. Celeborn supports storage type selection. HDD, SSD, and HDFS are
available for now.
2. Add new buffer size for HDFS file writers.
3. Worker support empty working dirs.
### Why are the changes needed?
Support HDFS only scenario.
### Does this PR introduce _any_ user-facing change?
NO.
### How was this patch tested?
UT and cluster.
Closes #1619 from FMX/CELEBORN-568.
Lead-authored-by: mingji <[email protected]>
Co-authored-by: Ethan Feng <[email protected]>
Signed-off-by: zky.zhoukeyong <[email protected]>
(cherry picked from commit 40760ede3a6d55e13be79724a86a4a8288fa297f)
Signed-off-by: zky.zhoukeyong <[email protected]>
---
README.md | 38 ++++++++++
.../celeborn/common/protocol/StorageInfo.java | 27 ++++---
.../org/apache/celeborn/common/CelebornConf.scala | 63 ++++++++++++----
.../apache/celeborn/common/CelebornConfSuite.scala | 21 ++++++
docs/configuration/client.md | 2 +-
docs/configuration/master.md | 3 +-
docs/configuration/network.md | 2 +-
docs/configuration/worker.md | 12 +--
docs/deploy.md | 36 +++++++++
.../master/clustermeta/AbstractMetaManager.java | 3 +-
.../celeborn/service/deploy/master/Master.scala | 14 ++--
.../deploy/master/SlotsAllocatorSuiteJ.java | 86 +++++++++++++++++-----
.../service/deploy/worker/storage/FileWriter.java | 3 +-
.../worker/storage/PartitionFilesSorter.java | 3 +-
.../service/deploy/worker/Controller.scala | 2 +-
.../service/deploy/worker/storage/FlushTask.scala | 2 +-
.../deploy/worker/storage/StorageManager.scala | 13 +++-
17 files changed, 262 insertions(+), 68 deletions(-)
diff --git a/README.md b/README.md
index 1b2558da4..f6b2ee6b5 100644
--- a/README.md
+++ b/README.md
@@ -113,8 +113,24 @@ celeborn.master.port 9097
celeborn.metrics.enabled true
celeborn.worker.flusher.buffer.size 256k
+
+# If Celeborn workers have local disks and HDFS. Following configs should be
added.
+# If Celeborn workers have local disks, use following config.
# Disk type is HDD by defaut.
celeborn.worker.storage.dirs /mnt/disk1:disktype=SSD,/mnt/disk2:disktype=SSD
+
+# If Celeborn workers don't have local disks. You can use HDFS.
+# Do not set `celeborn.worker.storage.dirs` and use following configs.
+celeborn.storage.activeTypes HDFS
+celeborn.worker.sortPartition.threads 64
+celeborn.worker.commitFiles.timeout 240s
+celeborn.worker.commitFiles.threads 128
+celeborn.master.slot.assign.policy roundrobin
+celeborn.rpc.askTimeout 240s
+celeborn.worker.flusher.hdfs.buffer.size 4m
+celeborn.worker.storage.hdfs.dir hdfs://<namenode>/celeborn
+celeborn.worker.replicate.fastFail.duration 240s
+
# If your hosts have disk raid or use lvm, set
celeborn.worker.monitor.disk.enabled to false
celeborn.worker.monitor.disk.enabled false
```
@@ -142,8 +158,26 @@ celeborn.master.ha.ratis.raft.server.storage.dir
/mnt/disk1/rss_ratis/
celeborn.metrics.enabled true
# If you want to use HDFS as shuffle storage, make sure that flush buffer size
is at least 4MB or larger.
celeborn.worker.flusher.buffer.size 256k
+
+# If Celeborn workers have local disks and HDFS. Following configs should be
added.
+# Celeborn will use local disks until local disk become unavailable to gain
the best performance.
+# Increase Celeborn's off-heap memory if Celeborn write to HDFS.
+# If Celeborn workers have local disks, use following config.
# Disk type is HDD by defaut.
celeborn.worker.storage.dirs /mnt/disk1:disktype=SSD,/mnt/disk2:disktype=SSD
+
+# If Celeborn workers don't have local disks. You can use HDFS.
+# Do not set `celeborn.worker.storage.dirs` and use following configs.
+celeborn.storage.activeTypes HDFS
+celeborn.worker.sortPartition.threads 64
+celeborn.worker.commitFiles.timeout 240s
+celeborn.worker.commitFiles.threads 128
+celeborn.master.slot.assign.policy roundrobin
+celeborn.rpc.askTimeout 240s
+celeborn.worker.flusher.hdfs.buffer.size 4m
+celeborn.worker.storage.hdfs.dir hdfs://<namenode>/celeborn
+celeborn.worker.replicate.fastFail.duration 240s
+
# If your hosts have disk raid or use lvm, set
celeborn.worker.monitor.disk.enabled to false
celeborn.worker.monitor.disk.enabled false
```
@@ -214,12 +248,16 @@ spark.celeborn.client.spark.shuffle.writer hash
# We recommend setting spark.celeborn.client.push.replicate.enabled to true to
enable server-side data replication
# If you have only one worker, this setting must be false
+# If your Celeborn is using HDFS, it's recommended to set this setting to false
spark.celeborn.client.push.replicate.enabled true
# Support for Spark AQE only tested under Spark 3
# we recommend setting localShuffleReader to false to get better performance
of Celeborn
spark.sql.adaptive.localShuffleReader.enabled false
+# If Celeborn is using HDFS
+spark.celeborn.worker.storage.hdfs.dir hdfs://<namenode>/celeborn
+
# we recommend enabling aqe support to gain better performance
spark.sql.adaptive.enabled true
spark.sql.adaptive.skewJoin.enabled true
diff --git
a/common/src/main/java/org/apache/celeborn/common/protocol/StorageInfo.java
b/common/src/main/java/org/apache/celeborn/common/protocol/StorageInfo.java
index 88a60ac70..b7bc0e878 100644
--- a/common/src/main/java/org/apache/celeborn/common/protocol/StorageInfo.java
+++ b/common/src/main/java/org/apache/celeborn/common/protocol/StorageInfo.java
@@ -18,12 +18,9 @@
package org.apache.celeborn.common.protocol;
import java.io.Serializable;
-import java.util.HashMap;
-import java.util.Map;
+import java.util.*;
public class StorageInfo implements Serializable {
- public static String UNKNOWN_DISK = "UNKNOWN_DISK";
-
public enum Type {
MEMORY(0),
HDD(1),
@@ -41,14 +38,16 @@ public class StorageInfo implements Serializable {
}
}
- public static Map<Integer, Type> typesMap =
- new HashMap() {
- {
- for (Type type : Type.values()) {
- put(type.value, type);
- }
- }
- };
+ public static String UNKNOWN_DISK = "UNKNOWN_DISK";
+ public static Map<Integer, Type> typesMap = new HashMap<>();
+ public static Set<String> typeNames = new HashSet<>();
+
+ static {
+ for (Type type : Type.values()) {
+ typesMap.put(type.value, type);
+ typeNames.add(type.name());
+ }
+ }
// Default storage Type is MEMORY.
private Type type = Type.MEMORY;
@@ -146,4 +145,8 @@ public class StorageInfo implements Serializable {
pbStorageInfo.getFinalResult(),
pbStorageInfo.getFilePath());
}
+
+ public static boolean validateStorageType(String str) {
+ return typeNames.contains(str);
+ }
}
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 af31918e8..8f5a4a752 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -498,6 +498,9 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable
with Logging with Se
// //////////////////////////////////////////////////////
def masterSlotAssignPolicy: SlotsAssignPolicy =
SlotsAssignPolicy.valueOf(get(MASTER_SLOT_ASSIGN_POLICY))
+
+ def hasHDFSStorage: Boolean =
get(ACTIVE_STORAGE_TYPES).contains(StorageInfo.Type.HDFS.name())
+
def masterSlotAssignLoadAwareDiskGroupNum: Int =
get(MASTER_SLOT_ASSIGN_LOADAWARE_DISKGROUP_NUM)
def masterSlotAssignLoadAwareDiskGroupGradient: Double =
get(MASTER_SLOT_ASSIGN_LOADAWARE_DISKGROUP_GRADIENT)
@@ -632,7 +635,8 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable
with Logging with Se
def workerCheckFileCleanTimeout: Long = get(WORKER_CHECK_FILE_CLEAN_TIMEOUT)
def workerHeartbeatTimeout: Long = get(WORKER_HEARTBEAT_TIMEOUT)
def workerReplicateThreads: Int = get(WORKER_REPLICATE_THREADS)
- def workerCommitThreads: Int = get(WORKER_COMMIT_THREADS)
+ def workerCommitThreads: Int =
+ if (hasHDFSStorage) Math.max(128, get(WORKER_COMMIT_THREADS)) else
get(WORKER_COMMIT_THREADS)
def workerShuffleCommitTimeout: Long = get(WORKER_SHUFFLE_COMMIT_TIMEOUT)
def minPartitionSizeToEstimate: Long = get(ESTIMATED_PARTITION_SIZE_MIN_SIZE)
def partitionSorterSortPartitionTimeout: Long =
get(PARTITION_SORTER_SORT_TIMEOUT)
@@ -846,10 +850,14 @@ class CelebornConf(loadDefaults: Boolean) extends
Cloneable with Logging with Se
(dir, maxCapacity, flushThread, diskType)
}
}.getOrElse {
- val prefix = workerStorageBaseDirPrefix
- val number = workerStorageBaseDirNumber
- (1 to number).map { i =>
- (s"$prefix$i", defaultMaxCapacity, workerHddFlusherThreads, HDD)
+ if (!hasHDFSStorage) {
+ val prefix = workerStorageBaseDirPrefix
+ val number = workerStorageBaseDirNumber
+ (1 to number).map { i =>
+ (s"$prefix$i", defaultMaxCapacity, workerHddFlusherThreads, HDD)
+ }
+ } else {
+ Seq.empty
}
}
}
@@ -898,6 +906,7 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable
with Logging with Se
// Flusher //
// //////////////////////////////////////////////////////
def workerFlusherBufferSize: Long = get(WORKER_FLUSHER_BUFFER_SIZE)
+ def workerHdfsFlusterBufferSize: Long = get(WORKER_HDFS_FLUSHER_BUFFER_SIZE)
def workerWriterCloseTimeoutMs: Long = get(WORKER_WRITER_CLOSE_TIMEOUT)
def workerHddFlusherThreads: Int = get(WORKER_FLUSHER_HDD_THREADS)
def workerSsdFlusherThreads: Int = get(WORKER_FLUSHER_SSD_THREADS)
@@ -1260,7 +1269,8 @@ object CelebornConf extends Logging {
buildConf("celeborn.rpc.askTimeout")
.categories("network")
.version("0.2.0")
- .doc("Timeout for RPC ask operations.")
+ .doc("Timeout for RPC ask operations. " +
+ "It's recommended to set at least `240s` when `HDFS` is enabled in
`celeborn.storage.activeTypes`")
.timeConf(TimeUnit.MILLISECONDS)
.createWithDefaultString("30s")
@@ -1765,10 +1775,13 @@ object CelebornConf extends Logging {
.withAlternative("celeborn.slots.assign.policy")
.categories("master")
.version("0.3.0")
- .doc("Policy for master to assign slots, Celeborn supports two types of
policy: roundrobin and loadaware.")
+ .doc("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`")
.stringConf
.transform(_.toUpperCase(Locale.ROOT))
- .checkValues(Set(SlotsAssignPolicy.ROUNDROBIN.name,
SlotsAssignPolicy.LOADAWARE.name))
+ .checkValues(Set(
+ SlotsAssignPolicy.ROUNDROBIN.name,
+ SlotsAssignPolicy.LOADAWARE.name))
.createWithDefault(SlotsAssignPolicy.ROUNDROBIN.name)
val MASTER_SLOT_ASSIGN_LOADAWARE_DISKGROUP_NUM: ConfigEntry[Int] =
@@ -2051,7 +2064,8 @@ object CelebornConf extends Logging {
val WORKER_REPLICATE_FAST_FAIL_DURATION: ConfigEntry[Long] =
buildConf("celeborn.worker.replicate.fastFail.duration")
.categories("worker")
- .doc("If a replicate request not replied during the duration, worker
will mark the replicate data request as failed.")
+ .doc("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`.")
.version("0.2.0")
.timeConf(TimeUnit.MILLISECONDS)
.createWithDefaultString("60s")
@@ -2080,7 +2094,8 @@ object CelebornConf extends Logging {
.withAlternative("celeborn.worker.commit.threads")
.categories("worker")
.version("0.3.0")
- .doc("Thread number of worker to commit shuffle data files
asynchronously.")
+ .doc("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`.")
.intConf
.createWithDefault(32)
@@ -2088,7 +2103,8 @@ object CelebornConf extends Logging {
buildConf("celeborn.worker.commitFiles.timeout")
.withAlternative("celeborn.worker.shuffle.commit.timeout")
.categories("worker")
- .doc("Timeout for a Celeborn worker to commit files of a shuffle.")
+ .doc("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`.")
.version("0.3.0")
.fallbackConf(RPC_ASK_TIMEOUT)
@@ -2105,7 +2121,8 @@ object CelebornConf extends Logging {
buildConf("celeborn.worker.sortPartition.threads")
.withAlternative("celeborn.worker.partitionSorter.threads")
.categories("worker")
- .doc("PartitionSorter's thread counts.")
+ .doc("PartitionSorter's thread counts. " +
+ "It's recommended to set at least `64` when `HDFS` is enabled in
`celeborn.storage.activeTypes`.")
.version("0.3.0")
.intConf
.createOptional
@@ -2127,6 +2144,14 @@ object CelebornConf extends Logging {
.bytesConf(ByteUnit.BYTE)
.createWithDefaultString("256k")
+ val WORKER_HDFS_FLUSHER_BUFFER_SIZE: ConfigEntry[Long] =
+ buildConf("celeborn.worker.flusher.hdfs.buffer.size")
+ .categories("worker")
+ .version("0.3.0")
+ .doc("Size of buffer used by a HDFS flusher.")
+ .bytesConf(ByteUnit.BYTE)
+ .createWithDefaultString("4m")
+
val WORKER_WRITER_CLOSE_TIMEOUT: ConfigEntry[Long] =
buildConf("celeborn.worker.writer.close.timeout")
.categories("worker")
@@ -2165,7 +2190,7 @@ object CelebornConf extends Logging {
.doc("Flusher's thread count used for write data to HDFS.")
.version("0.2.0")
.intConf
- .createWithDefault(4)
+ .createWithDefault(8)
val WORKER_FLUSHER_SHUTDOWN_TIMEOUT: ConfigEntry[Long] =
buildConf("celeborn.worker.flusher.shutdownTimeout")
@@ -2592,7 +2617,8 @@ object CelebornConf extends Logging {
.withAlternative("celeborn.push.replicate.enabled")
.categories("client")
.doc("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.")
+ "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`.")
.version("0.3.0")
.booleanConf
.createWithDefault(false)
@@ -3557,4 +3583,13 @@ object CelebornConf extends Logging {
.booleanConf
.createWithDefault(true)
+ val ACTIVE_STORAGE_TYPES: ConfigEntry[String] =
+ buildConf("celeborn.storage.activeTypes")
+ .categories("master", "worker")
+ .version("0.3.0")
+ .doc("Enabled storage levels. Available options: HDD,SSD,HDFS. ")
+ .stringConf
+ .transform(_.toUpperCase(Locale.ROOT))
+ .createWithDefault("HDD,SSD")
+
}
diff --git
a/common/src/test/scala/org/apache/celeborn/common/CelebornConfSuite.scala
b/common/src/test/scala/org/apache/celeborn/common/CelebornConfSuite.scala
index e7766040d..360ac12cc 100644
--- a/common/src/test/scala/org/apache/celeborn/common/CelebornConfSuite.scala
+++ b/common/src/test/scala/org/apache/celeborn/common/CelebornConfSuite.scala
@@ -186,4 +186,25 @@ class CelebornConfSuite extends CelebornFunSuite {
assert(conf.testAlternative == "rss")
}
+
+ test("Test empty working dir") {
+ val conf = new CelebornConf()
+ conf.set("celeborn.storage.activeTypes", "HDFS")
+ assert(conf.workerBaseDirs.isEmpty)
+
+ conf.set("celeborn.storage.activeTypes", "SDD,HDD,HDFS")
+ assert(conf.workerBaseDirs.isEmpty)
+
+ conf.set("celeborn.storage.activeTypes", "SDD,HDD")
+ assert(!conf.workerBaseDirs.isEmpty)
+ }
+
+ test("Test commit file threads") {
+ val conf = new CelebornConf()
+ conf.set("celeborn.storage.activeTypes", "HDFS")
+ assert(conf.workerCommitThreads === 128)
+
+ conf.set("celeborn.storage.activeTypes", "SDD,HDD")
+ assert(conf.workerCommitThreads === 32)
+ }
}
diff --git a/docs/configuration/client.md b/docs/configuration/client.md
index 25112c87e..badec4ae9 100644
--- a/docs/configuration/client.md
+++ b/docs/configuration/client.md
@@ -45,7 +45,7 @@ license: |
| celeborn.client.push.limit.strategy | SIMPLE | The strategy used to control
the push speed. Valid strategies are SIMPLE and SLOWSTART. the SLOWSTART
strategy is usually cooperate with congest control mechanism in the worker
side. | 0.3.0 |
| celeborn.client.push.maxReqsInFlight | 4 | Amount of Netty in-flight
requests per worker. The maximum memory is
`celeborn.client.push.maxReqsInFlight` * `celeborn.push.buffer.max.size` *
compression ratio(1 in worst case), default: 64Kib * 32 = 2Mib | 0.3.0 |
| celeborn.client.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.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. | 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.maxRetries | 5 | Max retry times for reviving
when celeborn push data failed. | 0.3.0 |
| celeborn.client.push.slowStart.initialSleepTime | 500ms | The initial sleep
time if the current max in flight requests is 0 | 0.3.0 |
diff --git a/docs/configuration/master.md b/docs/configuration/master.md
index fb7133ebc..8842b6c8c 100644
--- a/docs/configuration/master.md
+++ b/docs/configuration/master.md
@@ -31,6 +31,7 @@ license: |
| 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.policy | ROUNDROBIN | Policy for master to
assign slots, Celeborn supports two types of policy: roundrobin and loadaware.
| 0.3.0 |
+| 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.storage.activeTypes | HDD,SSD | Enabled storage levels. Available
options: HDD,SSD,HDFS. | 0.3.0 |
<!--end-include-->
diff --git a/docs/configuration/network.md b/docs/configuration/network.md
index 1784ca3ec..7cf89aebb 100644
--- a/docs/configuration/network.md
+++ b/docs/configuration/network.md
@@ -45,7 +45,7 @@ license: |
| 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 | 30s | Timeout for RPC ask operations. | 0.2.0 |
+| celeborn.rpc.askTimeout | 30s | 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.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 |
diff --git a/docs/configuration/worker.md b/docs/configuration/worker.md
index 0e57dd5e1..ccf4f29b3 100644
--- a/docs/configuration/worker.md
+++ b/docs/configuration/worker.md
@@ -22,10 +22,11 @@ license: |
| 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.activeTypes | HDD,SSD | Enabled storage levels. Available
options: HDD,SSD,HDFS. | 0.3.0 |
| celeborn.worker.bufferStream.threadsPerMountpoint | 8 | Threads count for
read buffer per mount point. | 0.3.0 |
| 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. | 0.3.0 |
-| celeborn.worker.commitFiles.timeout | <value of
celeborn.rpc.askTimeout> | Timeout for a Celeborn worker to commit files of
a shuffle. | 0.3.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 | <value of
celeborn.rpc.askTimeout> | 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.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 |
@@ -42,7 +43,8 @@ license: |
| 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.threads | 4 | Flusher's thread count used for
write data to HDFS. | 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 | 8 | Flusher's thread count per disk
used for write data to SSD disks. | 0.2.0 |
| celeborn.worker.flusher.threads | 2 | Flusher's thread count per disk for
unkown-type disks. | 0.2.0 |
@@ -75,7 +77,7 @@ license: |
| 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. | 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 |
@@ -84,7 +86,7 @@ license: |
| celeborn.worker.shuffle.partitionSplit.enabled | true | enable the partition
split on worker side | 0.3.0 |
| celeborn.worker.shuffle.partitionSplit.min | 1m | Min size for a partition
to split | 0.3.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. | 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 |
diff --git a/docs/deploy.md b/docs/deploy.md
index 9e600c011..bf5573ac5 100644
--- a/docs/deploy.md
+++ b/docs/deploy.md
@@ -39,8 +39,24 @@ celeborn.master.port 9097
celeborn.metrics.enabled true
celeborn.worker.flusher.buffer.size 256k
+
+# If Celeborn workers have local disks and HDFS. Following configs should be
added.
+# If Celeborn workers have local disks, use following config.
# Disk type is HDD by defaut.
celeborn.worker.storage.dirs /mnt/disk1:disktype=SSD,/mnt/disk2:disktype=SSD
+
+# If Celeborn workers don't have local disks. You can use HDFS.
+# Do not set `celeborn.worker.storage.dirs` and use following configs.
+celeborn.storage.activeTypes HDFS
+celeborn.worker.sortPartition.threads 64
+celeborn.worker.commitFiles.timeout 240s
+celeborn.worker.commitFiles.threads 128
+celeborn.master.slot.assign.policy roundrobin
+celeborn.rpc.askTimeout 240s
+celeborn.worker.flusher.hdfs.buffer.size 4m
+celeborn.worker.storage.hdfs.dir hdfs://<namenode>/celeborn
+celeborn.worker.replicate.fastFail.duration 240s
+
# If your hosts have disk raid or use lvm, set
celeborn.worker.monitor.disk.enabled to false
celeborn.worker.monitor.disk.enabled false
```
@@ -68,8 +84,24 @@ celeborn.master.ha.ratis.raft.server.storage.dir
/mnt/disk1/rss_ratis/
celeborn.metrics.enabled true
# If you want to use HDFS as shuffle storage, make sure that flush buffer size
is at least 4MB or larger.
celeborn.worker.flusher.buffer.size 256k
+
+# If Celeborn workers have local disks and HDFS. Following configs should be
added.
+# If Celeborn workers have local disks, use following config.
# Disk type is HDD by default.
celeborn.worker.storage.dirs /mnt/disk1:disktype=SSD,/mnt/disk2:disktype=SSD
+
+# If Celeborn workers don't have local disks. You can use HDFS.
+# Do not set `celeborn.worker.storage.dirs` and use following configs.
+celeborn.storage.activeTypes HDFS
+celeborn.worker.sortPartition.threads 64
+celeborn.worker.commitFiles.timeout 240s
+celeborn.worker.commitFiles.threads 128
+celeborn.master.slot.assign.policy roundrobin
+celeborn.rpc.askTimeout 240s
+celeborn.worker.flusher.hdfs.buffer.size 4m
+celeborn.worker.storage.hdfs.dir hdfs://<namenode>/celeborn
+celeborn.worker.replicate.fastFail.duration 240s
+
# If your hosts have disk raid or use lvm, set
celeborn.worker.monitor.disk.enabled to false
celeborn.worker.monitor.disk.enabled false
```
@@ -137,12 +169,16 @@ spark.celeborn.client.spark.shuffle.writer hash
# we recommend set spark.celeborn.client.push.replicate.enabled to true to
enable server-side data replication
# If you have only one worker, this setting must be false
+# If your Celeborn is using HDFS, it's recommended to set this setting to false
spark.celeborn.client.push.replicate.enabled true
# Support for Spark AQE only tested under Spark 3
# we recommend set localShuffleReader to false to get better performance of
Celeborn
spark.sql.adaptive.localShuffleReader.enabled false
+# If Celeborn is using HDFS
+spark.celeborn.worker.storage.hdfs.dir hdfs://<namenode>/celeborn
+
# we recommend enabling aqe support to gain better performance
spark.sql.adaptive.enabled true
spark.sql.adaptive.skewJoin.enabled true
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
index c9bfda809..d51079c99 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
@@ -202,7 +202,8 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
});
}
appDiskUsageMetric.update(estimatedAppDiskUsage);
- if (!blacklist.contains(worker) && disks.isEmpty()) {
+ // If using HDFSONLY mode, workers with empty disks should not be put into
blacklist.
+ if (!blacklist.contains(worker) && (disks.isEmpty() &&
!conf.hasHDFSStorage())) {
LOG.debug("Worker: {} num total slots is 0, add to blacklist", worker);
blacklist.add(worker);
} else if (availableSlots.get() > 0) {
diff --git
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
index 7d7811a1a..73182d50a 100644
---
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
+++
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
@@ -541,13 +541,7 @@ private[celeborn] class Master(
val slots =
masterSource.sample(MasterSource.OfferSlotsTime,
s"offerSlots-${Random.nextInt()}") {
statusSystem.workers.synchronized {
- if (slotsAssignPolicy == SlotsAssignPolicy.ROUNDROBIN) {
- SlotsAllocator.offerSlotsRoundRobin(
- availableWorkers,
- requestSlots.partitionIdList,
- requestSlots.shouldReplicate,
- requestSlots.shouldRackAware)
- } else {
+ if (slotsAssignPolicy == SlotsAssignPolicy.LOADAWARE &&
!conf.hasHDFSStorage) {
SlotsAllocator.offerSlotsLoadAware(
availableWorkers,
requestSlots.partitionIdList,
@@ -558,6 +552,12 @@ private[celeborn] class Master(
slotsAssignLoadAwareDiskGroupGradient,
loadAwareFlushTimeWeight,
loadAwareFetchTimeWeight)
+ } else {
+ SlotsAllocator.offerSlotsRoundRobin(
+ availableWorkers,
+ requestSlots.partitionIdList,
+ requestSlots.shouldReplicate,
+ requestSlots.shouldRackAware)
}
}
}
diff --git
a/master/src/test/java/org/apache/celeborn/service/deploy/master/SlotsAllocatorSuiteJ.java
b/master/src/test/java/org/apache/celeborn/service/deploy/master/SlotsAllocatorSuiteJ.java
index ef62ac935..1d7428116 100644
---
a/master/src/test/java/org/apache/celeborn/service/deploy/master/SlotsAllocatorSuiteJ.java
+++
b/master/src/test/java/org/apache/celeborn/service/deploy/master/SlotsAllocatorSuiteJ.java
@@ -38,7 +38,7 @@ import org.apache.celeborn.common.meta.WorkerInfo;
import org.apache.celeborn.common.protocol.PartitionLocation;
public class SlotsAllocatorSuiteJ {
- private List<WorkerInfo> prepareWorkers() {
+ private List<WorkerInfo> prepareWorkers(boolean hasDisks) {
long assumedPartitionSize = 64 * 1024 * 1024;
Random random = new Random();
@@ -67,9 +67,11 @@ public class SlotsAllocatorSuiteJ {
diskInfo1.maxSlots_$eq(diskInfo1.actualUsableSpace() /
assumedPartitionSize);
diskInfo2.maxSlots_$eq(diskInfo2.actualUsableSpace() /
assumedPartitionSize);
diskInfo3.maxSlots_$eq(diskInfo3.actualUsableSpace() /
assumedPartitionSize);
- disks1.put("/mnt/disk1", diskInfo1);
- disks1.put("/mnt/disk2", diskInfo2);
- disks1.put("/mnt/disk3", diskInfo3);
+ if (hasDisks) {
+ disks1.put("/mnt/disk1", diskInfo1);
+ disks1.put("/mnt/disk2", diskInfo2);
+ disks1.put("/mnt/disk3", diskInfo3);
+ }
Map<String, DiskInfo> disks2 = new HashMap<>();
DiskInfo diskInfo4 =
@@ -96,9 +98,11 @@ public class SlotsAllocatorSuiteJ {
diskInfo4.maxSlots_$eq(diskInfo4.actualUsableSpace() /
assumedPartitionSize);
diskInfo5.maxSlots_$eq(diskInfo5.actualUsableSpace() /
assumedPartitionSize);
diskInfo6.maxSlots_$eq(diskInfo6.actualUsableSpace() /
assumedPartitionSize);
- disks2.put("/mnt/disk1", diskInfo4);
- disks2.put("/mnt/disk2", diskInfo5);
- disks2.put("/mnt/disk3", diskInfo6);
+ if (hasDisks) {
+ disks2.put("/mnt/disk1", diskInfo4);
+ disks2.put("/mnt/disk2", diskInfo5);
+ disks2.put("/mnt/disk3", diskInfo6);
+ }
Map<String, DiskInfo> disks3 = new HashMap<>();
DiskInfo diskInfo7 =
@@ -125,9 +129,11 @@ public class SlotsAllocatorSuiteJ {
diskInfo7.maxSlots_$eq(diskInfo7.actualUsableSpace() /
assumedPartitionSize);
diskInfo8.maxSlots_$eq(diskInfo8.actualUsableSpace() /
assumedPartitionSize);
diskInfo9.maxSlots_$eq(diskInfo9.actualUsableSpace() /
assumedPartitionSize);
- disks3.put("/mnt/disk1", diskInfo7);
- disks3.put("/mnt/disk2", diskInfo8);
- disks3.put("/mnt/disk3", diskInfo9);
+ if (hasDisks) {
+ disks3.put("/mnt/disk2", diskInfo8);
+ disks3.put("/mnt/disk1", diskInfo7);
+ disks3.put("/mnt/disk3", diskInfo9);
+ }
ArrayList<WorkerInfo> workers = new ArrayList<>(3);
workers.add(new WorkerInfo("host1", 9, 10, 110, 113, disks1, null));
@@ -138,7 +144,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocateSlotsForEmptyPartitionId() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = Collections.emptyList();
final boolean shouldReplicate = true;
@@ -147,7 +153,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocateSlotsForSinglePartitionId() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = Collections.singletonList(0);
final boolean shouldReplicate = true;
@@ -156,7 +162,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocateSlotsForSinglePartitionIdWithoutReplicate() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = Collections.singletonList(0);
final boolean shouldReplicate = false;
@@ -165,7 +171,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocateSlotsForTwoPartitionIdsWithoutReplicate() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = Arrays.asList(0, 1);
final boolean shouldReplicate = false;
@@ -174,7 +180,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocateSlotsForThreePartitionIdsWithoutReplicate() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = Arrays.asList(0, 1, 2);
final boolean shouldReplicate = false;
@@ -183,7 +189,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocateSlotsForThreeReduceIdsWithReplicate() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = Arrays.asList(0, 1, 2);
final boolean shouldReplicate = true;
@@ -192,7 +198,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocate3000ReduceIdsWithReplicate() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = new ArrayList<>();
for (int i = 0; i < 3000; i++) {
partitionIds.add(i);
@@ -204,7 +210,7 @@ public class SlotsAllocatorSuiteJ {
@Test
public void testAllocate3000ReduceIdsWithoutReplicate() {
- final List<WorkerInfo> workers = prepareWorkers();
+ final List<WorkerInfo> workers = prepareWorkers(true);
final List<Integer> partitionIds = new ArrayList<>();
for (int i = 0; i < 3000; i++) {
partitionIds.add(i);
@@ -274,4 +280,48 @@ public class SlotsAllocatorSuiteJ {
: "Expect to fail to offer slots, but return " + slots.size() + "
slots.";
}
}
+
+ private void checkSlotsOnHDFS(
+ List<WorkerInfo> workers,
+ List<Integer> partitionIds,
+ boolean shouldReplicate,
+ boolean expectSuccess) {
+ String shuffleKey = "appId-1";
+ CelebornConf conf = new CelebornConf();
+ conf.set("celeborn.active.storage.levels", "HDFS");
+ Map<WorkerInfo, Tuple2<List<PartitionLocation>, List<PartitionLocation>>>
slots =
+ SlotsAllocator.offerSlotsRoundRobin(workers, partitionIds,
shouldReplicate, false);
+
+ int allocatedPartitionCount = 0;
+
+ for (Map.Entry<WorkerInfo, Tuple2<List<PartitionLocation>,
List<PartitionLocation>>>
+ workerToPartitions : slots.entrySet()) {
+ WorkerInfo workerInfo = workerToPartitions.getKey();
+ List<PartitionLocation> masterLocs = workerToPartitions.getValue()._1;
+ List<PartitionLocation> slaveLocs = workerToPartitions.getValue()._2();
+ allocatedPartitionCount += masterLocs.size();
+ allocatedPartitionCount += slaveLocs.size();
+ }
+ if (expectSuccess) {
+ Assert.assertEquals(slots.isEmpty(), false);
+ } else {
+ Assert.assertEquals(slots.isEmpty(), true);
+ }
+ if (shouldReplicate) {
+ Assert.assertEquals(allocatedPartitionCount, partitionIds.size() * 2);
+ } else {
+ Assert.assertEquals(allocatedPartitionCount, partitionIds.size());
+ }
+ }
+
+ @Test
+ public void testHDFSOnly() {
+ final List<WorkerInfo> workers = prepareWorkers(false);
+ final List<Integer> partitionIds = new ArrayList<>();
+ for (int i = 0; i < 3000; i++) {
+ partitionIds.add(i);
+ }
+ final boolean shouldReplicate = true;
+ checkSlotsOnHDFS(workers, partitionIds, shouldReplicate, true);
+ }
}
diff --git
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/FileWriter.java
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/FileWriter.java
index 1ce41c0fd..79c2d885b 100644
---
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/FileWriter.java
+++
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/FileWriter.java
@@ -98,14 +98,15 @@ public abstract class FileWriter implements DeviceObserver {
this.flushWorkerIndex = flusher.getWorkerIndex();
this.writerCloseTimeoutMs = conf.workerWriterCloseTimeoutMs();
this.splitThreshold = splitThreshold;
- this.flusherBufferSize = conf.workerFlusherBufferSize();
this.deviceMonitor = deviceMonitor;
this.splitMode = splitMode;
this.partitionType = partitionType;
this.rangeReadFilter = rangeReadFilter;
if (!fileInfo.isHdfs()) {
+ this.flusherBufferSize = conf.workerFlusherBufferSize();
channel =
FileChannelUtils.createWritableFileChannel(fileInfo.getFilePath());
} else {
+ this.flusherBufferSize = conf.workerHdfsFlusterBufferSize();
// We open the stream and close immediately because HDFS output stream
will
// create a DataStreamer that is a thread.
// If we reuse HDFS output stream, we will exhaust the memory soon.
diff --git
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
index 4994f7372..9ab516408 100644
---
a/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
+++
b/worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java
@@ -600,7 +600,8 @@ public class PartitionFilesSorter extends
ShuffleRecoverHelper {
private void initializeFiles() throws IOException {
if (isHdfs) {
hdfsOriginInput = StorageManager.hadoopFs().open(new
Path(originFilePath));
- hdfsSortedOutput = StorageManager.hadoopFs().create(new
Path(sortedFilePath));
+ hdfsSortedOutput =
+ StorageManager.hadoopFs().create(new Path(sortedFilePath), true,
256 * 1024);
} else {
originFileChannel =
FileChannelUtils.openReadableFileChannel(originFilePath);
sortedFileChannel =
FileChannelUtils.createWritableFileChannel(sortedFilePath);
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
index ca4b6dc9e..5093c05a1 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
@@ -151,7 +151,7 @@ private[deploy] class Controller(
return
}
- if (storageManager.healthyWorkingDirs().size <= 0 &&
storageManager.hdfsDir.isEmpty) {
+ if (storageManager.healthyWorkingDirs().size <= 0 && !conf.hasHDFSStorage)
{
val msg = "Local storage has no available dirs!"
logError(s"[handleReserveSlots] $msg")
context.reply(ReserveSlotsResponse(StatusCode.NO_AVAILABLE_WORKING_DIR,
msg))
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/FlushTask.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/FlushTask.scala
index d00a3353a..6c6bc64aa 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/FlushTask.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/FlushTask.scala
@@ -47,7 +47,7 @@ private[worker] class HdfsFlushTask(
val path: Path,
notifier: FlushNotifier) extends FlushTask(buffer, notifier) {
override def flush(): Unit = {
- val hdfsStream = StorageManager.hadoopFs.append(path)
+ val hdfsStream = StorageManager.hadoopFs.append(path, 256 * 1024)
hdfsStream.write(ByteBufUtil.getBytes(buffer))
hdfsStream.close()
}
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
index 2208497c6..ae29d79df 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
@@ -59,7 +59,7 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
(new File(workdir, conf.workerWorkingDir), maxSpace, flusherThread,
storageType)
}
- if (workingDirInfos.size <= 0) {
+ if (workingDirInfos.size <= 0 && !conf.hasHDFSStorage) {
throw new IOException("Empty working directory configuration!")
}
@@ -120,13 +120,13 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
deviceMonitor.startCheck()
val hdfsDir = conf.hdfsDir
- if (!hdfsDir.isEmpty) {
+ if (!hdfsDir.isEmpty && conf.hasHDFSStorage) {
logInfo(s"Initialize HDFS support with path ${hdfsDir}")
}
val hdfsPermission = new FsPermission("755")
val hdfsWriters = JavaUtils.newConcurrentHashMap[String, FileWriter]()
val (hdfsFlusher, _totalHdfsFlusherThread) =
- if (!hdfsDir.isEmpty) {
+ if (!hdfsDir.isEmpty && conf.hasHDFSStorage) {
val path = new Path(hdfsDir)
val scheme = path.toUri.getScheme
val disableCacheName = String.format("fs.%s.impl.disable.cache", scheme)
@@ -265,7 +265,7 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
partitionType: PartitionType,
rangeReadFilter: Boolean,
userIdentifier: UserIdentifier): FileWriter = {
- if (healthyWorkingDirs().size <= 0 && hdfsDir.isEmpty) {
+ if (healthyWorkingDirs().size <= 0 && !conf.hasHDFSStorage) {
throw new IOException("No available working dirs!")
}
@@ -665,6 +665,11 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
})
}
})
+ hdfsWriters.forEach(new BiConsumer[String, FileWriter] {
+ override def accept(t: String, u: FileWriter): Unit = {
+ u.flushOnMemoryPressure();
+ }
+ })
}
override def onPause(moduleName: String): Unit = {}