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 " +

Reply via email to