分布式系统(30):Spark/RDD 的血缘、物化与故障重算
迭代计算有两种基本恢复思路:每一轮都把完整中间状态写到可靠存储,或者保留“这份数据怎样算出来”的配方,丢失后只重算需要的分区。MapReduce 倾向前者,Spark RDD 把后者做成了核心抽象。差别不是“磁盘与内存”的口号,而是恢复证据放在哪里、何时必须物化。
分布式系统(29):Kafka 的日志、消费位点与事务边界第 29 篇用日志位点保存重放起点。RDD 没有用一个消费 offset 表示恢复位置,而是保存分区依赖图;shuffle、cache 和 checkpoint 又会在图上形成性质不同的边界。本篇用同一个四分区迭代负载比较这些边界。
RDD 保存的是分区与配方
RDD 论文把 RDD 定义为分区的只读记录集合。程序从可靠输入创建 RDD,再用 map、filter、join、reduceByKey 等转换得到新的 RDD。转换是惰性的;action 需要结果时,调度器才沿依赖图生成任务。
flowchart LR
I[可靠输入] --> M[map]
M --> F[filter]
F --> S[reduceByKey / shuffle]
S --> A[action]
C[(cache)] -.可丢失.-> F
K[(checkpoint)] -.截短血缘.-> S
“只读”很关键。某个分区不是被原地修补的共享内存,而是一个确定转换的结果。只要输入和转换仍可用,丢失的分区就能重新生成。lineage 是恢复配方,不是数据副本;缓存命中让执行更快,缓存丢失不应改变结果。
这套模型要求计算可重放。转换若读取不断变化的外部状态、依赖任务尝试次数,或在执行时修改外部系统,重算就可能得到不同值或产生重复效果。RDD 的容错结论首先覆盖 RDD 数据,不自动覆盖任意用户代码。
窄依赖与宽依赖决定重算扇出
论文按父分区如何被子分区使用来区分依赖。窄依赖中,一个父分区至多被一个子分区使用;map、filter 常可让一条流水线在同一分区内执行。宽依赖中,一个父分区的数据会按 key 分散给多个子分区,groupByKey、重新分区等操作通常需要 shuffle。
flowchart TB
subgraph Narrow[窄依赖]
P0[父0] --> C0[子0]
P1[父1] --> C1[子1]
P2[父2] --> C2[子2]
end
subgraph Wide[宽依赖 / shuffle]
A0[父0] --> B0[子0]
A0 --> B1[子1]
A1[父1] --> B0
A1 --> B1
end
窄依赖的丢失通常只向上追一个小分支。本地模型删除 map 后的分区 2,恢复只读取父分区 2。宽依赖的某个 reduce 分区则需要所有相关 map 分区的 bucket;如果这些 shuffle 文件也已丢失,就要重新执行相应 map task。
这不是“遇到 shuffle 就从头跑”的规则。Spark 会保留 shuffle 文件,当前 4.2.0 文档还列出外部 shuffle 服务、shuffle tracking、executor decommission 搬移和可靠 ShuffleDataIO 等选项。需要重算多少,取决于仍存在哪些 cache block、shuffle map output、checkpoint 与可靠输入。
stage 是调度边界,不等于可靠检查点
调度器可以把连续窄依赖合并为流水线 stage,宽依赖两侧分成不同 stage。上游 stage 产生 shuffle map output,下游 task 拉取属于自己分区的 bucket。
flowchart LR
R[read partitions] --> S1[Stage 1: map/filter]
S1 --> D[(shuffle files)]
D --> S2[Stage 2: reduce]
S2 --> O[result]
stage 完成表示一组 task 和依赖满足当前调度条件,不表示其输出已经写进独立可靠存储。shuffle 文件通常在 executor 本地磁盘,executor 消失或文件损坏时可能丢失。若文件由外部服务继续提供,下游可直接重取;若文件不可用,调度器重跑对应 map task。
因此需要分开四个词:lineage 是重算配方,stage 是调度切分,shuffle 是数据重分布,checkpoint 才是有意建立的持久恢复基线。把 stage boundary 称为 checkpoint,会把“任务已完成”误写成“状态已可靠物化”。
MapReduce 把轮次边界做成稳定文件
经典 MapReduce 的 map 中间结果先写 worker 本地磁盘,reduce 拉取并排序这些数据;最终 reduce 输出通过临时文件和原子重命名发布到分布式文件系统。worker 失败时,中间 map 输出可能重算,已完成 reduce 输出可从稳定文件读取。
迭代算法通常把第 (i) 轮输出作为第 (i+1) 轮输入。四轮计算因而形成四个作业,每轮结束都把完整 rank 向量物化:
flowchart LR
X0[(rank 0)] --> J1[MR job 1] --> X1[(rank 1)]
X1 --> J2[MR job 2] --> X2[(rank 2)]
X2 --> J3[MR job 3] --> X3[(rank 3)]
X3 --> J4[MR job 4] --> X4[(rank 4)]
好处是恢复边界清楚:第 4 轮失败,通常从稳定的第 3 轮输出继续。成本是每一轮都读写完整状态,并重复读取不变数据。RDD 可以缓存不变的图结构和当前 rank 分区,让相邻迭代不必把完整状态交给分布式文件系统。
这并不推出 RDD 在所有负载上更快。数据大于内存、lineage 很长、shuffle 占主导、执行器频繁回收或可靠存储足够快时,物化成本与重算风险会重新平衡。本篇不引用旧论文的性能倍数,也不拿 Python 计数模型代替集群基准。
cache 与 checkpoint 交换不同成本
persist/cache 让已经计算的分区留在内存或磁盘,供后续 action 重用。Spark 4.2.0 官方指南明确说明:缓存分区丢失时,可用创建它的转换重新计算。带副本的 storage level 可以更快继续,但所有标准级别的逻辑容错仍以重算为后盾。
checkpoint 则把某个 RDD 写到可靠存储,并让后续恢复不再穿过更早 lineage。它适合循环很多、血缘不断增长,或祖先重算代价太高的任务。它增加可靠写入、同步和空间成本,checkpoint 间隔也成为工程参数。
flowchart LR
I0[iter 0] --> I1[iter 1] --> I2[iter 2]
I2 --> CP[(checkpoint iter 2)]
CP --> I3[iter 3] --> I4[iter 4]
L[iter 4 partition lost] -.replay.-> CP
cache 主要优化正常路径,checkpoint 主要限制最坏恢复路径。把 cache 当可靠检查点,会低估 executor 丢失;把每个 RDD 都 checkpoint,又会退回频繁物化。
本地实验:同一结果,三种恢复账本
实验位于 examples/distributed-systems/spark30/check.py,用 PageRank 风格更新运行四个分区、四轮迭代。它不启动 Spark 或 Hadoop,只把计算次数和物化次数打印出来。
1 | |
Python 3.12.3 的正式输出中,三条路径得到相同数值结果。MapReduce 路径四轮写出 16 个持久分区;无 checkpoint 的 RDD 路径没有中间持久写入。模型随后假设最终 rank 分区丢失,且中间 rank/cache/shuffle 文件均不可用:没有 checkpoint 时需从第 0 轮重放 16 个 task,第二轮 checkpoint 后只重放第 3、4 轮的 8 个 task。
这个数字只属于固定模型。真实 Spark 若保留了部分 shuffle output 或 cache block,重算会更小;若 executor、fetch 与祖先数据同时丢失,重算可能更大。实验验证的是“恢复范围由尚存物化边界决定”,不是 Spark 的通用 task 数公式。
完整输出见结构化观察,来源与运行证据见实验证据,未验证边界见验证说明。
重算会重复 task,也会重复副作用
相同的 task attempt 可以因为 fetch failure、executor failure、stage 重跑或推测执行而再次运行。只要转换是确定且无副作用的,多次得到同一分区不会破坏 RDD 结果。外部 HTTP、邮件、数据库增量更新和不受约束的 accumulator 不是 RDD 分区。
sequenceDiagram
participant D as Driver
participant T1 as Task attempt 1
participant X as External API
participant T2 as Retry
D->>T1: compute partition 2
T1->>X: effect succeeds
Note over T1: result lost before driver accepts it
D->>T2: recompute partition 2
T2->>X: effect succeeds again
安全做法取决于 sink:用稳定业务键实现幂等写;把结果先写成可重复覆盖的分区文件,再由独立提交协议发布;或让支持事务的 sink 参与可恢复提交。仅仅把 cache() 加在 RDD 上,不会为外部调用增加一次性语义。
安全性、活性与工程边界
RDD 重算的安全直觉可以写成归纳:可靠输入分区确定;若每个 transformation 对确定父分区产生确定子分区,则沿 lineage 重算得到与首次执行相同的值。宽依赖只改变需要读取的父分区集合,不改变这一归纳结构。
活性要求 driver 仍可调度、可靠输入和 checkpoint 可读、集群有足够资源,失败没有持续发生。某个 executor 反复丢失、shuffle fetch 一直失败或外部存储不可达时,lineage 只能说明“怎样重算”,不能保证最终算完。
| 问题 | 首先检查 | 典型选择 |
|---|---|---|
| 同一 RDD 被多次 action 使用 | 分区是否值得复用 | persist,并选择 storage level |
| 血缘很长、祖先重算昂贵 | 最坏重放长度 | 定期可靠 checkpoint |
| shuffle 后 executor 会缩容 | shuffle 文件是否随 executor 消失 | shuffle service/tracking/decommission/可靠存储 |
| task 内有外部写入 | 重试是否会重复 | 幂等键、事务 sink、分区级提交 |
| 只比较“内存 vs 磁盘” | 恢复证据与物化边界 | 先画 DAG 和故障后的可用产物 |
两个推演练习
丢失的是 cache 还是 checkpoint。 第 20 轮 RDD 已 cache,第 10 轮做过可靠 checkpoint;第 20 轮的一个分区和它所需的 shuffle 文件都随 executor 丢失。恢复起点至少可以落在哪里?
cache 丢失不能作为可靠起点。若第 11–20 轮的其他中间产物也不可用,最迟从第 10 轮 checkpoint 沿 lineage 重算;实际只需多少 task,还要看仍在的分区和 shuffle map output。
副作用在 task 返回前成功。 task 写外部计数器后,executor 在结果被 driver 接受前崩溃。lineage 重算能否把外部计数回滚?
不能。lineage 只重建 RDD 分区。重试会再次执行 task 代码,计数器可能增加两次。需要幂等键、覆盖式结果或外部事务协议。
工程结论
MapReduce 用逐轮稳定物化换取短恢复路径;RDD 用分区 lineage 和选择性缓存减少正常路径 I/O,再用 checkpoint 给长血缘设置上限。shuffle 是宽依赖的数据交换与调度边界,不天然等于可靠物化边界。
下一篇进入 Ray。RDD 的节点主要是数据分区及转换,Ray 把远程函数返回值直接暴露成 distributed future;失败时同样会重算 DAG,也同样不能自动撤销任务里的外部副作用。
参考资料
- Zaharia et al., 2012, Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing:§1–§3、§5。
- Dean and Ghemawat, 2004, MapReduce: Simplified Data Processing on Large Clusters:§3.1、§3.3、§4.2。
- Apache Spark 4.2.0 RDD Programming Guide:Overview、Shuffle Operations、RDD Persistence、Accumulators,访问 2026-09-26。
- Apache Spark 4.2.0 Job Scheduling:Stages、Dynamic Resource Allocation、Graceful Decommission,访问 2026-09-26。
