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

很多团队写批处理任务的第一反应是 MapReduce,但真正跑起来才发现:一个简单的 WordCount 要拆成 Map 和 Reduce 两个阶段,中间结果全部落盘 HDFS,稍复杂的迭代算法要串起十几个 MR 作业,每一个都在读写磁盘。Spark 用 DAG 执行引擎 + 基于内存的计算 把这个问题彻底重做了一遍。本文从原理到实战,把 Spark 最容易被问、最容易踩坑的部分一次讲透。

一、为什么需要 Spark:MapReduce 的三大痛点

MapReduce 是第一代批处理引擎,但它天生”腿短”:

痛点 MapReduce 的表现 Spark 的解法
中间结果落盘 每个 MR 作业的输出写 HDFS,下游再读一遍 中间结果保留在内存,必要时才溢写磁盘
表达能力弱 只有 Map / Reduce 两个算子,Join、排序要手写 80+ 算子(flatMapreduceByKeyjoin…),代码量降一个量级
延迟高 JVM 冷启动 + 磁盘 IO,分钟级延迟 DAG 调度 + 线程级任务模型,秒级~亚秒级

一句话总结:MapReduce 把计算当”流水线上的搬运工”,Spark 把计算当”可以反复使用的工作台”。数据一旦加载进集群,后续十几步转换都在内存工作台上进行,只在必要时才落盘。

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
2
3
4
5
6
val lines = sc.textFile("hdfs://data/logs")            // Transformation:不执行
val errors = lines.filter(_.contains("ERROR")) // Transformation:不执行
val wc = errors.flatMap(_.split(" "))
.map((_, 1))
.reduceByKey(_ + _) // Transformation:仍不执行
wc.collect() // Action:触发真正计算
类别 触发时机 代表算子 返回值
Transformation 惰性,只记录血缘 map filter flatMap reduceByKey join distinct 新 RDD
Action 立即触发 Job collect count take first saveAsTextFile foreach 非 RDD(值/落盘)

每个 Action 生成一个 Job。代码里写了三个 collect,就是三个独立的 Job,前面的血缘会被重复计算——这就是为什么 “一个应用一个 Action” 是常见优化原则,或者用 persist() 把公共中间结果缓存起来。

三、宽窄依赖与 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
2
3
4
5
6
7
sc.textFile("data.txt")          // Stage 0 开始
.flatMap(_.split(" ")) // 窄
.map((_, 1)) // 窄
.reduceByKey(_ + _) // 宽依赖 → 在此切分 Stage
.map(kv => (kv._2, kv._1)) // Stage 1
.sortByKey(ascending = false) // 宽依赖 → 切分
.collect() // Stage 2(Action 触发 Job)

整条链被切成 3 个 Stage:前两个窄操作被”打包”进同一个 Stage 流水线执行,reduceByKeysortByKey 各自切开一个新 Stage。

四、任务调度:从 Job 到 Task 的完整链路

一个 Action 触发后,调度链路如下:

1
2
3
4
5
Application
└── Job(每个 Action 一个)
└── Stage(DAGScheduler 按宽依赖切分)
└── TaskSet(Stage 内所有分区的任务集合)
└── Task(TaskScheduler 分发到 Executor 执行)

关键角色:

  • 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
2
小文件数 = Map Task 数 × Reduce Task 数
1000 map × 1000 reduce = 100 万个小文件

文件句柄耗尽、磁盘随机 IO 打爆、GC 频繁。优化后的 Consolidate 机制把同一 Executor 上 Core 的输出合并为 CPU 核数 × Reduce Task 数,但依然治标不治本。

5.2 SortShuffle:现在的默认机制(1.2+)

SortShuffle 的核心思路:每个 Map Task 只输出一个按分区有序的数据文件 + 一个索引文件

  1. 数据写入内存排序结构(PartitionedAppendOnlyMap),超过阈值后溢写成临时文件;
  2. 最终把所有溢写文件归并成一个数据文件,按 Reduce 分区有序排列;
  3. 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
2
3
4
5
6
Executor 堆内内存
├── Reserved Memory:固定 300MB
├── User Memory:(1 - spark.memory.fraction) × 堆,存用户数据结构、RDD 依赖等
└── Unified Memory(spark.memory.fraction = 0.6)
├── Execution Memory:Shuffle / Join / Sort / 聚合 —— 不被抢占
└── Storage Memory(storageFraction = 0.5):缓存 RDD / 广播变量 —— 可被驱逐落盘

动态借用的不对称规则是面试重点:两块内存空闲时可以互相借用;但 Execution 占用的部分绝不会被归还(中途释放会导致数据结构失效),Storage 借用的部分在 Execution 需要时会被驱逐(缓存块溢写磁盘或丢弃)。所以 Executor OOM 大多发生在 Execution 侧——Shuffle 聚合是内存大户。

堆外内存通过 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
2
3
4
sc.setCheckpointDir("hdfs://checkpoints")
val result = longLineageRdd.persist()
result.checkpoint() // 必须在 Action 之前设置
result.count() // 触发时顺带写入 checkpoint

实践口诀:persistcheckpoint——否则 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
2
import org.apache.spark.sql.functions.broadcast
val joined = bigDF.join(broadcast(smallDF), "user_id") // 小表全量广播到各 Executor

阈值由 spark.sql.autoBroadcastJoinThreshold(默认 10MB)控制,超过时需手动 broadcast() 或调大阈值(注意广播本身占内存)。

方案三:加盐两阶段聚合。reduceByKey 类聚合的倾斜 Key 加随机前缀,先局部聚合打散热点,再去前缀做全局聚合:

1
2
3
4
5
6
// 第一阶段:key 加 0~N 随机前缀,局部聚合
val partial = rdd.map { case (k, v) => ((Random.nextInt(10) + "_" + k), v) }
.reduceByKey(_ + _)
// 第二阶段:去掉前缀,全局聚合
val global = partial.map { case (pk, v) => (pk.split("_")(1), v) }
.reduceByKey(_ + _)

方案四:提高 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 负责缓冲与解耦,计算引擎负责处理)。

十、生产避坑清单

  1. 并行度别拍脑袋:理想状态是 Task 数 = 总 Core 数的 2~3 倍;spark.sql.shuffle.partitions=200 对小数据是浪费、对大数据是瓶颈,按单分区 128MB~256MB 反推。
  2. 别用 collect() 拉全量数据:Driver OOM 的头号来源,调试用 take(n) / limit(n),落盘用 write
  3. 小文件问题:频繁写 Hive/对象存储会产生海量小文件,写入前 coalesce(n) 合并分区(保序减少数据移动用 coalesce,允许全量重分布用 repartition)。
  4. Kryo 序列化spark.serializer=org.apache.spark.serializer.KryoSerializer,自定义类注册 registerKryoClasses,比 Java 序列化快 10 倍、体积小一半。
  5. 缓存不是越多越好MEMORY_ONLY 缓存被频繁淘汰重算时反而更慢;缓存前先评估重用次数与分区大小,用 unpersist() 及时释放。
  6. GC 调优顺序:先调内存分配(fraction / overhead),再看 GC 日志选收集器;spark.executor.cores=4~5 能减少 HDFS 客户端并发争抢。
  7. 广播变量只读:Executor 端修改广播变量不会同步回 Driver,也看不到——它本质是只读快照。
  8. Spill 不等于故障:Shuffle 溢写是正常兜底机制;持续大量 Spill 才说明内存给少了或倾斜了,结合 Web UI 的 Spill (memory/disk) 指标判断。