乐于分享
好东西不私藏

【极速学习spark源码】RDD 五大特性的源码体现

【极速学习spark源码】RDD 五大特性的源码体现

基于 Spark 4.2,分支 branch-4.2

阅读时长约 10 分钟 · 入门到中级


背景

Spark 官方文档说 RDD 有五个特性(Five Properties), 面试也经常考。

这个文章做一件事:打开 RDD.scala,把五个特性逐个指到具体的代码行。


一、RDD 注释里的原文

打开 RDD.scala:71-75,Spark 自己写了五个特性:

 *  - A list of partitions
 *  - A function for computing each split
 *  - A list of dependencies on other RDDs
 *  - Optionally, a Partitioner for key-value RDDs (e.g. to say that the RDD is hash-partitioned)
 *  - Optionally, a list of preferred locations to compute each split on (e.g. block locations for
 *    an HDFS file)

对应的抽象类定义在 RDD.scala:84

abstract class RDD[TClassTag](    @transient private var scSparkContext,    @transient private var depsSeq[Dependency[_]]  ) extends Serializable with Logging

下面逐个拆开。


二、特性 1:一组 Partitions

含义:一个 RDD 由一组 Partition 组成,每个 Partition 是数据的逻辑分片。

源码体现

  • RDD.scala:125:抽象方法 getPartitions

    protected def getPartitionsArray[Partition]
  • RDD.scala:296:final 方法 partitions,带 checkpoint 逻辑

    final def partitionsArray[Partition] = {  // checkpoint or lazily compute  getPartitions}
  • Partition.scala:23:Partition trait

    trait Partition extends Serializable {  def index: Int}

Partition 本身不存数据。它只是一个标识——告诉 Spark "第 N 个分片在这里"。

子类实现举例

  • ParallelCollectionRDD.scala:96getPartitions 把本地集合切成多个 slice

  • MapPartitionsRDD.scala:54getPartitions 直接返回父 RDD 的 partitions(窄转换不改分区)

  • ShuffledRDD.scala:92getPartitions 根据 partitioner.numPartitions 生成新分区


三、特性 2:每个 Partition 的计算函数

含义:给定一个 Partition,RDD 知道怎么算出它的数据。

源码体现

  • RDD.scala:116:抽象方法 compute

    def compute(split: Partition, context: TaskContext): Iterator[T]
  • RDD.scala:334:final 方法 iterator,统一处理 cache/checkpoint

    final def iterator(split: Partition, context: TaskContext): Iterator[T] = {  if (storageLevel != StorageLevel.NONE) {    getOrCompute(split, context)  } else {    computeOrReadCheckpoint(split, context)  }}

iterator 是 Task 执行时实际调用的入口(见上一篇文章 ResultTask.runTask 里的 rdd.iterator(partition, context))。它先查缓存,没有缓存才调 compute

子类实现举例

  • ParallelCollectionRDD.scala:101compute 直接从 Partition 的 iterator 读数据

  • MapPartitionsRDD.scala:56compute 用函数 f 包装父 RDD 的 iterator

    override def compute(split: Partition, context: TaskContext): Iterator[U] =  f(context, split.index, firstParent[T].iterator(split, context))
  • ShuffledRDD.scala:102compute 通过 ShuffleManager reader 读取 shuffle 输出


四、特性 3:对其他 RDD 的依赖列表

含义:RDD 知道自己依赖谁。这个依赖列表构成了 RDD 血缘(lineage),也是容错和 Stage 划分的基础。

源码体现

  • RDD.scala:131getDependencies,默认返回构造参数 deps

    protected def getDependencies: Seq[Dependency[_]] = deps
  • RDD.scala:260:final 方法 dependencies,带 checkpoint 逻辑

    final def dependencies: Seq[Dependency[_]] = {  // checkpoint or lazily compute  getDependencies}
  • Dependency.scala:41Dependency 抽象类

    abstract classDependency[TextendsSerializable{  def rdd: RDD[T]}

每个 Dependency 对象指向一个父 RDD,就是血缘图中的一条边。

窄依赖和宽依赖

  • Dependency.scala:52NarrowDependency,子 partition 只依赖少量父 partition

  • Dependency.scala:58getParents(partitionId) 描述子到父的映射

  • Dependency.scala:84ShuffleDependency,宽依赖,需要 shuffle

  • Dependency.scala:266OneToOneDependency,最典型的窄依赖

子类实现举例

  • MapPartitionsRDD.scala:50extends RDD[U](prev) —— 用单父构造器,自动产生 OneToOneDependencyRDD.scala:104

  • ShuffledRDD.scala:78getDependencies 返回 List(new ShuffleDependency(...))

  • ParallelCollectionRDD.scala:90extends RDD[T](sc, Nil) —— 没有父依赖,deps = Nil


五、特性 4:可选的 Partitioner

含义:对于 key-value 类型的 RDD,可以声明分区策略(hash / range 等)。这影响 shuffle 输出的分区方式,也影响 join 等操作的优化。

源码体现

  • RDD.scala:139partitioner 字段

    @transient val partitioner: Option[Partitioner] = None

大多数 RDD 的 partitioner 是 None。只有经过 shuffle 或显式 repartition 的 RDD 才有。

子类实现举例

  • ShuffledRDD.scala:90

    override val partitioner = Some(part)
  • MapPartitionsRDD.scala:52:如果 preservesPartitioning 为 true,继承父 partitioner

    override val partitioner = if (preservesPartitioning) firstParent[T].partitioner else None

ShuffleDependency 也持有 partitioner(Dependency.scala:86),保证 shuffle 读写双方的分区策略一致。


六、特性 5:可选的 Preferred Locations

含义:RDD 可以建议"某个 Partition 最好在哪个机器上算"。这是数据本地性的基础。

源码体现

  • RDD.scala:136getPreferredLocations,默认返回空

    protected def getPreferredLocations(split: Partition): Seq[String] = Nil
  • RDD.scala:323:final 方法 preferredLocations

    final def preferredLocations(split: Partition): Seq[String] = {  getPreferredLocations(split)}

子类实现举例

  • ShuffledRDD.scala:96getPreferredLocations 通过 MapOutputTracker 查询 shuffle 数据块的位置

    override def getPreferredLocations(partition: Partition): Seq[String] = {  val tracker = SparkEnv.get.mapOutputTracker.asInstanceOf[MapOutputTrackerMaster]  tracker.getPreferredLocationsForShuffle(part, partition.index)}
  • ParallelCollectionRDD.scala:105:从 locationPrefs 读取位置偏好

数据本地性的调度策略在 TaskSetManager 里使用:优先分配 Task 到 PROCESS_LOCAL,其次 NODE_LOCAL,依次降级到 ANY


七、五大特性一览表

#特性RDD.scala 位置抽象/默认子类典型实现
1PartitionsgetPartitions:125 / partitions:296抽象ParallelCollectionRDD 切片;MapPartitionsRDD 继承父
2Computecompute:116 / iterator:334抽象MapPartitionsRDD 包装函数;ShuffledRDD 读 shuffle
3DependenciesgetDependencies:131 / dependencies:260默认返回构造参数ShuffledRDD 返回 ShuffleDependency
4Partitionerpartitioner:139默认 NoneShuffledRDD 声明 partitioner
5Preferred LocationsgetPreferredLocations:136 / preferredLocations:323默认 NilShuffledRDD 查 MapOutputTracker

八、三个子类对比

子类PartitionsComputeDependenciesPartitionerPreferred Locations
ParallelCollectionRDD切本地集合 (:96)从 Partition 读 (:101)Nil(无父 RDD)(:90)NonelocationPrefs (:105)
MapPartitionsRDD继承父 (:54)f(父.iterator) (:56)OneToOneDependency(自动)(:50)继承父或 None (:52)默认 Nil
ShuffledRDDpartitioner 决定 (:92)读 shuffle (:102)ShuffleDependency (:78)声明 partitioner (:90)MapOutputTracker (:96)

这三个子类几乎覆盖了所有场景:

  • 源头 RDD(如 ParallelCollectionRDD / HadoopRDD):自己造分区、自己算数据

  • 窄转换 RDD(如 MapPartitionsRDD):直接继承分区和依赖,只在 compute 里加一层函数

  • 宽转换 RDD(如 ShuffledRDD):创建新的分区、新的依赖、新的 partitioner


九、实战:用源码验证五大特性

object TestRDD {  def main(args: Array[String]): Unit = {    val spark = SparkSession      .builder()      .appName("RDD-DEMO")      .master("local[4]")      .config("spark.serializer""org.apache.spark.serializer.JavaSerializer")      .getOrCreate();    val sc = spark.sparkContext    sc.setLogLevel("WARN")    val rdd1 = sc.parallelize(1 to 1004)    val rdd2 = rdd1.map(_ * 2)    val pairs = sc.parallelize(Seq(("a"1), ("b"2), ("a"3)), 2)    val rdd3 = pairs.reduceByKey(_ + _)    // 特性 1: partitions    println(s"rdd1 partitions: ${rdd1.getNumPartitions}")    println(s"rdd3 partitions: ${rdd3.getNumPartitions}")    // 特性 3: dependencies    println(s"rdd1 deps: ${rdd1.dependencies.map(_.getClass.getSimpleName)}")    println(s"rdd2 deps: ${rdd2.dependencies.map(_.getClass.getSimpleName)}")    println(s"rdd3 deps: ${rdd3.dependencies.map(_.getClass.getSimpleName)}")    // 特性 4: partitioner    println(s"rdd1 partitioner: ${rdd1.partitioner}")    println(s"rdd3 partitioner: ${rdd3.partitioner}")    // 特性 5: preferredLocations    println(s"rdd1 preferred: ${rdd1.preferredLocations(rdd1.partitions(0))}")  }}

输出:

rdd1 partitions: 4
rdd3 partitions: 2
rdd1 deps: List()
rdd2 deps: List(OneToOneDependency)
rdd3 deps: List(ShuffleDependency)
rdd1 partitioner: None
rdd3 partitioner: Some(org.apache.spark.HashPartitioner@2)
rdd1 preferred: List()

十、总结

RDD 的五大特性不是面试八股,而是 Spark 调度系统的地基:

  • Partitions 决定了并行度

  • Compute 决定了每个 Task 做什么

  • Dependencies 决定了 Stage 怎么切、数据怎么传

  • Partitioner 决定了 shuffle 输出怎么分

  • Preferred Locations 决定了 Task 往哪台机器放

从源码角度看,RDD 的设计非常统一:五个特性对应五个可覆盖的方法/字段,子类只需选择性地覆盖自己关心的部分。


每天花费5分钟学习spark,让你技术之路走得更稳、更快。

喜欢的点个关注。