这篇笔记整理了用 Spark 的一些心得,覆盖了面试常考的核心知识点,也包括了不少线上真实踩过的坑。


一、Spark Core 到底是什么

先给一个我自己的理解:Spark Core 是整个 Spark 生态的"发动机"。不管是后面你用的 Spark SQL、Spark Streaming 还是 MLlib,底层的执行引擎都是 Core 这一套东西。它解决的核心问题就两个:

  1. 如何把一个大计算任务拆成小块,分布到多台机器上并行跑
  2. 跑的过程中某个节点挂了,怎么自动恢复而不是从头再来

这两个问题分别对应 分布式计算容错机制,也是 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 使用。典型代表:mapfilterunionmapPartitions

父 RDD: [A][B][C]        子 RDD: [A'][B'][C']
         |                  |
         +--一对一映射------+

窄依赖的特点是可以在同一个 Stage 内流水线执行,不需要等待父 RDD 全部算完,父 RDD 一个 partition 算完直接传给子 RDD 的算子继续处理。

宽依赖(Wide Dependency):父 RDD 的 partition 会被子 RDD 的多个 partition 使用。典型代表:groupByKeyreduceByKeysortByKeyjoin(非 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 有机会看到完整的计算图,做全局优化。比如你把 filtermap 连续调用,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。但有两个点可以参考:

  1. 尽量用 DataFrame 替代 RDD 做结构化数据处理,性能差距可能有 10 倍以上
  2. 如果必须用 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 分钟),大概率是数据倾斜。

解决:

  1. 加盐(Salting):给倾斜的 key 加随机前缀,先局部聚合再全局聚合
  2. 自定义 Partitioner:让数据更均匀地分布
  3. 两阶段聚合reduceByKeygroupByKey 更不容易倾斜,因为 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 怎么分的了。


如果这篇文章对你有帮助,记得点赞收藏呀~

更多推荐