Spark RDD任务并行度Part1:文件读取并行度源码解读
Spark RDD 文件读取并行度源码解读
生产中数据常存储在HDFS文件系统中,那么当Spark RDD在读取hdfs文件时的并行度是如何确定的呢?
接下来本文会以最常见的读取文本文件为例来讲解Spark RDD是如何划分分区的。
textFile()源码解读
读取文件文件常使用textFile方法。
org/apache/spark/SparkContext.scala
def textFile(
path: String,
minPartitions: Int = defaultMinPartitions): RDD[String] = withScope {
assertNotStopped()
hadoopFile(path, classOf[TextInputFormat], classOf[LongWritable], classOf[Text],
minPartitions).map(pair => pair._2.toString).setName(path)
}
由上述textFile源码可知该方法要传入两个参数:① 文件路径;② minPartitions,最小分区数,如果不传入值默认使用defaultMinPartitions方法的返回值。
方法体中是调用了hadoopFile方法。下面将深入探究defaultMinPartitions方法和hadoopFile方法。
defaultMinPartitions()
1、进入defaultMinPartitions方法:
org/apache/spark/SparkContext.scala
def defaultMinPartitions: Int = math.min(defaultParallelism, 2)
2、进入defaultParallelism方法:
org/apache/spark/scheduler/TaskSchedulerImpl.scala
def defaultParallelism: Int = {
assertNotStopped()
taskScheduler.defaultParallelism
}
3、继续进入defaultParallelism方法,如果集群模式(yarn集群或standlone模式),taskScheduler的实现类是CoarseGrainedSchedulerBackend
org/apache/spark/scheduler/cluster/CoarseGrainedSchedulerBackend.scala
override def defaultParallelism(): Int = {
conf.getInt("spark.default.parallelism", math.max(totalCoreCount.get(), 2))
}
// 从配置中获取名为 spark.default.parallelism 的整数值;如果未设置该配置项,则返回 math.max(totalCoreCount.get(), 2)。totalCoreCount,表示当前集群中所有 Executor 可用的总 CPU 核心数
从上述逻辑可知,当没有在textFile方法中手动指定minPartitions值时,minPartitions=min(defaultParallelism, 2),defaultParallelism的值分两种情况:
① 设置spark.default.parallelism参数时,defaultParallelism=spark.default.parallelism;
② 没有设置spark.default.parallelism参数时,defaultParallelism=math.max(totalCoreCount.get(), 2);
也就是说除非设置了spark.default.parallelism=1,否则 minPartitions=2。
侧面说明了
spark.default.parallelism参数并不会在读取文件过程中决定其分区数,而是会被2“硬限制”给覆盖了。
hadoopFile()
hadoopFile方法的关键形参有:
① inputFormatClass:读取Hadoop 文件的InputFormat类;(从上述textFile可知,传入的inputFormatClass是TextInputFormat类,因此textFile读取文件是按行读取)
注:Spark可以无缝集成Hadoop的 InputFormat,这使得它可以读取任何Hadoop支持的数据格式(如SequenceFile, Avro, Parquet等),只需要使用相应的InputFormat来读取文件。
② keyClass:键(Key)的 Class 类型;
③ valueClass:值(Value)的 Class 类型;
④ minPartitions:最小分区数。
方法体中会创建HadoopRDD类。
org/apache/spark/SparkContext.scala
// withScope:将操作包装在 Spark 的监控作用域中
def hadoopFile[K, V](
path: String,
inputFormatClass: Class[_ <: InputFormat[K, V]],
keyClass: Class[K],
valueClass: Class[V],
minPartitions: Int = defaultMinPartitions): RDD[(K, V)] = withScope {
assertNotStopped()
// This is a hack to enforce loading hdfs-site.xml.
// See SPARK-11227 for details.
FileSystem.getLocal(hadoopConfiguration)
// A Hadoop configuration can be about 10 KiB, which is pretty big, so broadcast it.
val confBroadcast = broadcast(new SerializableConfiguration(hadoopConfiguration))
val setInputPathsFunc = (jobConf: JobConf) => FileInputFormat.setInputPaths(jobConf, path)
new HadoopRDD(
this,
confBroadcast,
Some(setInputPathsFunc),
inputFormatClass,
keyClass,
valueClass,
minPartitions).setName(path)
}
1、HadoopRDD类
HadoopRDD类中含有一个getPartitions方法,是重写RDD类中的getPartitions方法,该方法在 DAGScheduler创建 Stage 时会被调用,用来确定分区数量,即并行度,关键逻辑是在getSplits方法中。
getSplits方法是InputFormat接口的方法,创建HadoopRDD传入的InputFormat类默认是TextInputFormat类。因此进入到TextInputFormat类的getSplits方法中就可以看到分区数量如何确定了。
org/apache/spark/rdd/HadoopRDD.scala
class HadoopRDD[K, V](
sc: SparkContext,
broadcastedConf: Broadcast[SerializableConfiguration],
initLocalJobConfFuncOpt: Option[JobConf => Unit],
inputFormatClass: Class[_ <: InputFormat[K, V]],
keyClass: Class[K],
valueClass: Class[V],
minPartitions: Int)
extends RDD[(K, V)](sc, Nil) with Logging {
...
override def getPartitions: Array[Partition] = {
...
val allInputSplits = getInputFormat(jobConf).getSplits(jobConf, minPartitions) // 获取文件分片.step into
val inputSplits = if (ignoreEmptySplits) {
allInputSplits.filter(_.getLength > 0)
} else {
allInputSplits
}
...
val array = new Array[Partition](inputSplits.size)
for (i <- 0 until inputSplits.size) {
array(i) = new HadoopPartition(id, i, inputSplits(i))
}
array
}
...
}
2、FileInputFormat类
TextInputFormat继承FileInputFormat类,getSplits方法在FileInputFormat中定义,进入getSplits中就能看到切割分片的核心逻辑了。
注意:MR读取hdfs文件使用的FileInputFormat全路径是
org/apache/hadoop/mapreduce/lib/input/FileInputFormat.java。参考:
MapReduce 任务并行度源码解析
org.apache.hadoop.mapred.FileInputFormat.java
public InputSplit[] getSplits(JobConf job, int numSplits)
throws IOException {
FileStatus[] stats = listStatus(job); // 获取输入路径下的所有文件FileStatus
// Save the number of input files for metrics/loadgen
job.setLong(NUM_INPUT_FILES, stats.length);
long totalSize = 0; // compute total size
List<FileStatus> files = new ArrayList<>(stats.length);
for (FileStatus file: stats) { // check we have valid files
if (file.isDirectory()) {
if (!ignoreDirs) {
throw new IOException("Not a file: "+ file.getPath());
}
} else {
files.add(file);
totalSize += file.getLen(); // 循环遍历每一个文件,获取所有文件的大小
}
}
long goalSize = totalSize / (numSplits == 0 ? 1 : numSplits); // numSplits一般默认=2
long minSize = Math.max(job.getLong(org.apache.hadoop.mapreduce.lib.input.
FileInputFormat.SPLIT_MINSIZE, 1), minSplitSize); // 获取最小切分大小
// generate splits
ArrayList<FileSplit> splits = new ArrayList<FileSplit>(numSplits);
NetworkTopology clusterMap = new NetworkTopology();
for (FileStatus file: files) { // 遍历每一个文件
Path path = file.getPath();
long length = file.getLen();
if (length != 0) {
FileSystem fs = path.getFileSystem(job);
BlockLocation[] blkLocations;
if (file instanceof LocatedFileStatus) {
blkLocations = ((LocatedFileStatus) file).getBlockLocations();
} else {
blkLocations = fs.getFileBlockLocations(file, 0, length);
}
if (isSplitable(fs, path)) { // 判断是否可分割
long blockSize = file.getBlockSize();// 获取hdfs块大小
long splitSize = computeSplitSize(goalSize, minSize, blockSize); // 计算分片大小。核心
// 根据分片大小对单个文件进行切割,同时会判断文件剩下未切的大小是否大于切片大小的1.1倍,不大于1.1倍就只划分一块切片,大于1.1倍就继续切割。
long bytesRemaining = length;
while (((double) bytesRemaining)/splitSize > SPLIT_SLOP) {
String[][] splitHosts = getSplitHostsAndCachedHosts(blkLocations,
length-bytesRemaining, splitSize, clusterMap);
splits.add(makeSplit(path, length-bytesRemaining, splitSize,
splitHosts[0], splitHosts[1]));
bytesRemaining -= splitSize;
}
if (bytesRemaining != 0) {
String[][] splitHosts = getSplitHostsAndCachedHosts(blkLocations, length
- bytesRemaining, bytesRemaining, clusterMap);
splits.add(makeSplit(path, length - bytesRemaining, bytesRemaining,
splitHosts[0], splitHosts[1]));
}
} else {
if (LOG.isDebugEnabled()) {
// Log only if the file is big enough to be splitted
if (length > Math.min(file.getBlockSize(), minSize)) {
LOG.debug("File is not splittable so no parallelization "
+ "is possible: " + file.getPath());
}
}
String[][] splitHosts = getSplitHostsAndCachedHosts(blkLocations,0,length,clusterMap);
splits.add(makeSplit(path, 0, length, splitHosts[0], splitHosts[1]));// 文件不可切割时,整个文件作为一个分片
}
} else {
//Create empty hosts array for zero length files
splits.add(makeSplit(path, 0, length, new String[0]));
}
}
return splits.toArray(new FileSplit[splits.size()]);
}
从上述源码中可知,切片大小splitSize逻辑如下:
long splitSize = computeSplitSize(goalSize, minSize, blockSize);
protected long computeSplitSize(long goalSize, long minSize,long blockSize) {
return Math.max(minSize, Math.min(goalSize, blockSize));
}
① goalSize:目标文件大小
用所有文件大小/传入的numSplits,从上述的调用链中numSplits取值分为两种场景:
① 如果最初使用textFile()方法传入了值,则取该值;
② 如果没有传入,除非设置了spark.default.parallelism=1,否则为2。
long goalSize = totalSize / (numSplits == 0 ? 1 : numSplits);
② minSize:最小切片大小
long minSize = Math.max(job.getLong(org.apache.hadoop.mapreduce.lib.input.FileInputFormat.SPLIT_MINSIZE, 1), minSplitSize); // 获取最小切分大小
其中SPLIT_MINSIZE = "mapreduce.input.fileinputformat.split.minsize"
minSplitSize如果没有通过setMinSplitSize设置,默认为1。
private long minSplitSize = 1
protected void setMinSplitSize(long minSplitSize) {this.minSplitSize = minSplitSize;}
③ blockSize:hdfs集群存储块大小
结论
综上,在实际生产中,数据基本都是支持可切割的且数据量往往很大,计算出的goalSize往往是比blockSize大,因此一般情况下splitSize=blockSize。
Spark读取hdfs并行度测试
为验证上述源码解读逻辑的正确性,用两个例子debug查看各核心变量取值是否符合预期。
文件准备
在hdfs路径下/input/spark_parallelism_test2上传两个文件:SogouQ2.txt 115M,smalltable.txt 12.4M。(受限于配置集群资源,设置hdfs块大小为32M)

案例一:text()方法不传入minPartitions参数
此测试案例text()方法不传入参数,在提交Application时设置spark.default.parallelism参数。
val rdd1: RDD[String] = sc.textFile(args(0))
开发程序并开启远程调试
① 开发Spark程序
用经典的wordcount案例进行调试,具体开发流程参考:IDEA开发Spark应用
② 开启远程调试
在Spark客户端提交WordCount,设置spark.default.parallelism=8:
/opt/module/spark-3.1.2/bin/spark-submit \
--master yarn \
--deploy-mode client \
--conf "spark.driver.extraJavaOptions=-agentlib:jdwp=transport=dt_socket,server=y,suspend=y,address=5005" \
--conf spark.default.parallelism=8 \
--class com.spark.core.WordCount \
./wc_args.jar hdfs://hadoop100:8020/input/spark_parallelism_test2 hdfs://hadoop100:8020/output
在IDEA配置Remote JVM Debug。

具体流程参考:IDEA远程调试Spark 源码与应用
关键参数
1、 minPartitions
即使设置了spark.default.parallelism=8,也不起作用,minPartitions还是=2

2、计算totalSize
获取提交路径下的文件列表、每个文件的大小、块位置(block locations)、修改时间等等,统计所有文件的大小totalSize。

3、计算goalSize和minSize
① numSplits=2,goalSize是totalSize的一半
② 没有设置 mapreduce.input.fileinputformat.split.minsize参数也没有调用setMinSplitSize()设置minSplitSize,minSplitSize取默认值1。

4、计算切片大小splitSize
splitSize=Math.max(1, Math.min(66832092, 33554432)),得到splitSize=33554432,与blockSize相等。

5、切片文件
根据splitSize对文件切片,第一个文件SogouQ2被切割成4个分片,第二个文件smalltable.txt单独形成一个分片,符合预期。

打开 Spark UI,可以看到读取HDFS文件启动了5个task,与上述切片个数相同,且各task读取文件大小与上述切片文件大小近似相等。


案例二:text()方法传入minPartitions参数
在程序中textFile传入minPartitions=6;
val rdd1: RDD[String] = sc.textFile(args(0),6)
将程序打包命名为wc_args2.jar,提交应用,开启远程调试
/opt/module/spark-3.1.2/bin/spark-submit \
--master yarn \
--deploy-mode client \
--conf "spark.driver.extraJavaOptions=-agentlib:jdwp=transport=dt_socket,server=y,suspend=y,address=5005" \
--conf spark.default.parallelism=8 \
--class com.spark.core.WordCount \
./wc_args2.jar hdfs://hadoop100:8020/input/spark_parallelism_test2 hdfs://hadoop100:8020/output2
查看各关键变量:
1、 minPartitions
设置了spark.default.parallelism=8,依旧不起作用,minPartitions=6。

2、计算goalSize和minSize
① numSplits=6,goalSize=totalSize/6=22277354;
② 没有设置 mapreduce.input.fileinputformat.split.minsize参数也没有调用setMinSplitSize()设置minSplitSize,minSplitSize取默认值1。

3、计算切片大小splitSize
splitSize=Math.max(1, Math.min(22277354, 33554432)),得到splitSize=22277354,此时与blockSize不等。

4、切片文件
根据splitSize对文件切片,第一个文件SogouQ2被切割成6个分片,第二个文件smalltable.txt单独形成一个分片,符合预期。

打开 Spark UI,可以看到读取HDFS文件启动了7个task,与上述切片个数相同,且各task读取文件大小与上述切片文件大小近似相等。

更多推荐
所有评论(0)