Spark 核心原理与实战:RDD、DAG 调度、Shuffle 与内存管理

Spark 核心原理与实战:RDD、DAG 调度、Shuffle 与内存管理
神经蛙很多团队写批处理任务的第一反应是 MapReduce,但真正跑起来才发现:一个简单的 WordCount 要拆成 Map 和 Reduce 两个阶段,中间结果全部落盘 HDFS,稍复杂的迭代算法要串起十几个 MR 作业,每一个都在读写磁盘。Spark 用 DAG 执行引擎 + 基于内存的计算 把这个问题彻底重做了一遍。本文从原理到实战,把 Spark 最容易被问、最容易踩坑的部分一次讲透。
一、为什么需要 Spark:MapReduce 的三大痛点
MapReduce 是第一代批处理引擎,但它天生”腿短”:
| 痛点 | MapReduce 的表现 | Spark 的解法 |
|---|---|---|
| 中间结果落盘 | 每个 MR 作业的输出写 HDFS,下游再读一遍 | 中间结果保留在内存,必要时才溢写磁盘 |
| 表达能力弱 | 只有 Map / Reduce 两个算子,Join、排序要手写 | 80+ 算子(flatMap、reduceByKey、join…),代码量降一个量级 |
| 延迟高 | JVM 冷启动 + 磁盘 IO,分钟级延迟 | DAG 调度 + 线程级任务模型,秒级~亚秒级 |
Spark 生态覆盖了统一的计算场景:Spark Core(RDD 内核)、Spark SQL(结构化数据 + Catalyst 优化器)、Spark Streaming / Structured Streaming(流计算)、MLlib(机器学习)、GraphX(图计算)。部署上支持 Local、Standalone、YARN、Kubernetes 四种模式,生产以 YARN / K8s 为主。
二、RDD:Spark 的基石
2.1 RDD 是什么
RDD(Resilient Distributed Dataset,弹性分布式数据集)是一个不可变、可分区、可并行计算的元素集合。”弹性”体现在三层:
- 存储弹性:内存不够自动落盘,磁盘不够自动重算;
- 计算弹性:失败自动通过血缘(Lineage)重算恢复;
- 分区弹性:分区数随时可以重新划分(
repartition/coalesce)。
源码里 RDD 的注释写明了五大特性,面试高频,务必能背出主干:
| 特性 | 含义 | 一句话理解 |
|---|---|---|
| A list of partitions | 一组分区列表 | 并行度的基本单位 |
| A function for computing each split | 每个分区一个计算函数 | 作用在分区而非全量数据上 |
| A list of dependencies on other RDDs | 依赖(血缘)列表 | 容错重算的依据 |
| Optionally, a Partitioner | 可选分区器(K-V RDD) | Hash / Range,决定数据分布 |
| Optionally, preferred locations | 可选优先位置 | 移动计算不移动数据 |
2.2 算子体系:Transformation 与 Action
RDD 算子分两类,这个区分决定了作业什么时候真正执行:
1 | val lines = sc.textFile("hdfs://data/logs") // Transformation:不执行 |
| 类别 | 触发时机 | 代表算子 | 返回值 |
|---|---|---|---|
| Transformation | 惰性,只记录血缘 | map filter flatMap reduceByKey join distinct |
新 RDD |
| Action | 立即触发 Job | collect count take first saveAsTextFile foreach |
非 RDD(值/落盘) |
三、宽窄依赖与 DAG:Stage 划分的分界线
3.1 宽依赖 vs 窄依赖
| 依赖类型 | 定义 | 典型算子 | 特点 |
|---|---|---|---|
| 窄依赖 | 父 RDD 的每个分区最多被子 RDD 的一个分区使用 | map filter union mapPartitions |
分区可流水线(pipeline)执行,一个分区的失败只需重算父分区 |
| 宽依赖 | 子 RDD 的分区依赖父 RDD 的多个分区 | groupByKey reduceByKey join(非同分区) |
必须等父 Stage 全部就绪,伴随 Shuffle |
DAG 划分 Stage 的规则:从后往前,遇到宽依赖就切一刀。 每个切出来的阶段是一个 Stage,Stage 内部尽量流水线执行,Shuffle 是两个 Stage 之间的数据交换边界。
3.2 一个例子看懂 Stage 划分
1 | sc.textFile("data.txt") // Stage 0 开始 |
整条链被切成 3 个 Stage:前两个窄操作被”打包”进同一个 Stage 流水线执行,reduceByKey 和 sortByKey 各自切开一个新 Stage。
四、任务调度:从 Job 到 Task 的完整链路
一个 Action 触发后,调度链路如下:
1 | Application |
关键角色:
- DAGScheduler:负责把血缘 DAG 按 Stage 划分、维护 Stage 依赖、提交”没有父依赖”的 Stage;同时处理 ShuffleMapStage 与 ResultStage 的区分。
- TaskScheduler:把 TaskSet 里的 Task 按本地性级别(
PROCESS_LOCAL → NODE_LOCAL → ANY)分发给 Executor,支持推测执行(Speculative Execution):某个 Task 明显慢于同 Stage 均值时,另起一个备份 Task,谁先完成用谁的结果。 - Executor:以线程池方式执行 Task(MapReduce 是进程级,启动开销大),这是 Spark 秒级延迟的另一个原因。
五、Shuffle:性能问题的 80% 都在这
Shuffle 是宽依赖的数据交换过程——按 Key 把数据重新分发到不同的下游分区。它是 Spark 里最昂贵、也最值得调优的环节。
5.1 HashShuffle:为什么被淘汰
未优化的 HashShuffle 中,每个 Map Task 为每个 Reduce Task 生成一个文件:
1 | 小文件数 = Map Task 数 × Reduce Task 数 |
文件句柄耗尽、磁盘随机 IO 打爆、GC 频繁。优化后的 Consolidate 机制把同一 Executor 上 Core 的输出合并为 CPU 核数 × Reduce Task 数,但依然治标不治本。
5.2 SortShuffle:现在的默认机制(1.2+)
SortShuffle 的核心思路:每个 Map Task 只输出一个按分区有序的数据文件 + 一个索引文件。
- 数据写入内存排序结构(
PartitionedAppendOnlyMap),超过阈值后溢写成临时文件; - 最终把所有溢写文件归并成一个数据文件,按 Reduce 分区有序排列;
- Reduce Task 根据索引文件偏移量拉取属于自己的数据段。
Bypass 机制是 SortShuffle 的”轻量通道”:当 Reduce Task 数 ≤ spark.shuffle.sort.bypassMergeThreshold(默认 200)且算子不需要 map 端聚合(如 groupBy,而非 reduceByKey)时,跳过排序,直接按分区写文件再归并——用 HashShuffle 的写法拿到单文件的好处。
5.3 Reduce 端拉取与关键参数
Reduce 端通过 HTTP 拉取数据,相关参数直接影响稳定性:
| 参数 | 默认值 | 说明 |
|---|---|---|
spark.reducer.maxSizeInFlight |
48m | 单次拉取缓冲区,太小网络往返多 |
spark.shuffle.io.maxRetries |
3 | 拉取失败重试次数,生产建议 5+ |
spark.shuffle.io.retryWait |
5s | 重试间隔,生产建议 15s+(maxRetries × retryWait 要能扛住 FullGC) |
spark.sql.shuffle.partitions |
200 | SQL 场景 Reduce 并行度,按数据量调整 |
六、内存管理:统一内存管理模型
Spark 1.6+ 采用 UnifiedMemoryManager(统一内存管理),Executor 堆内空间划分如下:
1 | Executor 堆内内存 |
堆外内存通过 spark.memory.offHeap.enabled 开启,用 spark.executor.memoryOverhead 预留堆外空间(YARN/K8s 下默认为堆的 10%),容器总内存 = 堆 + 堆外 + overhead,配小了会被 YARN 直接 kill。
七、持久化与容错:Lineage 与 Checkpoint
7.1 缓存:persist / cache
| 级别 | 空间 | 说明 |
|---|---|---|
MEMORY_ONLY(cache 默认) |
内存 | 内存不够则部分分区不缓存,用时重算 |
MEMORY_AND_DISK |
内存+磁盘 | 内存放不下的溢写磁盘,最常用 |
MEMORY_ONLY_SER |
内存 | 序列化存储,省空间但耗 CPU |
DISK_ONLY |
磁盘 | 完全落盘 |
缓存本身不切断血缘——缓存分区丢失时,仍然沿血缘重算。
7.2 Checkpoint:斩断血缘
血缘太长时(迭代算法、长链路转换),单分区失败的重算代价会指数级膨胀。checkpoint 把 RDD 完整写入高可靠存储(HDFS)并清空血缘,失败后直接从 checkpoint 恢复:
1 | sc.setCheckpointDir("hdfs://checkpoints") |
实践口诀:先 persist 再 checkpoint——否则 checkpoint 会把血缘完整重算一遍,等于白算两次。
八、数据倾斜:定位与治理套路
8.1 现象与定位
- 现象:绝大多数 Task 几秒跑完,个别 Task 卡几十分钟甚至 OOM;
- 定位:Web UI 的 Stage 页看 Task 的 Shuffle Read Size 分布是否长尾;
df.groupBy(key).count().orderBy("count", ascending=False)检查热点 Key。
8.2 治理方案(按优先级)
方案一:过滤或单独处理异常 Key。 如果倾斜来自 null、"" 等脏数据,先过滤;如果是个别真实大 Key,把它们捞出来单独跑,再和主流程结果合并。
方案二:广播 Join(小表广播)。 大小表 Join 时避免 Shuffle:
1 | import org.apache.spark.sql.functions.broadcast |
阈值由 spark.sql.autoBroadcastJoinThreshold(默认 10MB)控制,超过时需手动 broadcast() 或调大阈值(注意广播本身占内存)。
方案三:加盐两阶段聚合。 对 reduceByKey 类聚合的倾斜 Key 加随机前缀,先局部聚合打散热点,再去前缀做全局聚合:
1 | // 第一阶段:key 加 0~N 随机前缀,局部聚合 |
方案四:提高 Shuffle 并行度。 调大 spark.sql.shuffle.partitions,让热点 Key 的数据被切得更碎,缓解(而非根治)倾斜。
九、横向对比:Spark vs MapReduce vs Flink
| 维度 | MapReduce | Spark | Flink |
|---|---|---|---|
| 计算模型 | 批处理 | 批处理为主(微批流) | 流处理为主(流批一体) |
| 中间结果 | 落盘 HDFS | 内存优先 | 基于流,数据持续流动 |
| 延迟 | 分钟级 | 秒级~亚秒级 | 毫秒级 |
| 容错 | 任务级重跑 | Lineage 重算 + Checkpoint | 分布式快照(Chandy-Lamport) |
| 适用场景 | 超大规模离线归档 | 离线数仓、ETL、机器学习 | 实时数仓、CEP、风控 |
选型经验:离线批处理选 Spark,实时流处理选 Flink,两者不是替代关系。Kafka + Spark/Flink 是流式架构最常见的组合(Kafka 负责缓冲与解耦,计算引擎负责处理)。
十、生产避坑清单
- 并行度别拍脑袋:理想状态是 Task 数 = 总 Core 数的 2~3 倍;
spark.sql.shuffle.partitions=200对小数据是浪费、对大数据是瓶颈,按单分区 128MB~256MB 反推。 - 别用
collect()拉全量数据:Driver OOM 的头号来源,调试用take(n)/limit(n),落盘用write。 - 小文件问题:频繁写 Hive/对象存储会产生海量小文件,写入前
coalesce(n)合并分区(保序减少数据移动用coalesce,允许全量重分布用repartition)。 - Kryo 序列化:
spark.serializer=org.apache.spark.serializer.KryoSerializer,自定义类注册registerKryoClasses,比 Java 序列化快 10 倍、体积小一半。 - 缓存不是越多越好:
MEMORY_ONLY缓存被频繁淘汰重算时反而更慢;缓存前先评估重用次数与分区大小,用unpersist()及时释放。 - GC 调优顺序:先调内存分配(fraction / overhead),再看 GC 日志选收集器;
spark.executor.cores=4~5能减少 HDFS 客户端并发争抢。 - 广播变量只读:Executor 端修改广播变量不会同步回 Driver,也看不到——它本质是只读快照。
- Spill 不等于故障:Shuffle 溢写是正常兜底机制;持续大量 Spill 才说明内存给少了或倾斜了,结合 Web UI 的 Spill (memory/disk) 指标判断。











