【Spark Core 笔记】RDD 原理、调度流程和线上踩坑
这篇笔记整理了用 Spark 的一些心得,覆盖了面试常考的核心知识点,也包括了不少线上真实踩过的坑。
一、Spark Core 到底是什么
先给一个我自己的理解:Spark Core 是整个 Spark 生态的"发动机"。不管是后面你用的 Spark SQL、Spark Streaming 还是 MLlib,底层的执行引擎都是 Core 这一套东西。它解决的核心问题就两个:
- 如何把一个大计算任务拆成小块,分布到多台机器上并行跑
- 跑的过程中某个节点挂了,怎么自动恢复而不是从头再来
这两个问题分别对应 分布式计算 和 容错机制,也是 Spark 区别于 MapReduce 最核心的竞争力。
二、RDD —— 一切计算的起点
2.1 RDD 的本质
RDD(Resilient Distributed Dataset,弹性分布式数据集),名字很长,但其实理解起来不复杂。你可以把它想象成一个被切成了很多份的不可变数据集合,每一份存到不同的节点上,你对它的每一次操作都会生成一个新的 RDD。
这里有几个关键词需要注意:
| 特性 | 说明 |
|---|---|
| 不可变(Immutable) | RDD 一旦被创建就不能修改,transform 操作会产生新的 RDD。这和 Java 的 String 类似。 |
| 分区(Partitioned) | 数据被切分成多个 partition,每个 partition 作为一个 Task 在不同节点上执行。 |
| 依赖关系(Lineage) | 每个 RDD 都记录了它是怎么从父 RDD 算出来的,这就是血统/依赖链。 |
| 算子(Lazy Evaluation) | Transformation 不会立即执行,只有遇到 Action 才会真正触发计算。 |
| 持久化(Persist/Cache) | 可以把中间结果缓存到内存或磁盘,避免重复计算。 |
2.2 RDD vs DataFrame:不是谁更好,是不同阶段的选择
很多新手会纠结:既然 DataFrame 有 Catalyst 优化器,那还用 RDD 干嘛?
| 维度 | RDD | DataFrame |
|---|---|---|
| 优化层面 | 无自动优化,手写逻辑是啥就执行啥 | Catalyst 自动优化执行计划 |
| 类型安全 | 编译期类型检查,操作的是 JVM 对象 | 运行时才能发现类型错误 |
| 灵活性 | 高,可以精确控制每个 partition 的处理 | 低,SQL 表达式覆盖不到的只能 UDF |
| 性能 | 序列化开销大(JVM 对象) | Tungsten 二进制格式,内存效率高 |
| 适用场景 | 细粒度控制、非结构化数据、自定义优化 | 结构化数据分析、ETL、数仓场景 |
一个实用的判断标准:如果你的数据已经是结构化的(JSON、Parquet、Hive 表),直接上 DataFrame/Dataset;如果你需要精确控制数据的分布、分区策略,或者做一些非标准的数据处理(比如图计算),RDD 仍然是最灵活的选择。
2.3 宽窄依赖 —— Stage 划分的核心依据
这是面试必考题,也是理解 Spark 执行流程的关键。
窄依赖(Narrow Dependency):父 RDD 的每个 partition 最多被子 RDD 的一个 partition 使用。典型代表:map、filter、union、mapPartitions。
父 RDD: [A][B][C] 子 RDD: [A'][B'][C']
| |
+--一对一映射------+
窄依赖的特点是可以在同一个 Stage 内流水线执行,不需要等待父 RDD 全部算完,父 RDD 一个 partition 算完直接传给子 RDD 的算子继续处理。
宽依赖(Wide Dependency):父 RDD 的 partition 会被子 RDD 的多个 partition 使用。典型代表:groupByKey、reduceByKey、sortByKey、join(非 co-partitioned 情况下)。
父 RDD: [A][B] 子 RDD: [结果1][结果2][结果3][结果4]
| \ / | |
| \ / | |
+----+Shuffle----+----+----+
宽依赖的特点是必须等父 RDD 所有 partition 的数据都准备好,经过 Shuffle 才能继续,所以它是 Stage 的边界线。
2.4 DAG 怎么切 Stage?Task 怎么生成?
DAGScheduler 切 Stage 的规则:从最后一个 RDD 往前推,遇到宽依赖就切一刀。
RDD A -> map -> RDD B -> map -> RDD C -> groupByKey -> RDD D -> map -> RDD E -> collect()
^宽依赖切一刀^ ^Action触发^
Stage 0: A.map().map() (窄依赖流水线)
Stage 1: C.groupByKey() (Shuffle)
Stage 2: D.map() -> E (窄依赖流水线)
Task 的生成:每个 Stage 内,根据 RDD 的 partition 数量生成对应数量的 Task。比如 Stage 0 有 100 个 partition,就会生成 100 个 Task 分发到 Executor 上并行执行。
三、Transformation 和 Action —— 懒执行的秘密
3.1 核心区别
| 类型 | 是否触发执行 | 返回值 | 例子 |
|---|---|---|---|
| Transformation | 否(记录血统) | 新的 RDD | map, filter, flatMap, reduceByKey |
| Action | 是(触发 Job 提交) | 非 RDD 类型(值/集合/Unit) | collect, count, saveAsTextFile, foreach |
3.2 懒执行怎么实现的?
Spark 内部维护了一个 DAG(有向无环图)。每次调用 Transformation,并不会真正跑计算,而是往 DAG 里添加一个节点,记录"这个 RDD 是从哪个父 RDD 通过什么算子变来的"。只有当碰到 Action 时,Spark 才会把这个 DAG 从头到尾扫描一遍,优化后切成 Stage,提交到集群执行。
这个设计的好处:Spark 有机会看到完整的计算图,做全局优化。比如你把 filter 和 map 连续调用,Spark 可以合并成一个 Stage 流水线执行,而不是真的写两遍中间结果。
3.3 常见的坑
坑 1:在 Transformation 里做外部 IO
// 错误示范:Transformation 不会立即执行,你的日志可能打不出来
rdd.map { record =>
println("处理记录: " + record) // 在 Executor 端执行,但你可能在日志里看不到
writeToExternalDB(record) // 这也是懒执行的,不会立刻写
}
正确做法:如果需要确保副作用执行,包在 Action 里或者配合 foreachPartition。
坑 2:连续调 Action 导致重复计算
val rdd = source.map(...).filter(...).groupByKey()
val count = rdd.count() // Action 1:触发完整计算
val result = rdd.collect() // Action 2:又触发一遍完整计算!
两次 Action 之间 RDD 的血统没有缓存,所以 count() 和 collect() 都会从头算一遍。解决方式是在中间加 rdd.cache() 或 rdd.persist()。
四、任务调度全流程 —— 从代码到集群执行
4.1 先搞清楚 5 个角色各自干啥
| 角色 | 核心职责 | 类比 |
|---|---|---|
| Driver | 解析用户代码、构建 DAG 逻辑执行计划、调度 Task、收集结果 | 项目经理,不动手干活,只指挥 |
| Cluster Manager | 管理整个集群的 CPU/内存资源,按需分配启动 Executor | HR,负责招人分配 |
| Executor | 线程池执行 Task,BlockManager 存储中间数据(cache/shuffle),定期汇报心跳 | 干活的工程师 |
| DAGScheduler | 从后往前推 RDD 依赖链,遇到宽依赖(Shuffle)就切一刀,把 Stage 包装成 TaskSet | 任务拆分专家 |
| TaskScheduler | 接收 TaskSet,按"数据本地性"原则分发 Task,监控执行状态,失败自动重试(默认 4 次) | 工单派发系统 |
4.2 完整执行流程(9 步,跟着走一遍就懂了)
Step 1: 用户 spark-submit 提交 Application
→ 比如: spark-submit --class MyApp --master yarn myapp.jar
Step 2: Driver 启动(主控进程,通常跑在集群某台机器的 Container 里)
→ 解析用户代码,构建 DAG 逻辑执行计划
→ 维护 RDD 之间的依赖关系(Lineage)
→ 注意: 此时只是解析,还没有真正执行!
Step 3: Driver 向 Cluster Manager 申请资源
→ CM 可以是 YARN / Kubernetes / Standalone
→ 请求在各个 Worker 节点上启动 Executor 进程
Step 4: Cluster Manager 在 Worker 节点启动 Executor
→ Executor(工作节点)启动后向 Driver 注册
→ Executor 内部有线程池 + BlockManager(存 cache/shuffle 数据)
→ 多个 Worker 节点同时启动,形成分布式执行集群
Step 5: 代码中遇到第一个 Action,触发真正计算
→ Action 比如: collect() / count() / saveAsTextFile()
→ 之前所有 Transformation 都是懒执行的,只有 Action 才触发计算
→ 一个 Action = 一个 Job
Step 6: DAGScheduler 开始切分 Stage
→ 从最后一个 RDD 往前推,遍历依赖链
→ 遇到宽依赖(Shuffle)就切一刀
→ 每个 Stage 包装成一个 TaskSet
→ 示例:
RDD A → map → RDD B → map → RDD C → groupByKey → RDD D → map → RDD E → collect()
窄依赖流水线 宽依赖切Stage 窄依赖流水线
|<───── Stage 0 ──────>|<── Stage 1 ──>|<── Stage 2 ──>|
Step 7: TaskScheduler 分发 Task
→ 接收 DAGScheduler 发来的 TaskSet
→ 按"数据本地性"原则分发(优先发到数据所在节点)
→ 监控 Task 执行状态,失败自动重试(默认 4 次)
Step 8: Executor 执行 Task
→ 线程池里的线程执行具体的 Task
→ 处理完后把结果返回给 Driver
→ 中间数据需要缓存的写入 BlockManager
Step 9: Driver 汇总结果返回用户
→ 收集所有 Executor 返回的 Task 结果
→ 聚合后返回给用户程序
Task 怎么生成:每个 Stage 根据 RDD 的 partition 数量拆分成多个 Task。比如 Stage 0 有 100 个 partition,就会生成 100 个 Task。
4.4 Job / Stage / Task 的关系(面试必背)
记忆口诀:
一个 Action 触发一个 Job
一个 Job 按宽依赖切成多个 Stage
一个 Stage 按分区数拆成多个 Task
一个 Application 里有多少个 Job?取决于代码里调了多少个 Action。比如:
rdd.count() // Job 1
rdd.take(10) // Job 2
rdd.saveAsTextFile("hdfs://...") // Job 3
这三个 Action 会依次触发 3 个 Job。
五、容错机制 —— 节点挂了怎么办
5.1 Lineage 血统重算
RDD 的每个 partition 都记录了它的"家谱"——从数据源开始,经过哪些 Transformation 一步步变过来的。如果某个 partition 的数据丢失(节点挂了),Spark 会根据这个血统信息,找到上游还活着的数据,重新计算这个 partition。
原始数据(HDFS) -> map -> filter -> reduceByKey -> 结果
^ ^
节点挂了! lineage 回溯到 HDFS 重新读取 -> map -> filter -> 恢复
Lineage 的优势:不需要额外存储副本,靠逻辑记录就能恢复,节省存储。
Lineage 的劣势:如果血统链太长,重算代价太大。比如经历了 50 个 Transformation 才到这一步,重算可能要跑很久。
5.2 Checkpoint —— 切断长血统
Checkpoint 的本质是把 RDD 的当前状态持久化到可靠的分布式存储(比如 HDFS、S3),并且切断它和上游的依赖关系。
sparkContext.setCheckpointDir("hdfs:///checkpoints")
val longLineageRDD = ... // 经过了很多 transformation
longLineageRDD.checkpoint() // 触发写入 HDFS
执行 Checkpoint 后,这个 RDD 的血统被截断。如果后续某个 partition 丢失,只需要从 Checkpoint 位置恢复,不需要从头重算。
5.3 Lineage 和 Checkpoint 的配合策略
实际生产中我的建议:
| 场景 | 策略 |
|---|---|
| 血统链短(< 10 个 Transformation) | 纯 Lineage 重算即可 |
| 血统链中等(10-20 个) | 在中间关键节点 persist(StorageLevel.MEMORY_AND_DISK) |
| 血统链长(> 20 个)或迭代计算 | 定期 checkpoint() 切断血统 |
| 数据源头本身不稳定(如外部 API) | 第一步就先 cache 或 checkpoint |
六、Tungsten 优化 —— Spark 3.x 的性能杀手锏
6.1 为什么要做 Tungsten?
传统 Spark(1.x 时代)的一个大问题是JVM 对象开销太大。一个 int 本来 4 字节就够了,但包装成 Integer 对象在 JVM 里要占 16 字节甚至更多。在百亿级数据处理时,内存压力巨大,GC 频繁触发,任务经常因为内存溢出被杀掉。
Tungsten 就是 Spark 团队针对这个问题做的底层重构。
6.2 UnsafeRow:把对象拍扁成二进制
Tungsten 的核心思路:用堆外内存(Off-Heap)存二进制数据,消除 JVM 对象的开销。
传统方式(JVM 对象):
Row { Integer age(16字节), String name(引用+对象), Double score(24字节) }
-> 实际内存占用:100+ 字节
Tungsten UnsafeRow:
[ age(4字节) | name长度(4字节) | name指针(8字节) | score(8字节) ]
-> 紧凑的二进制布局,直接存堆外内存
UnsafeRow 是 Spark SQL/Dataset 的底层存储格式,它的名字里带 “Unsafe” 是因为它绕过了 JVM 的类型安全检查,直接用 sun.misc.Unsafe 操作内存地址,换来了极致的性能。
6.3 三个阶段(Phase 1/2/3)分别做了什么
| Phase | 核心优化 |
|---|---|
| Phase 1 | 内存管理优化,引入 Page 管理堆外内存,替代 JVM GC 分配 |
| Phase 2 | UnsafeRow 紧凑二进制格式,消除对象头和指针开销 |
| Phase 3 | Whole-Stage Codegen,把多个算子编译成一段 JVM 字节码,消除虚函数调用 |
Phase 3 是最狠的。传统执行模式下,每个算子(filter、project、join)都是一个独立的函数调用,数据在算子之间传递时要经过大量的虚函数调用和序列化。而 Whole-Stage Codegen 会把这些算子合并编译成一段连续的 Java 代码,就像手写了一个 for 循环,把 filter、map 全塞进去,编译成字节码直接执行。
你可以这样验证是否触发了 codegen:
df.explain(true) // 看执行计划里的 * 标记,带 * 的表示 codegen 后的阶段
6.4 对 RDD 用户的启示
Tungsten 主要是 DataFrame/Dataset 层面的优化,RDD 用户享受不到 Catalyst 和 Whole-Stage Codegen。但有两个点可以参考:
- 尽量用 DataFrame 替代 RDD 做结构化数据处理,性能差距可能有 10 倍以上
- 如果必须用 RDD,注意控制 partition 内的数据量,避免单个 partition 太大导致内存压力
七、Pipeline / Operator Chain —— 什么时候断开?
7.1 流水线执行的条件
窄依赖的算子会被 Spark 合并到同一个 Task 内流水线执行。比如:
rdd.map(...).filter(...).map(...)
这三个算子会合并成一个 pipeline,在一个 Task 内完成,中间不需要把数据落盘。
7.2 什么情况下 pipeline 会断开?
| 操作 | 是否断开 Pipeline | 原因 |
|---|---|---|
repartition(n) |
一定断 | 触发 Shuffle,Stage 边界 |
coalesce(n) |
可能断 | 如果 shuffle=false 且 n < 当前分区数,不断;否则断 |
groupByKey() |
一定断 | 宽依赖,Shuffle |
sortByKey() |
一定断 | 宽依赖,Shuffle |
checkpoint() |
一定断 | 强制物化,切断 lineage |
7.3 实战技巧:控制 pipeline 长度
Pipeline 太长会导致单个 Task 执行时间过长,如果中途失败重算代价大;Pipeline 太短又会导致 Task 数量过多,调度开销大。
一个经验法则:
- 单个 Stage 内的 pipeline 操作控制在 5-15 个为宜
- 如果超过 15 个,考虑在中间加一个
persist(),把计算成果存下来 - 如果少于 3 个,检查是否分区策略不合理,产生了过多的小 Task
八、线上常见问题与解决方案
8.1 OOM(内存溢出)—— 最头疼的问题
场景 1:单个 partition 数据量太大
报错:java.lang.OutOfMemoryError: Java heap space
或 Executor lost / Container killed by YARN
原因:某个 key 的数据高度倾斜,导致一个 partition 的数据远超其他。
解决思路:
// 1. 先给 RDD 加随机前缀打散倾斜的 key
val prefixRdd = skewedRdd.map { case (key, value) =>
val prefix = Random.nextInt(10)
(s"${prefix}_$key", value)
}
// 2. 做聚合
val aggregated = prefixRdd.reduceByKey(...)
// 3. 去掉前缀二次聚合
val result = aggregated.map { case (prefixedKey, value) =>
val realKey = prefixedKey.substring(prefixedKey.indexOf("_") + 1)
(realKey, value)
}.reduceByKey(...)
场景 2:Broadcast 变量太大
报错:Size exceeds Integer.MAX_VALUE
解决:
// 检查 broadcast 的数据大小,如果超过几 GB,改用 map-side join
// 或者调整 spark.broadcast.blockSize(默认 4MB,可适当调大)
spark.conf.set("spark.broadcast.blockSize", "16m")
8.2 Shuffle 阶段慢 —— 数据倾斜
判断方法:看 Spark UI 的 Stage 页面,如果某个 Task 的执行时间远超其他 Task(比如其他 10 秒,这个 30 分钟),大概率是数据倾斜。
解决:
- 加盐(Salting):给倾斜的 key 加随机前缀,先局部聚合再全局聚合
- 自定义 Partitioner:让数据更均匀地分布
- 两阶段聚合:
reduceByKey比groupByKey更不容易倾斜,因为reduceByKey会在 map 端先做预聚合
8.3 小文件问题
Shuffle 后如果 partition 数设置得太大,会产生大量小文件,HDFS namenode 压力大,后续读取也慢。
解决:
// 输出前合并小文件
result.coalesce(10).saveAsTextFile("hdfs:///output")
// 或者
result.repartition(10).saveAsTextFile("hdfs:///output")
注意:coalesce 只能减少分区数,不能增加;repartition 可以增加也可以减少,但一定会触发 Shuffle。
8.4 Spark UI 看不到或端口冲突
本地调试时 Spark UI 默认开在 4040 端口,如果同时跑了多个 Spark 程序,端口会递增(4041、4042…)。线上环境可以通过 spark.history.server 查看历史任务的 UI。
九、生产环境调参经验
这里列几个我调过的、真正有效果的参数:
| 参数 | 建议值 | 说明 |
|---|---|---|
spark.sql.adaptive.enabled |
true |
Spark 3.x 默认开启 AQE,自动调整分区数,强烈建议开着 |
spark.sql.adaptive.coalescePartitions.enabled |
true |
AQE 自动合并小分区 |
spark.serializer |
org.apache.spark.serializer.KryoSerializer |
比 Java 序列化快得多,内存占用也更小 |
spark.shuffle.service.enabled |
true |
动态分配资源时必开 |
spark.sql.shuffle.partitions |
根据数据量调整,默认 200 通常太小 | 小数据 50,大数据 1000+ |
十、面试核心考点速查表
把前面讲的内容整理成面试问答形式,方便背诵:
Q:RDD 的五大特性?
A:不可变、分区、依赖关系、Transformation 懒执行、持久化。
Q:宽窄依赖区别?
A:窄依赖是父 partition 对应子 partition 一对一/多对一,不经过 Shuffle;宽依赖是父 partition 对应子 partition 一对多,必须经过 Shuffle。DAGScheduler 在宽依赖处切分 Stage。
Q:Transformation 和 Action 的区别?
A:Transformation 记录 lineage 不触发执行,返回新 RDD;Action 触发 runJob,把任务提交到集群执行,返回非 RDD 结果。
Q:一个 Application 有几个 Job?
A:等于代码中 Action 的调用次数。每个 Action 触发一个 Job。
Q:Task / Stage / Job 的关系?
A:一个 Action = 一个 Job;Job 按宽依赖切分成多个 Stage;Stage 按 RDD 的分区数拆分成多个 Task。
Q:Spark 容错怎么实现?
A:主要靠 Lineage 血统重算,配合 Checkpoint 切断长 lineage。Task 失败会自动重试 4 次。
Q:Pipeline 什么时候断开?
A:宽依赖处一定断开;repartition 一定断;coalesce 在不触发 Shuffle 时不断。
Q:Tungsten 优化是什么?
A:Spark 的底层执行优化,包括堆外内存管理(Phase 1)、UnsafeRow 紧凑二进制格式(Phase 2)、Whole-Stage Codegen 生成字节码(Phase 3)。
写在最后
Spark Core 的内容远不止这些,但上面这些覆盖了生产环境最常用的核心知识点。建议配合 Spark UI 反复观察任务执行过程,看多了自然就理解 Stage 怎么切、Task 怎么分的了。
如果这篇文章对你有帮助,记得点赞收藏呀~
更多推荐
所有评论(0)