This is an automated email from the ASF dual-hosted git repository.
angerszhuuuu 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 611b0f539 [CELEBORN-718] Support override Hadoop Conf by Celeborn Conf
with `celeborn.hadoop.` prefix
611b0f539 is described below
commit 611b0f539f94f8147498d9d00c92b8548a5848fc
Author: Angerszhuuuu <[email protected]>
AuthorDate: Tue Jun 27 17:00:47 2023 +0800
[CELEBORN-718] Support override Hadoop Conf by Celeborn Conf with
`celeborn.hadoop.` prefix
### What changes were proposed in this pull request?
Celeborn generate hadoop configuration should respect Celeborn conf
### Why are the changes needed?
In spark client side we should write like `spark.celeborn.hadoop.xxx.xx`
In server side we should write like `celeborn.hadoop.xxx.xxx`
### Does this PR introduce _any_ user-facing change?
### How was this patch tested?
Closes #1629 from AngersZhuuuu/CELEBORN-719.
Authored-by: Angerszhuuuu <[email protected]>
Signed-off-by: Angerszhuuuu <[email protected]>
(cherry picked from commit a2b215bd477f7d923c4de447bbc711f863806c03)
Signed-off-by: Angerszhuuuu <[email protected]>
---
client/src/main/java/org/apache/celeborn/client/ShuffleClient.java | 3 ++-
.../org/apache/celeborn/common/util}/CelebornHadoopUtils.scala | 2 +-
docs/migration.md | 4 ++++
.../service/deploy/master/network/CelebornRackResolver.scala | 2 +-
.../celeborn/service/deploy/worker/storage/StorageManager.scala | 6 +++---
5 files changed, 11 insertions(+), 6 deletions(-)
diff --git a/client/src/main/java/org/apache/celeborn/client/ShuffleClient.java
b/client/src/main/java/org/apache/celeborn/client/ShuffleClient.java
index 2ec83dc82..5e833f5fe 100644
--- a/client/src/main/java/org/apache/celeborn/client/ShuffleClient.java
+++ b/client/src/main/java/org/apache/celeborn/client/ShuffleClient.java
@@ -30,6 +30,7 @@ import org.apache.celeborn.common.CelebornConf;
import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.protocol.PartitionLocation;
import org.apache.celeborn.common.rpc.RpcEndpointRef;
+import org.apache.celeborn.common.util.CelebornHadoopUtils$;
import org.apache.celeborn.common.write.PushState;
/**
@@ -83,7 +84,7 @@ public abstract class ShuffleClient {
if (null == hdfsFs) {
synchronized (ShuffleClient.class) {
if (null == hdfsFs) {
- Configuration hdfsConfiguration = new Configuration();
+ Configuration hdfsConfiguration =
CelebornHadoopUtils$.MODULE$.newConfiguration(conf);
// enable fs cache to avoid too many fs instances
hdfsConfiguration.set("fs.hdfs.impl.disable.cache", "false");
hdfsConfiguration.set("fs.viewfs.impl.disable.cache", "false");
diff --git
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/utils/CelebornHadoopUtils.scala
b/common/src/main/scala/org/apache/celeborn/common/util/CelebornHadoopUtils.scala
similarity index 96%
rename from
master/src/main/scala/org/apache/celeborn/service/deploy/master/utils/CelebornHadoopUtils.scala
rename to
common/src/main/scala/org/apache/celeborn/common/util/CelebornHadoopUtils.scala
index 109a7afbb..39b70ff18 100644
---
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/utils/CelebornHadoopUtils.scala
+++
b/common/src/main/scala/org/apache/celeborn/common/util/CelebornHadoopUtils.scala
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.celeborn.service.deploy.master.utils
+package org.apache.celeborn.common.util
import org.apache.hadoop.conf.Configuration
diff --git a/docs/migration.md b/docs/migration.md
index cd97080a8..d19647a0a 100644
--- a/docs/migration.md
+++ b/docs/migration.md
@@ -55,3 +55,7 @@ license: |
at
org.apache.celeborn.common.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:76)
at
org.apache.celeborn.common.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:110)
```
+
+ - Since 0.3.0, Celeborn supports overriding Hadoop
configuration(`core-site.xml`, `hdfs-site.xml`, etc.) from Celeborn
configuration with the additional prefix `celeborn.hadoop.`.
+ On Spark client side, user should set Hadoop configuration like
`spark.celeborn.hadoop.foo=bar`, note that `spark.hadoop.foo=bar` does not take
effect;
+ on Flink client and Celeborn Master/Worker side, user should set like
`celeborn.hadoop.foo=bar`.
\ No newline at end of file
diff --git
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/network/CelebornRackResolver.scala
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/network/CelebornRackResolver.scala
index c47cf8f23..83882e842 100644
---
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/network/CelebornRackResolver.scala
+++
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/network/CelebornRackResolver.scala
@@ -28,7 +28,7 @@ import org.apache.hadoop.util.ReflectionUtils
import org.apache.celeborn.common.CelebornConf
import org.apache.celeborn.common.internal.Logging
-import org.apache.celeborn.service.deploy.master.utils.CelebornHadoopUtils
+import org.apache.celeborn.common.util.CelebornHadoopUtils
class CelebornRackResolver(celebornConf: CelebornConf) extends Logging {
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 f66d20425..2208497c6 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
@@ -41,7 +41,7 @@ import org.apache.celeborn.common.meta.{DeviceInfo, DiskInfo,
DiskStatus, FileIn
import org.apache.celeborn.common.metrics.source.AbstractSource
import org.apache.celeborn.common.protocol.{PartitionLocation,
PartitionSplitMode, PartitionType}
import org.apache.celeborn.common.quota.ResourceConsumption
-import org.apache.celeborn.common.util.{JavaUtils, PbSerDeUtils, ThreadUtils,
Utils}
+import org.apache.celeborn.common.util.{CelebornHadoopUtils, JavaUtils,
PbSerDeUtils, ThreadUtils, Utils}
import org.apache.celeborn.service.deploy.worker._
import
org.apache.celeborn.service.deploy.worker.memory.MemoryManager.MemoryPressureListener
import
org.apache.celeborn.service.deploy.worker.storage.StorageManager.hadoopFs
@@ -129,8 +129,8 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
if (!hdfsDir.isEmpty) {
val path = new Path(hdfsDir)
val scheme = path.toUri.getScheme
- val disableCacheName = String.format("fs.%s.impl.disable.cache", scheme);
- val hdfsConfiguration = new Configuration
+ val disableCacheName = String.format("fs.%s.impl.disable.cache", scheme)
+ val hdfsConfiguration = CelebornHadoopUtils.newConfiguration(conf)
hdfsConfiguration.set("dfs.replication", "2")
hdfsConfiguration.set(disableCacheName, "false")
logInfo("Celeborn will ignore cluster settings " +