1 一句话主线
用户代码 → DAG → 按宽依赖切 Stage → Stage 内并行跑 Task → Shuffle 落盘作为 Stage 边界 → 下游 Stage 拉取。
理解 Spark 性能问题,本质就是理解这条链路上”哪一环把并行变成了串行、把内存变成了磁盘、把本地变成了网络”。
| 层级 | 触发者 | 数量决定因素 |
|---|---|---|
| Application | spark-submit 一次提交 |
1 |
| Job | 每个 Action 算子(count/collect/save) |
Action 个数 |
| Stage | 每个 宽依赖(Shuffle) | Shuffle 次数 + 1 |
| Task | Stage 内每个 partition | 上游分区数 / spark.sql.shuffle.partitions |
概念对齐与并发度换算(Executor / core / partition 的关系)见 Spark task split block。
2 DAG 的构建:惰性求值
2.1 Transformation 只记账,Action 才算账
RDD 上的 map/filter/join 等 Transformation 不触发计算,只在 Driver 端记录”血缘”(Lineage)——即当前 RDD 由哪个父 RDD、经何种变换得来。只有遇到 Action 时,DAGScheduler 才回溯血缘生成执行计划并提交。
这带来两个关键工程后果:
- 可优化:Driver 拿到的是完整算子链,可以做流水线融合(多个窄依赖算子合并进一个 Task,中间结果不落地,见 Spark Codegen)。
- 可容错:不需要 Checkpoint,丢失的 partition 可以按血缘重算。代价是血缘过长时重算成本高——这是
cache/checkpoint的存在理由。
调试血缘:rdd.toDebugString 可打印依赖链与 Stage 边界。
2.2 三层调度对象
1 | Driver |
面试常考的分工边界:Stage 级失败重试属于 DAGScheduler,Task 级重试与推测执行属于 TaskScheduler。
3 Stage 划分:宽窄依赖是唯一标准
3.1 判定规则
看父 RDD 的一个 partition 是否被多个子 partition 使用:
| 依赖类型 | 特征 | 典型算子 | 是否 Stage 边界 |
|---|---|---|---|
| 窄依赖 Narrow | 父 partition → 唯一子 partition | map/filter/mapPartitions/union、co-partitioned 的 join |
否,可流水线执行 |
| 宽依赖 Wide / Shuffle | 父 partition → 多个子 partition,需按 key 重分布 | groupByKey/reduceByKey/join/distinct/repartition |
是 |
3.2 划分算法(反向回溯)
DAGScheduler 从 Action 对应的 final RDD 出发逆向遍历血缘:
- 遇窄依赖 → 把该 RDD 并入当前 Stage,继续向父回溯;
- 遇宽依赖 → 在此断开,当前 Stage 结束(成为 ResultStage 或 ShuffleMapStage),为父 RDD 新建一个 ShuffleMapStage,递归处理;
- 结果是一个 Stage 级 DAG,无父 Stage 的先执行,有依赖的按拓扑序等待。
因此有 Stage 数 = Shuffle 次数 + 1(单 Job 内、无分支时)。
3.3 两类 Stage
- ShuffleMapStage:输出写入 Shuffle 文件,供下游拉取;输出位置注册到
MapOutputTracker。 - ResultStage:执行 Action 本身,结果返回 Driver 或写外部存储。
一个附带收益:ShuffleMapStage 的输出会被复用。同一 RDD 被多个 Job 使用时,已完成的 ShuffleMapStage 可被跳过(Spark UI 中显示 Skipped Stages),这是”为什么第二次 Action 变快了”的答案。
4 Shuffle:Stage 边界上发生了什么
Shuffle 是跨 Stage 的数据重分布,也是性能问题的主要发源地
——它同时引入了磁盘 I/O、网络传输、序列化三重开销,并且是必须全部完成才能进入下游的同步屏障。
4.1 Write 侧(上游 ShuffleMapTask)
每个 map task 按分区器(默认 HashPartitioner)计算目标 reduce 分区,写出一个数据文件 + 一个索引文件(Sort Shuffle):
| 写入器 | 触发条件 | 行为 |
|---|---|---|
BypassMergeSortShuffleWriter |
分区数少(默认 ≤ 200)且无 map 端聚合 | 每分区一临时文件,最后合并,跳过排序 |
UnsafeShuffleWriter |
无聚合、无排序、支持序列化重定位 | 在序列化后的二进制上排序,Tungsten 路径 |
SortShuffleWriter |
兜底(需 map 端聚合/排序时) | 内存缓冲 + 溢写(spill)+ 归并 |
关键点:内存不足时会 spill 到磁盘,spill 次数是判断 Shuffle 是否健康的重要指标(Spark UI 的 Spill (Memory)/Spill (Disk) 两列)。
4.2 Read 侧( 下游 ShuffleReduceTask )
- 向 Driver 的
MapOutputTracker查询自己那个分区的数据都在哪些 Executor 上; - 通过 Netty 并发拉取(
spark.reducer.maxSizeInFlight限制在途数据量); - 需聚合则边拉边聚合,需排序则用
ExternalSorter(同样可能 spill)。
4.3 从执行模型看 Shuffle 的三类病症
| 症状 | 执行模型层面的解释 | 对策方向 |
|---|---|---|
| 少数 Task 拖尾(99/100 completed 卡住) | 分区器把大量同 key 数据打到同一 reduce 分区 → 单 Task 数据量畸大 | 见 Spark 数据倾斜:加盐打散、两阶段聚合、Broadcast Join 消除 Shuffle |
| Shuffle 数据量巨大 / 磁盘打满 | 宽依赖过多,且未做 map 端预聚合 | reduceByKey 替代 groupByKey(map 端 combine);提前列裁剪与过滤下推 |
| 分区数不合理 | spark.sql.shuffle.partitions 默认 200,与实际数据量脱节:过大产生大量空 Task 调度开销,过小则单 Task OOM |
开启 AQE 让运行时动态合并/拆分分区 |
Shuffle 全过程图见 Spark Shuffle。
5 AQE:让静态计划变成运行时自适应
静态计划在编译期一次定稿,只能依赖表的统计信息做代价估算,一旦统计信息缺失或过期,就会选错 Shuffle、Join 策略或分区数。
AQE(Adaptive Query Execution,Spark 3.0 引入)把优化时机推迟到运行时,由 spark.sql.adaptive.enabled 控制(3.0 默认关闭,3.2 起默认开启)。
它的决策点正是 Stage 边界——Shuffle 写完后,Driver 首次掌握各分区的真实字节数,据此重写尚未执行的下游计划:
- 动态合并分区(Coalesce Partitions):把 200 个小分区合并成少数几个,消除空 Task;
- 动态切换 Join 策略:发现一侧真实体积远小于阈值,把 Sort-Merge Join 改为 Broadcast Join,直接省掉一次 Shuffle;
- 动态处理倾斜 Join(Skew Join):把超大分区自动拆成多个子分区并复制另一侧数据。
其中”动态切换 Join 策略”是运行时在三种物理 Join 之间改写,方向与触发条件:
| 切换 | 方向 | 触发条件 |
|---|---|---|
| SortMergeJoin → BroadcastHashJoin | 升级(省一次 Shuffle) | 一侧 Shuffle 输出 < spark.sql.adaptive.autoBroadcastJoinThreshold |
| SortMergeJoin → ShuffledHashJoin | 切换(省排序) | 各分区可装入 spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold |
| BroadcastHashJoin → SortMergeJoin | 降级(纠错) | 广播侧实际超过阈值,由 DemoteBroadcastHashJoin 退回 |
注意:Skew Join 是分区级倾斜治理,不属于 Join 策略切换。
这也解释了 AQE 的能力边界:它只能在 Stage 边界生效,无法优化 Stage 内部。
执行计划层的规则优化(列裁剪、谓词下推、Join Reorder)见 Spark SQL。
6 面试速答模板
问:讲讲 Spark 的执行模型。
从 Action 触发说起:Transformation 只构建血缘 DAG,Action 触发 DAGScheduler 逆向回溯,以宽依赖为界切出 Stage——窄依赖能流水线融合所以合并进同一 Stage,宽依赖需要跨节点按 key 重分布所以必须断开。每个 Stage 按 partition 展开成 TaskSet 交给 TaskScheduler 分发,Stage 之间通过 Shuffle 文件衔接,是个同步屏障。所以性能优化的落点很明确:减少 Shuffle 次数、减少 Shuffle 数据量、让 Shuffle 后的分区分布均匀,这三条分别对应 Broadcast Join、map 端预聚合与列裁剪、以及倾斜治理和 AQE。
追问:为什么宽依赖一定要切 Stage?
因为子 partition 的输入依赖多个父 partition,必须等父 Stage 全部 task 完成才能确定数据完整——无法像窄依赖那样一条记录来了就往下走。这个”全部完成”就是同步屏障,也是长尾 Task 会阻塞整个作业的根因。
7 可迁移的方法论
这套模型不止用于 Spark,它是所有分布式执行引擎的通用骨架,横向迁移时只需替换名词:
| Spark | Flink | GPU 推理 |
|---|---|---|
| Stage 边界(同步屏障) | 算子链断点 / Checkpoint Barrier | Kernel 间的同步点 |
| Task 长尾 | 反压(Backpressure)与热点 subtask | 静态 Batching 的木桶效应 |
| Shuffle 磁盘 + 网络开销 | 网络 Shuffle / 状态访问 | HBM 读写(FlashAttention 优化的对象) |
| Executor 内存 OOM | TaskManager 堆外内存超限 | 显存 OOM |
| Spark UI 定位 Top 耗时 Stage | Flink Web UI 反压面板 | Nsight Systems Kernel 时间线 |
共同的排查动作是同构的:采集基线 → 找到最慢的那个执行单元 → 判断它慢在计算/内存/I-O/通信哪一层 → 单变量验证优化收益。
与 Flink 的抽象差异对比见 Flink VS Spark;面试中的应用场景见 AI 开发技术图谱 1.1.1。
8 参考资料
- Apache Spark 官方文档 —— RDD Programming Guide / Tuning Guide
- 《Spark 性能优化指南——高级篇》(美团技术团队)
- Matei Zaharia et al., Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing, NSDI 2012