Spark RDD任务并行度Part2:宽窄依赖并行度源码解读
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 参数等有关。
更多推荐
所有评论(0)