RDD依赖关系

RDD 依赖关系:窄依赖快(无 Shuffle)、宽依赖慢(要 Shuffle),Stage 按宽依赖切

RDD 的依赖关系记录了一个 RDD 是怎么从父 RDD 算出来的。它干三件事:

作用 说明
容错 分区丢了,沿着依赖关系找到父 RDD,只重算丢失的分区
Stage 划分 宽依赖是”切 Stage”的边界,窄依赖可以合并到一个 Stage
性能优化 宽依赖 = Shuffle,能避免就避免

一句话:依赖关系 = 血缘 + 调度边界 + 容错依据

窄依赖 vs 宽依赖

窄依赖 宽依赖
子分区依赖几个父分区 1 个(或少数几个) 多个
有没有 Shuffle 没有
能不能在同一 Stage 不能(必须切 Stage)
容错代价 低(只重算一个父分区) 高(重算多个父分区)
典型算子 mapfilterflatMap reduceByKeygroupByKeyjoin

窄依赖的例子:map

父 RDD 有 2 个分区,map 之后子 RDD 也是 2 个分区,一一对应。

父分区0 → 子分区0
父分区1 → 子分区1

宽依赖的例子:reduceByKey

父 RDD 的多个分区的数据,按 key 重新分布后,可能落到同一个子分区。

父分区0 (a, b) ──┐
                 ├── 子分区0 (a)
父分区1 (a, c) ──┘
                 ├── 子分区1 (b, c)

宽依赖意味着数据要跨节点重新分布 → Shuffle。

怎么看依赖关系

val rdd = sc.textFile("data.txt")
  .flatMap(_.split(" "))
  .map((_, 1))
  .reduceByKey(_ + _)

// 看血缘(树形结构)
println(rdd.toDebugString)

输出:

(2) ShuffledRDD[3] at reduceByKey
 +-(2) MapPartitionsRDD[2] at map
    |  MapPartitionsRDD[1] at flatMap
    |  data.txt MapPartitionsRDD[0] at textFile
    |  data.txt HadoopRDD[0] at textFile

从底往上看:HadoopRDDMapPartitionsRDD(窄) → MapPartitionsRDD(窄) → ShuffledRDD(宽)。

reduceByKey 那里出现了 ShuffledRDD,说明这里切了 Stage。

Stage 怎么切

Stage 的切分规则:从后往前,遇到宽依赖就切。

操作链:textFile → flatMap → map → reduceByKey → collect

Stage 0(窄依赖链):textFile → flatMap → map
Stage 1(宽依赖):reduceByKey(触发 Shuffle)
Stage 2(行动):collect

Stage 0 里的三个算子在一个 Stage 里流水线执行,不用等、不用 Shuffle。Stage 1 要等 Stage 0 的数据全部算完,然后做 Shuffle。

窄依赖合并到同一个 Stage,宽依赖触发新 Stage。

容错:分区丢了怎么办

窄依赖:

子分区0 丢了,只需要重算父分区0,其他分区不受影响。

父分区0 → 子分区0  ← 丢了,只重算父分区0
父分区1 → 子分区1  ← 正常

宽依赖:

子分区0 丢了,可能需要重算多个父分区(所有贡献数据的父分区)。

父分区0 ──┐
          ├── 子分区0  ← 丢了,要重算父分区0、父分区1、父分区2...
父分区1 ──┤
          │
父分区2 ──┘

所以宽依赖的容错代价更高。

实战:怎么利用依赖关系优化代码

1. 减少宽依赖

//  宽依赖,所有数据 Shuffle
rdd.groupByKey().mapValues(_.sum)

//  窄依赖 + 预聚合,Shuffle 数据量小很多
rdd.reduceByKey(_ + _)

2. 对宽依赖后的 RDD 做缓存

宽依赖后的 RDD,如果被多个 Action 使用,缓存它避免重复 Shuffle。

val rdd = sc.textFile("data.txt")
  .flatMap(_.split(" "))
  .map((_, 1))
  .reduceByKey(_ + _)   // 宽依赖,触发 Shuffle
  .cache()              // ← 缓存,后面多个 Action 不用重复 Shuffle

rdd.count()
rdd.collect()

3. 长血缘链的 RDD 缓存

如果 RDD 的血缘链很长(几十个转换),缓存中间结果可以避免失败时从头重算。

val rdd = sc.textFile("data.txt")
  .flatMap(...)
  .map(...)
  .filter(...)
  .map(...)   // 长链
  .cache()    // ← 缓存中间结果
  .map(...)   // 后续操作