Spark RDD 宽窄依赖并行度源码解读

在上一节 Spark RDD任务并行度Part1:文件读取并行度源码解读 中,我们从源码出发知晓了Spark RDD在读取hdfs文本文件时的并行度生成逻辑,但这只是第一步,后续肯定是会对读取的数据进行处理的,那么后续计算的并行度是如何确定的?

本篇同样会通过解读源码来揭晓这个答案的。

Spark RDD间存在着依赖关系,正是由于不同的依赖导致数据并行度会有所不同。因此,首先要了解Spark RDD的依赖是如何定义的。

RDD依赖

Spark 使用Dependency 抽象类定义RDD依赖关系。

org/apache/spark/Dependency.scala

abstract class Dependency[T] extends Serializable {
  def rdd: RDD[T]
}

该类中只定义了一个方法rdd(),返回当前RDD依赖的父RDD。

RDD依赖分为窄依赖和宽依赖(Shuffle依赖)两种,实现类继承关系如下。

在这里插入图片描述

窄依赖

如果下游RDD与上游RDD分区是一对一关系,那么该RDD与其上游之间的依赖关系属于窄依赖。窄依赖的基类是NarrowDependency抽象类。

abstract class NarrowDependency[T](_rdd: RDD[T]) extends Dependency[T] {
  def getParents(partitionId: Int): Seq[Int]

  override def rdd: RDD[T] = _rdd
}

NarrowDependency类带有一个构造方法参数_rdd,并重写了rdd()方法,让其返回_rdd,它就是当前RDD依赖的父RDD。另外,它还定义了一个用来返回某一分区partitionId依赖的所有父RDD分区ID抽象方法getParents()。该方法由NarrowDependency的子类实现,分别为OneToOneDependency(一对一依赖)和RangeDependency(范围依赖)。

从源码可知这两个实现类子RDD的分区与依赖的分RDD分区都是一对一的关系,其中OneToOneDependency的父子RDD分区ID严格相同。

因此窄依赖的分区数是保持和上游一致的

class OneToOneDependency[T](rdd: RDD[T]) extends NarrowDependency[T](rdd) {
  override def getParents(partitionId: Int): List[Int] = List(partitionId)
}


class RangeDependency[T](rdd: RDD[T], inStart: Int, outStart: Int, length: Int)
  extends NarrowDependency[T](rdd) {

  override def getParents(partitionId: Int): List[Int] = {
    if (partitionId >= outStart && partitionId < outStart + length) {
      List(partitionId - outStart + inStart)
    } else {
      Nil
    }
  }
}

这个结论同样可从另一个角度得出。窄依赖算子创建的RDD如MapPartitionsRDD获取分区的方法partitions中,定义了窄依赖的RDD与上游依赖RDD的分区数是一致的。

org/apache/spark/rdd/MapPartitionsRDD.scala

final def partitions: Array[Partition] = {
    checkpointRDD.map(_.partitions).getOrElse {
        if (partitions_ == null) {
            stateLock.synchronized {
                if (partitions_ == null) {
                    partitions_ = getPartitions
                    partitions_.zipWithIndex.foreach { case (partition, index) =>
                        require(partition.index == index,
                                s"partitions($index).partition == ${partition.index}, but it should equal $index")
                    }
                }
            }
        }
        partitions_
    }
}

override def getPartitions: Array[Partition] = firstParent[T].partitions

org/apache/spark/rdd/RDD.scala


  protected[spark] def firstParent[U: ClassTag]: RDD[U] = {
    dependencies.head.rdd.asInstanceOf[RDD[U]]
  }
  
  final def dependencies: Seq[Dependency[_]] = {
    checkpointRDD.map(r => List(new OneToOneDependency(r))).getOrElse {
      if (dependencies_ == null) {
        stateLock.synchronized {
          if (dependencies_ == null) {
            dependencies_ = getDependencies
          }
        }
      }
      dependencies_
    }
  }

  protected def getDependencies: Seq[Dependency[_]] = deps

宽依赖(shuffle依赖)

如果下游RDD的一个分区会对应上游RDD的多个分区,这种依赖关系就是宽依赖(或者称为shuffle依赖),实现类是ShuffleDependency。

但下游RDD的各分区将具体依赖上游RDD的哪些分区?或者说上游RDD如何确定每个分区的数据分配到下游RDD的哪些分区呢?

这些都需要由ShuffleDependency类中的构造方法参数Partitioner类完成。

class ShuffleDependency[K: ClassTag, V: ClassTag, C: ClassTag](
    @transient private val _rdd: RDD[_ <: Product2[K, V]],
    val partitioner: Partitioner,
    val serializer: Serializer = SparkEnv.get.serializer,
    val keyOrdering: Option[Ordering[K]] = None,
    val aggregator: Option[Aggregator[K, V, C]] = None,
    val mapSideCombine: Boolean = false,
    val shuffleWriterProcessor: ShuffleWriteProcessor = new ShuffleWriteProcessor)
  extends Dependency[Product2[K, V]] {

  override def rdd: RDD[Product2[K, V]] = _rdd.asInstanceOf[RDD[Product2[K, V]]]
  private[spark] val keyClassName: String = reflect.classTag[K].runtimeClass.getName
  private[spark] val valueClassName: String = reflect.classTag[V].runtimeClass.getName
  ...
  val shuffleId: Int = _rdd.context.newShuffleId()
  val shuffleHandle: ShuffleHandle = _rdd.context.env.shuffleManager.registerShuffle(
    shuffleId, this)
	...
}

RDD分区计算器(Partitioner)

Partitioner类是一个抽象类,定义了分区计算器的规范。

org/apache/spark/Partitioner.scala

abstract class Partitioner extends Serializable {
  def numPartitions: Int
  def getPartition(key: Any): Int
}

numPartitions()方法用于获取下游分区总数;getPartitions()方法用于将输入的key映射到下游RDD分区ID,范围是0~numPartitions-1。

Partitioner实现类

Partitioner有很多官方实现类,生产中最常用的就是HashPartitioner类,该类定义了构造器参数partitions,用来作为分区数。

在重写的getPartition()方法中,会将key的hashCode值和分区数numPartitions进行取模运算,这样就将上游RDD数据映射下游RDD [0,numPartitions - 1]分区中了。

class HashPartitioner(partitions: Int) extends Partitioner {
  require(partitions >= 0, s"Number of partitions ($partitions) cannot be negative.")

  def numPartitions: Int = partitions

  def getPartition(key: Any): Int = key match {
    case null => 0
    case _ => Utils.nonNegativeMod(key.hashCode, numPartitions)
  }
   ...
}

defaultPartitioner()方法

Partitioner还带有一个伴生对象,定义了defaultPartitioner()方法,它是在没有显示设置分区器时返回默认的分区逻辑。

输入参数有:

rdd: RDD[_]:主 RDD(通常是调用 Shuffle 操作的 RDD)。

others: RDD[_]*:其他参与 Shuffle 的 RDD(如 join中的右 RDD)。

③ 返回值Partitioner:最终确定的统一分区器。

object Partitioner {
    
  def defaultPartitioner(rdd: RDD[_], others: RDD[_]*): Partitioner = {
    // 主 RDD 和其他 RDD 合并为一个序列 
    val rdds = (Seq(rdd) ++ others)
    // 从输入的所有RDD中筛选出已存在分区器且分区数大于0的 RDD
    val hasPartitioner = rdds.filter(_.partitioner.exists(_.numPartitions > 0))
	
    // 进一步找出分区数最多的 RDD(有效最大分区器)
    val hasMaxPartitioner: Option[RDD[_]] = if (hasPartitioner.nonEmpty) {
      Some(hasPartitioner.maxBy(_.partitions.length))
    } else {
      None
    }
	// 如果用户定义了spark.default.parallelism 参数,默认分区数就会采用该参数的值;否则使用所有输入RDD分区数的最大值
    val defaultNumPartitions = if (rdd.context.conf.contains("spark.default.parallelism")) {
      rdd.context.defaultParallelism
    } else {
      rdds.map(_.partitions.length).max
    }
	// 如果存在有效最大分区器,且(该最大分区器是“合格”的 or 其分区数大于默认分区数),就使用该最大分区器。
    if (hasMaxPartitioner.nonEmpty && (isEligiblePartitioner(hasMaxPartitioner.get, rdds) ||
        defaultNumPartitions <= hasMaxPartitioner.get.getNumPartitions)) {
      hasMaxPartitioner.get.partitioner.get
    } else {// 否则创建HashPartitioner,分区数=defaultNumPartitions
      new HashPartitioner(defaultNumPartitions)
    }
  }
    
    // 最大分区器是“合格”:该分区器的分区数与所有 RDD 最大分区数的对数差小于 1(即分区数处于同一数量级)。
  private def isEligiblePartitioner(
     hasMaxPartitioner: RDD[_],
     rdds: Seq[RDD[_]]): Boolean = {
    val maxPartitions = rdds.map(_.partitions.length).max
    log10(maxPartitions) - log10(hasMaxPartitioner.getNumPartitions) < 1
  }
}

RDD算子

通过上面的源码我们知道了RDD依赖分宽窄依赖,那各RDD算子是如何生成宽窄依赖的呢?

窄依赖算子

窄依赖算子有很多,如map、filter、flatMap、union 等等,下面以map算子为例,从源码角度看看依赖是如何生成的。

map()算子源码如下,创建了一个MapPartitionsRDD。

  def map[U: ClassTag](f: T => U): RDD[U] = withScope {
    val cleanF = sc.clean(f)
    new MapPartitionsRDD[U, T](this, (_, _, iter) => iter.map(cleanF))
  }

MapPartitionsRDD构造器参数prev是传入当前RDD,即上游RDD。MapPartitionsRDD继承了RDD类,在创建MapPartitionsRDD实例时会调用父类RDD的构造函数,并传入prev。

private[spark] class MapPartitionsRDD[U: ClassTag, T: ClassTag](
    var prev: RDD[T],
    f: (TaskContext, Int, Iterator[T]) => Iterator[U],  // (TaskContext, partition index, iterator)
    preservesPartitioning: Boolean = false,
    isFromBarrier: Boolean = false,
    isOrderSensitive: Boolean = false)
  extends RDD[U](prev)

在RDD中有一个辅助构造器只需要传入一个RDD类型的参数,因此在创建MapPartitionsRDD对象时会调用该辅助构造器,该辅助构造器又会调用RDD主构造器,设置RDD构造器参数deps=List(new OneToOneDependency(prev)),由此可知创建的MapPartitionsRDD与上游RDD是窄依赖关系。

abstract class RDD[T: ClassTag](
    @transient private var _sc: SparkContext,
    @transient private var deps: Seq[Dependency[_]]
  ) extends Serializable with Logging {
    ...
    
    def this(@transient oneParent: RDD[_]) =
    this(oneParent.context, List(new OneToOneDependency(oneParent)))
    ...
    protected def getDependencies: Seq[Dependency[_]] = deps
    ...
}

宽依赖算子

同样宽依赖算子有很多,如reduceByKey、sortByKey、join等等,下面以reduceByKey算子为例,从源码角度看看依赖是如何生成的。

提供了三种不同的reduceByKey重载方法,但最终都是会调用combineByKeyWithClassTag方法。

  // 形参:① 分区器,一般是自定义分区器时用;② 聚合规则
  def reduceByKey(partitioner: Partitioner, func: (V, V) => V): RDD[(K, V)] = self.withScope {
    combineByKeyWithClassTag[V]((v: V) => v, func, func, partitioner)
  }
 // 形参:① 分区数,分区规则默认使用HashPartitioner;② 聚合规则
  def reduceByKey(func: (V, V) => V, numPartitions: Int): RDD[(K, V)] = self.withScope {
    reduceByKey(new HashPartitioner(numPartitions), func)
  }
 // 形参:只有 聚合规则。使用默认分区器
  def reduceByKey(func: (V, V) => V): RDD[(K, V)] = self.withScope {
    reduceByKey(defaultPartitioner(self), func)
  }

进入到combineByKeyWithClassTag方法内部,发现最后会创建ShuffledRDD。

org/apache/spark/rdd/PairRDDFunctions.scala

new ShuffledRDD[K, V, C](self, partitioner)
    .setSerializer(serializer)
    .setAggregator(aggregator)
    .setMapSideCombine(mapSideCombine)
}

同样ShuffledRDD有一个构造器参数prev是传入当前RDD,即上游RDD。也继承了RDD类,在创建MapPartitionsRDD实例时会调用父类RDD的主构造函数,deps传的是空值。而是在调用getDependencies方法时会创建ShuffleDependency。

org/apache/spark/rdd/ShuffledRDD.scala

class ShuffledRDD[K: ClassTag, V: ClassTag, C: ClassTag](
    @transient var prev: RDD[_ <: Product2[K, V]],
    part: Partitioner)
  extends RDD[(K, C)](prev.context, Nil) {
    ...
    override def getDependencies: Seq[Dependency[_]] = {
    val serializer = userSpecifiedSerializer.getOrElse {
      val serializerManager = SparkEnv.get.serializerManager
      if (mapSideCombine) {
        serializerManager.getSerializer(implicitly[ClassTag[K]], implicitly[ClassTag[C]])
      } else {
        serializerManager.getSerializer(implicitly[ClassTag[K]], implicitly[ClassTag[V]])
      }
    }
    List(new ShuffleDependency(prev, part, serializer, keyOrdering, aggregator, mapSideCombine))
  }
      ...
}

测试案例验证

为验证上述源码解读逻辑的正确性,debug两个wordcount案例查看各核心变量取值是否符合预期。

文件准备

在hdfs路径下/input/spark_parallelism_test2上传两个文件:SogouQ2.txt 115M,smalltable.txt 12.4M。(受限于配置集群资源,设置hdfs块大小为32M)

在这里插入图片描述

wordcount案例开发参考:IDEA开发Spark应用

注:以下两个案例reduceByKey()算子都不传入分区数。

案例一:不设置spark.default.parallelism 参数

客户端提交提交应用:

/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" \
--class com.spark.core.WordCount \
./wc_args2.jar hdfs://hadoop100:8020/input/spark_parallelism_test2 hdfs://hadoop100:8020/output

在IDEA配置Remote JVM Debug。具体流程参考:IDEA远程调试Spark 源码与应用

窄依赖分区数

读取hdfs文件创建rdd1,使用Evaluate Expression功能,输入rdd1.partitions.length,查看rdd1的分区数=7;

在这里插入图片描述

在rdd1基础上执行flatMap算子生成rdd2,两者间是窄依赖关系,查看rdd2的分区数=7;

在这里插入图片描述

同样执行完map后,分区数还是=7。

在这里插入图片描述

宽依赖分区数

reduceByKey()方法没有传入分区数,会执行defaultPartitioner()。由于在之前没有执行过宽依赖算子,因此hasPartitioner为NULL,hasMaxPartitioner也会为NULL;没有设置spark.default.parallelism 参数,defaultNumPartitions=所依赖RDD(这里是rdd3)分区数的最大值(取值是7)。

在这里插入图片描述

同样可以通过调用result.partitions.length获取shuffle后的分区,结果也是7。

在这里插入图片描述

由于shuffle分区数是7,最终也生成了7个文件。

在这里插入图片描述

案例二:设置spark.default.parallelism 参数

客户端提交提交应用,设置spark.default.parallelism=4:

/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=4 \
--class com.spark.core.WordCount \
./wc_args2.jar hdfs://hadoop100:8020/input/spark_parallelism_test2 hdfs://hadoop100:8020/output2

窄依赖分区数

与上面案例一样,都是7。

在这里插入图片描述

宽依赖分区数

reduceByKey()方法没有传入分区数,会执行defaultPartitioner()。由于在之前没有执行过宽依赖算子,因此hasPartitioner为NULL,hasMaxPartitioner也会为NULL;设置了spark.default.parallelism 参数,defaultNumPartitions=spark.default.parallelism,取值为4。

在这里插入图片描述

调用result.partitions.length获取shuffle后的分区,结果也是4。

在这里插入图片描述

最终生成的文件个数也是4。

在这里插入图片描述

总结

通过上述源码解读和测试案例验证,我们知道了Spark RDD窄依赖分区数是不会发生变化;宽依赖分区数与用户是否手动设置、上游RDD分区数、spark.default.parallelism 参数等有关。

参考

耿嘉安. Spark内核设计的艺术:架构设计与实现[M]. 北京: 机械工业出版社, 2018

Spark Core源码精读计划19 | RDD的依赖与分区逻辑

总结

通过上述源码解读和测试案例验证,我们知道了Spark RDD窄依赖分区数是不会发生变化;宽依赖分区数与用户是否手动设置、上游RDD分区数、spark.default.parallelism 参数等有关。

更多推荐