RDD依赖关系
RDD 依赖关系:窄依赖快(无 Shuffle)、宽依赖慢(要 Shuffle),Stage 按宽依赖切
RDD 的依赖关系记录了一个 RDD 是怎么从父 RDD 算出来的。它干三件事:
| 作用 | 说明 |
|---|---|
| 容错 | 分区丢了,沿着依赖关系找到父 RDD,只重算丢失的分区 |
| Stage 划分 | 宽依赖是”切 Stage”的边界,窄依赖可以合并到一个 Stage |
| 性能优化 | 宽依赖 = Shuffle,能避免就避免 |
一句话:依赖关系 = 血缘 + 调度边界 + 容错依据
窄依赖 vs 宽依赖
| 窄依赖 | 宽依赖 | |
|---|---|---|
| 子分区依赖几个父分区 | 1 个(或少数几个) | 多个 |
| 有没有 Shuffle | 没有 | 有 |
| 能不能在同一 Stage | 能 | 不能(必须切 Stage) |
| 容错代价 | 低(只重算一个父分区) | 高(重算多个父分区) |
| 典型算子 | map、filter、flatMap |
reduceByKey、groupByKey、join |
窄依赖的例子: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
从底往上看:HadoopRDD → MapPartitionsRDD(窄) → 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(...) // 后续操作