Spark2.4.0 RDD依赖关系源码分析

本文假设我们已经对spark有一定的使用经验,并对spark常见对一些名词有一定的理解。

Spark DAG图中的Stage的划分依据是宽依赖和窄依赖,遇到宽依赖则会划分新的stage。今天主要分析下spark依赖的源码实现,spark依赖定义在org.apache.spark.Dependency中,如图所示:

image.png
类名 类型 说明
org.apache.spark.Dependency abstract 所有依赖的抽象基类
org.apache.spark.ShuffleDependency class 宽依赖实现类
org.apache.spark.NarrowDependency abstract 所有窄依赖的抽象基类
org.apache.spark.OneToOneDependency class 一对一的窄依赖实现类
org.apache.spark.RangeDependency class 有边界的范围窄依赖实现类

依赖基类 org.apache.spark.Dependency

所有的依赖均继承于该类,包括宽依赖(也称ShufferDependency)和窄依赖

/**
 * :: DeveloperApi ::
 * Base class for dependencies.
 */
@DeveloperApi
abstract class Dependency[T] extends Serializable {
  def rdd: RDD[T]
}

窄依赖 org.apache.spark.NarrowDependency

窄依赖(Narrow Dependency):如果子RDD的每个分区依赖于父RDD的少量分区,那么这种关系我们称为窄依赖。换言之,每一个父RDD中的Partition最多被子RDD的1个Partition所使用。

image.png

窄依赖分两种:一对一的依赖关系(OneToOneDependency)和范围依赖关系(RangeDependency)。
图中map,filter,join with inputs co-partitioned(对输入进行协同划分的join操作,典型的reduceByKey,先按照key分组然后shuffle write的时候一个父分区对应一个子分区)属于一对一的依赖关系;union属于范围依赖关系。

窄依赖基类

其中getParents方法定义了如何获取父RDD对应的partitionId

/**
 * :: DeveloperApi ::
 * Base class for dependencies where each partition of the child RDD depends on a small number
 * of partitions of the parent RDD. Narrow dependencies allow for pipelined execution.
 */
@DeveloperApi
abstract class NarrowDependency[T](_rdd: RDD[T]) extends Dependency[T] {
  /**
   * Get the parent partitions for a child partition.
   * @param partitionId a partition of the child RDD
   * @return the partitions of the parent RDD that the child partition depends upon
   */
  def getParents(partitionId: Int): Seq[Int]

  override def rdd: RDD[T] = _rdd
}

一对一的依赖 org.apache.spark.OneToOneDependency

直译就是一对一的依赖关系。在这种关系下,父RDD的partitionId和子RDD的partitionId相等。

/**
 * :: DeveloperApi ::
 * Represents a one-to-one dependency between partitions of the parent and child RDDs.
 */
@DeveloperApi
class OneToOneDependency[T](rdd: RDD[T]) extends NarrowDependency[T](rdd) {
  override def getParents(partitionId: Int): List[Int] = List(partitionId)
}

范围依赖关系 org.apache.spark.RangeDependency

有范围的依赖关系是指父RDD的partitionId和子RDD的partitionId在一定范围区间内一一对应。

/**
 * :: DeveloperApi ::
 * Represents a one-to-one dependency between ranges of partitions in the parent and child RDDs.
 * @param rdd the parent RDD
 * @param inStart the start of the range in the parent RDD
 * @param outStart the start of the range in the child RDD
 * @param length the length of the range
 */
@DeveloperApi
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
    }
  }
}

算法图解


image.png

宽依赖 org.apache.spark.ShuffleDependency

宽依赖(Wide Dependency | Shuffle Dependency):从英文名可以看出,宽依赖是存在shuffle的。换言之,一个父RDD的partition会被多个子RDD的partition所使用。

image.png

源码解读

/**
 * :: DeveloperApi ::
 * Represents a dependency on the output of a shuffle stage. Note that in the case of shuffle,
 * the RDD is transient since we don't need it on the executor side.
 *
 * @param _rdd the parent RDD
 * @param partitioner partitioner used to partition the shuffle output
 * @param serializer [[org.apache.spark.serializer.Serializer Serializer]] to use. If not set
 *                   explicitly then the default serializer, as specified by `spark.serializer`
 *                   config option, will be used.
 * @param keyOrdering key ordering for RDD's shuffles
 * @param aggregator map/reduce-side aggregator for RDD's shuffle
 * @param mapSideCombine whether to perform partial aggregation (also known as map-side combine)
 * @param shuffleWriterProcessor the processor to control the write behavior in ShuffleMapTask
 */
@DeveloperApi
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]] {

  if (mapSideCombine) {
    require(aggregator.isDefined, "Map-side combine without Aggregator specified!")
  }
  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
  // Note: It's possible that the combiner class tag is null, if the combineByKey
  // methods in PairRDDFunctions are used instead of combineByKeyWithClassTag.
  private[spark] val combinerClassName: Option[String] =
    Option(reflect.classTag[C]).map(_.runtimeClass.getName)

  val shuffleId: Int = _rdd.context.newShuffleId()

  val shuffleHandle: ShuffleHandle = _rdd.context.env.shuffleManager.registerShuffle(
    shuffleId, _rdd.partitions.length, this)

  _rdd.sparkContext.cleaner.foreach(_.registerShuffleForCleanup(this))
}

我们看到shuffle依赖会生成对应的shuffleId作为标示,另外,还会在shuffleManager中进行注册。

RDD依赖测试

测试方案:前面我们说groupByKey会产生shuffle生成宽依赖,所以groupByKey返回的rdd的依赖类型应该是ShuffleDependency。
测试代码

val rdd = sc.parallelize( 1 to 10, 5)
// map产生一对一窄依赖
val mapRdd = rdd.map(key => (key, 1))
// 窄依赖类
val narrowDependency = mapRdd.dependencies.head
// groupByKey产生宽依赖
val groupRdd = mapRdd.groupByKey
// 宽依赖类
val dependency = groupRdd.dependencies.head
val rdd2 = sc.parallelize( 10 to 100, 10)
// union产生范围窄依赖
val unionRdd = rdd.union(rdd2)
// 范围窄依赖
val rangeDependence = unionRdd.dependencies.head

测试结果符合预期

[hadoop@jms-master-01 ~]$ spark-shell
2019-03-19 23:46:25 WARN  NativeCodeLoader:62 - Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
Spark context Web UI available at http://jms-master-01:4040
Spark context available as 'sc' (master = local[*], app id = local-1553010403826).
Spark session available as 'spark'.
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 2.4.0
      /_/

Using Scala version 2.11.12 (Java HotSpot(TM) 64-Bit Server VM, Java 1.8.0_191)
Type in expressions to have them evaluated.
Type :help for more information.

scala> val rdd = sc.parallelize( 1 to 10, 5)
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:24

scala> val mapRdd = rdd.map(key => (key, 1))
mapRdd: org.apache.spark.rdd.RDD[(Int, Int)] = MapPartitionsRDD[1] at map at <console>:25

scala> val narrowDependency = mapRdd.dependencies.head
narrowDependency: org.apache.spark.Dependency[_] = org.apache.spark.OneToOneDependency@55882ff2

scala> val groupRdd = mapRdd.groupByKey
groupRdd: org.apache.spark.rdd.RDD[(Int, Iterable[Int])] = ShuffledRDD[2] at groupByKey at <console>:25

scala> val dependency = groupRdd.dependencies.head
dependency: org.apache.spark.Dependency[_] = org.apache.spark.ShuffleDependency@219aab91

scala> val rdd2 = sc.parallelize( 10 to 100, 10)
rdd2: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[3] at parallelize at <console>:24

scala> val unionRdd = rdd.union(rdd2)
unionRdd: org.apache.spark.rdd.RDD[Int] = UnionRDD[4] at union at <console>:27

scala> val rangeDependence = unionRdd.dependencies.head
rangeDependence: org.apache.spark.Dependency[_] = org.apache.spark.RangeDependency@5afefa87

scala>
最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 219,539评论 6 508
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 93,594评论 3 396
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 165,871评论 0 356
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,963评论 1 295
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,984评论 6 393
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,763评论 1 307
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,468评论 3 420
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 39,357评论 0 276
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,850评论 1 317
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 38,002评论 3 338
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 40,144评论 1 351
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,823评论 5 346
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 41,483评论 3 331
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 32,026评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 33,150评论 1 272
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 48,415评论 3 373
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 45,092评论 2 355