0%

Spark 核心执行模型

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 才回溯血缘生成执行计划并提交。

这带来两个关键工程后果:

  1. 可优化:Driver 拿到的是完整算子链,可以做流水线融合(多个窄依赖算子合并进一个 Task,中间结果不落地,见 Spark Codegen)。
  2. 可容错:不需要 Checkpoint,丢失的 partition 可以按血缘重算。代价是血缘过长时重算成本高——这是 cache/checkpoint 的存在理由。

调试血缘:rdd.toDebugString 可打印依赖链与 Stage 边界。

2.2 三层调度对象

1
2
3
4
Driver
├── DAGScheduler —— 面向 Stage:切分 DAG、提交 TaskSet、处理 Stage 失败重试
├── TaskScheduler —— 面向 Task:把 Task 分发到 Executor、处理推测执行/黑名单
└── SchedulerBackend —— 与集群管理器(YARN/K8s)交互申请资源

面试常考的分工边界: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 出发逆向遍历血缘:

  1. 遇窄依赖 → 把该 RDD 并入当前 Stage,继续向父回溯;
  2. 遇宽依赖 → 在此断开,当前 Stage 结束(成为 ResultStage 或 ShuffleMapStage),为父 RDD 新建一个 ShuffleMapStage,递归处理;
  3. 结果是一个 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 )

  1. 向 Driver 的 MapOutputTracker 查询自己那个分区的数据都在哪些 Executor 上;
  2. 通过 Netty 并发拉取(spark.reducer.maxSizeInFlight 限制在途数据量);
  3. 需聚合则边拉边聚合,需排序则用 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 首次掌握各分区的真实字节数,据此重写尚未执行的下游计划:

  1. 动态合并分区(Coalesce Partitions):把 200 个小分区合并成少数几个,消除空 Task;
  2. 动态切换 Join 策略:发现一侧真实体积远小于阈值,把 Sort-Merge Join 改为 Broadcast Join,直接省掉一次 Shuffle;
  3. 动态处理倾斜 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 参考资料

  1. Apache Spark 官方文档 —— RDD Programming Guide / Tuning Guide
  2. 《Spark 性能优化指南——高级篇》(美团技术团队)
  3. Matei Zaharia et al., Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing, NSDI 2012