From 827ee04aad51a3ef77aaf65923d98389cea7c9dc Mon Sep 17 00:00:00 2001 From: baoloongmao Date: Mon, 18 Nov 2024 16:04:32 +0800 Subject: [PATCH 1/7] Put sparkConf as extra properties while client request accessCluster --- .../apache/spark/shuffle/RssSparkConfig.java | 13 +++++++++++++ .../shuffle/manager/RssShuffleManagerBase.java | 17 ++--------------- .../shuffle/DelegationRssShuffleManager.java | 1 + 3 files changed, 16 insertions(+), 15 deletions(-) diff --git a/client-spark/common/src/main/java/org/apache/spark/shuffle/RssSparkConfig.java b/client-spark/common/src/main/java/org/apache/spark/shuffle/RssSparkConfig.java index 734dedcccd..2d767980de 100644 --- a/client-spark/common/src/main/java/org/apache/spark/shuffle/RssSparkConfig.java +++ b/client-spark/common/src/main/java/org/apache/spark/shuffle/RssSparkConfig.java @@ -17,6 +17,8 @@ package org.apache.spark.shuffle; +import java.util.HashMap; +import java.util.Map; import java.util.Set; import scala.Tuple2; @@ -512,4 +514,15 @@ public static RssConf toRssConf(SparkConf sparkConf) { } return rssConf; } + + public static Map sparkConfToMap(SparkConf sparkConf) { + Map map = new HashMap<>(); + + for (Tuple2 tuple : sparkConf.getAll()) { + String key = tuple._1; + map.put(key, tuple._2); + } + + return map; + } } diff --git a/client-spark/common/src/main/java/org/apache/uniffle/shuffle/manager/RssShuffleManagerBase.java b/client-spark/common/src/main/java/org/apache/uniffle/shuffle/manager/RssShuffleManagerBase.java index 8e921c66e7..c1b697cbde 100644 --- a/client-spark/common/src/main/java/org/apache/uniffle/shuffle/manager/RssShuffleManagerBase.java +++ b/client-spark/common/src/main/java/org/apache/uniffle/shuffle/manager/RssShuffleManagerBase.java @@ -34,8 +34,6 @@ import java.util.function.Supplier; import java.util.stream.Collectors; -import scala.Tuple2; - import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Maps; import com.google.common.collect.Sets; @@ -1064,7 +1062,7 @@ protected void registerShuffleServers( } LOG.info("Start to register shuffleId {}", shuffleId); long start = System.currentTimeMillis(); - Map sparkConfMap = sparkConfToMap(getSparkConf()); + Map sparkConfMap = RssSparkConfig.sparkConfToMap(getSparkConf()); serverToPartitionRanges.entrySet().stream() .forEach( entry -> { @@ -1095,7 +1093,7 @@ protected void registerShuffleServers( } LOG.info("Start to register shuffleId[{}]", shuffleId); long start = System.currentTimeMillis(); - Map sparkConfMap = sparkConfToMap(getSparkConf()); + Map sparkConfMap = RssSparkConfig.sparkConfToMap(getSparkConf()); Set>> entries = serverToPartitionRanges.entrySet(); entries.stream() @@ -1141,15 +1139,4 @@ public boolean isRssStageRetryForFetchFailureEnabled() { public SparkConf getSparkConf() { return sparkConf; } - - public Map sparkConfToMap(SparkConf sparkConf) { - Map map = new HashMap<>(); - - for (Tuple2 tuple : sparkConf.getAll()) { - String key = tuple._1; - map.put(key, tuple._2); - } - - return map; - } } diff --git a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java index bb8ed3a901..792c8744f2 100644 --- a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java +++ b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java @@ -123,6 +123,7 @@ private boolean tryAccessCluster() { Map extraProperties = Maps.newHashMap(); extraProperties.put( ACCESS_INFO_REQUIRED_SHUFFLE_NODES_NUM, String.valueOf(assignmentShuffleNodesNum)); + extraProperties.putAll(RssSparkConfig.sparkConfToMap(sparkConf)); Set assignmentTags = RssSparkShuffleUtils.getAssignmentTags(sparkConf); try { From f538b5696664649fe563577e674a5ab81490fdad Mon Sep 17 00:00:00 2001 From: baoloongmao Date: Wed, 20 Nov 2024 23:16:21 +0800 Subject: [PATCH 2/7] Address suggestion, remove the spark prefix to be common --- .../apache/spark/shuffle/DelegationRssShuffleManager.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java index 792c8744f2..7cf8701aaf 100644 --- a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java +++ b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java @@ -123,7 +123,11 @@ private boolean tryAccessCluster() { Map extraProperties = Maps.newHashMap(); extraProperties.put( ACCESS_INFO_REQUIRED_SHUFFLE_NODES_NUM, String.valueOf(assignmentShuffleNodesNum)); - extraProperties.putAll(RssSparkConfig.sparkConfToMap(sparkConf)); + // Put all spark conf into extra properties, except which length is longer than 100 + // to avoid extra properties too long. + RssSparkConfig.toRssConf(sparkConf).getAll().stream() + .filter(entry -> StringUtils.length((String) entry.getValue()) < 100) + .forEach(entry -> extraProperties.put(entry.getKey(), (String) entry.getValue())); Set assignmentTags = RssSparkShuffleUtils.getAssignmentTags(sparkConf); try { From 9b016096e307bfe4672c59e88f18a03b6b4d4733 Mon Sep 17 00:00:00 2001 From: baoloongmao Date: Mon, 25 Nov 2024 14:35:28 +0800 Subject: [PATCH 3/7] Add a exclude properties config --- .../spark/shuffle/DelegationRssShuffleManager.java | 13 +++++++++++++ .../spark/shuffle/DelegationRssShuffleManager.java | 10 +++++++++- .../apache/uniffle/common/config/RssClientConf.java | 6 ++++++ 3 files changed, 28 insertions(+), 1 deletion(-) diff --git a/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java b/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java index af9295e559..602bc8088b 100644 --- a/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java +++ b/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java @@ -17,6 +17,7 @@ package org.apache.spark.shuffle; +import java.util.List; import java.util.Map; import java.util.Set; @@ -32,6 +33,8 @@ import org.apache.uniffle.client.impl.grpc.CoordinatorGrpcRetryableClient; import org.apache.uniffle.client.request.RssAccessClusterRequest; import org.apache.uniffle.client.response.RssAccessClusterResponse; +import org.apache.uniffle.common.config.RssClientConf; +import org.apache.uniffle.common.config.RssConf; import org.apache.uniffle.common.exception.RssException; import org.apache.uniffle.common.rpc.StatusCode; import org.apache.uniffle.common.util.Constants; @@ -124,6 +127,16 @@ private boolean tryAccessCluster() { extraProperties.put( ACCESS_INFO_REQUIRED_SHUFFLE_NODES_NUM, String.valueOf(assignmentShuffleNodesNum)); + RssConf rssConf = RssSparkConfig.toRssConf(sparkConf); + List excludeProperties = + rssConf.get(RssClientConf.RSS_CLIENT_REPORT_EXCLUDE_PROPERTIES); + // Put all spark conf into extra properties, except which length is longer than 100 + // to avoid extra properties too long. + rssConf.getAll().stream() + .filter(entry -> StringUtils.length((String) entry.getValue()) < 100) + .filter(entry -> !excludeProperties.contains(entry.getKey())) + .forEach(entry -> extraProperties.put(entry.getKey(), (String) entry.getValue())); + Set assignmentTags = RssSparkShuffleUtils.getAssignmentTags(sparkConf); try { if (coordinatorClient != null) { diff --git a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java index 7cf8701aaf..c136ab11ef 100644 --- a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java +++ b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java @@ -17,6 +17,7 @@ package org.apache.spark.shuffle; +import java.util.List; import java.util.Map; import java.util.Set; @@ -32,6 +33,8 @@ import org.apache.uniffle.client.impl.grpc.CoordinatorGrpcRetryableClient; import org.apache.uniffle.client.request.RssAccessClusterRequest; import org.apache.uniffle.client.response.RssAccessClusterResponse; +import org.apache.uniffle.common.config.RssClientConf; +import org.apache.uniffle.common.config.RssConf; import org.apache.uniffle.common.exception.RssException; import org.apache.uniffle.common.rpc.StatusCode; import org.apache.uniffle.common.util.Constants; @@ -123,10 +126,15 @@ private boolean tryAccessCluster() { Map extraProperties = Maps.newHashMap(); extraProperties.put( ACCESS_INFO_REQUIRED_SHUFFLE_NODES_NUM, String.valueOf(assignmentShuffleNodesNum)); + + RssConf rssConf = RssSparkConfig.toRssConf(sparkConf); + List excludeProperties = + rssConf.get(RssClientConf.RSS_CLIENT_REPORT_EXCLUDE_PROPERTIES); // Put all spark conf into extra properties, except which length is longer than 100 // to avoid extra properties too long. - RssSparkConfig.toRssConf(sparkConf).getAll().stream() + rssConf.getAll().stream() .filter(entry -> StringUtils.length((String) entry.getValue()) < 100) + .filter(entry -> !excludeProperties.contains(entry.getKey())) .forEach(entry -> extraProperties.put(entry.getKey(), (String) entry.getValue())); Set assignmentTags = RssSparkShuffleUtils.getAssignmentTags(sparkConf); diff --git a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java index 6d311a549a..5375a2fe5e 100644 --- a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java +++ b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java @@ -303,4 +303,10 @@ public class RssClientConf { .withDescription( "The block id manager class of server for this application, " + "the implementation of this interface to manage the shuffle block ids"); + public static final ConfigOption> RSS_CLIENT_REPORT_EXCLUDE_PROPERTIES = + ConfigOptions.key("rss.client.reportExcludeProperties") + .stringType() + .asList() + .noDefaultValue() + .withDescription("the report exclude properties could be configured by this option"); } From 1c93abfe25405872e875998b4a75657499a23bb3 Mon Sep 17 00:00:00 2001 From: baoloongmao Date: Mon, 25 Nov 2024 18:37:00 +0800 Subject: [PATCH 4/7] Set default value --- .../java/org/apache/uniffle/common/config/RssClientConf.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java index 5375a2fe5e..4d10db067a 100644 --- a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java +++ b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java @@ -307,6 +307,6 @@ public class RssClientConf { ConfigOptions.key("rss.client.reportExcludeProperties") .stringType() .asList() - .noDefaultValue() + .defaultValues("hdfs.inputPaths") .withDescription("the report exclude properties could be configured by this option"); } From 9a3b25dba5ac723017b33280bde1d7d9743dfab0 Mon Sep 17 00:00:00 2001 From: baoloongmao Date: Mon, 25 Nov 2024 19:12:15 +0800 Subject: [PATCH 5/7] Remove weired config value to keep empty list by default --- .../java/org/apache/uniffle/common/config/RssClientConf.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java index 4d10db067a..5b6d92cebc 100644 --- a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java +++ b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java @@ -307,6 +307,6 @@ public class RssClientConf { ConfigOptions.key("rss.client.reportExcludeProperties") .stringType() .asList() - .defaultValues("hdfs.inputPaths") + .defaultValues() .withDescription("the report exclude properties could be configured by this option"); } From 8d3f2ac3371aefbf82b2af7c826fa0251a07adff Mon Sep 17 00:00:00 2001 From: baoloongmao Date: Mon, 25 Nov 2024 20:06:39 +0800 Subject: [PATCH 6/7] Remove larger than 100 filter --- .../org/apache/spark/shuffle/DelegationRssShuffleManager.java | 3 --- .../org/apache/spark/shuffle/DelegationRssShuffleManager.java | 3 --- 2 files changed, 6 deletions(-) diff --git a/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java b/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java index 602bc8088b..027507a47a 100644 --- a/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java +++ b/client-spark/spark2/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java @@ -130,10 +130,7 @@ private boolean tryAccessCluster() { RssConf rssConf = RssSparkConfig.toRssConf(sparkConf); List excludeProperties = rssConf.get(RssClientConf.RSS_CLIENT_REPORT_EXCLUDE_PROPERTIES); - // Put all spark conf into extra properties, except which length is longer than 100 - // to avoid extra properties too long. rssConf.getAll().stream() - .filter(entry -> StringUtils.length((String) entry.getValue()) < 100) .filter(entry -> !excludeProperties.contains(entry.getKey())) .forEach(entry -> extraProperties.put(entry.getKey(), (String) entry.getValue())); diff --git a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java index c136ab11ef..ac1c9a0fd2 100644 --- a/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java +++ b/client-spark/spark3/src/main/java/org/apache/spark/shuffle/DelegationRssShuffleManager.java @@ -130,10 +130,7 @@ private boolean tryAccessCluster() { RssConf rssConf = RssSparkConfig.toRssConf(sparkConf); List excludeProperties = rssConf.get(RssClientConf.RSS_CLIENT_REPORT_EXCLUDE_PROPERTIES); - // Put all spark conf into extra properties, except which length is longer than 100 - // to avoid extra properties too long. rssConf.getAll().stream() - .filter(entry -> StringUtils.length((String) entry.getValue()) < 100) .filter(entry -> !excludeProperties.contains(entry.getKey())) .forEach(entry -> extraProperties.put(entry.getKey(), (String) entry.getValue())); From 6c16c8c6021de549dbb525c81c766bad3d0cd8d2 Mon Sep 17 00:00:00 2001 From: baoloongmao Date: Tue, 26 Nov 2024 14:32:17 +0800 Subject: [PATCH 7/7] Add a blank line --- .../java/org/apache/uniffle/common/config/RssClientConf.java | 1 + 1 file changed, 1 insertion(+) diff --git a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java index 5b6d92cebc..79b41afcc5 100644 --- a/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java +++ b/common/src/main/java/org/apache/uniffle/common/config/RssClientConf.java @@ -303,6 +303,7 @@ public class RssClientConf { .withDescription( "The block id manager class of server for this application, " + "the implementation of this interface to manage the shuffle block ids"); + public static final ConfigOption> RSS_CLIENT_REPORT_EXCLUDE_PROPERTIES = ConfigOptions.key("rss.client.reportExcludeProperties") .stringType()