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读取文件大小与上述切片文件大小近似相等。

在这里插入图片描述

更多推荐