SparkCore 之 Spark 持久化类算子详解
作者:starzy | 日期:2026-07-30
关键词:Spark、RDD 持久化、Cache、Persist、Checkpoint、StorageLevel、大数据性能优化
一、为什么需要持久化(Why)
在大数据计算中,我们常会遇到这样一个场景:某个 RDD 或 DataFrame 被多个下游 Action 算子重复使用。默认情况下,Spark 的 RDD 是惰性计算且不存储中间结果的——每次 Action 触发时,都会从源头重新计算整个 DAG 链路。
这不仅浪费计算资源,在数据量大时还会导致严重的性能问题。Spark 的**持久化(Persistence)**机制正是为解决这一问题而设计的——将中间结果缓存到内存或磁盘,供后续 Action 直接复用。
一句话总结:持久化就是让你的 RDD/DataFrame "记住"自己的计算结果,避免重复计算。
二、持久化核心概念(What)
2.1 持久化类算子全景
Spark 提供了三个核心持久化算子:
| 算子 | 本质 | 是否立即执行 | 默认存储级别 |
|---|---|---|---|
cache() | persist() 的简化版 | ❌ 懒执行 | MEMORY_AND_DISK(RDD)MEMORY_AND_DISK_DESER(Dataset) |
persist() | 真正的持久化入口 | ❌ 懒执行 | 可指定任意 StorageLevel |
unpersist() | 从缓存中移除 | ✅ 立即执行 | — |
关键认知:
cache()和persist()都是懒执行的 Transformation 算子,真正的缓存动作发生在第一个 Action 被触发时。
2.2 为何不是 Action?
很多人会困惑:cache() 既然要"存储数据",为什么设计为 Transformation 而不是 Action?
答案在于 Spark 的 DAG 优化。如果 cache() 是 Action,它会破坏 DAG 的完整性,导致 Spark 无法做全局优化。将缓存标记为 Transformation,Spark 可以在构建执行计划时就将"缓存点"纳入规划,实现更优的 Stage 划分。
三、StorageLevel 详解(What)
StorageLevel 是 Spark 持久化的核心配置,决定了数据的存储位置(内存/磁盘/堆外)和序列化方式。
3.1 StorageLevel 五维度模型
每个 StorageLevel 由 5 个布尔参数组合而成:
class StorageLevel private(
private var _useDisk: Boolean, // 是否使用磁盘
private var _useMemory: Boolean, // 是否使用堆内存
private var _useOffHeap: Boolean, // 是否使用堆外内存
private var _deserialized: Boolean, // 是否以反序列化(对象)形式存储
private var _replication: Int = 1 // 副本因子(默认1)
)3.2 完整 StorageLevel 对照表
| StorageLevel | 内存 | 磁盘 | 堆外 | 序列化 | 副本 | 典型场景 |
|---|---|---|---|---|---|---|
MEMORY_ONLY | ✅ | ❌ | ❌ | ❌ | 1 | 最快,适合内存充裕场景 |
MEMORY_ONLY_SER | ✅ | ❌ | ❌ | ✅ | 1 | 内存有限时推荐,节省空间 |
MEMORY_AND_DISK | ✅ | ✅ | ❌ | ❌ | 1 | RDD 默认,平衡性能与可靠性 |
MEMORY_AND_DISK_SER | ✅ | ✅ | ❌ | ✅ | 1 | 空间效率最高的安全方案 |
DISK_ONLY | ❌ | ✅ | ❌ | ✅ | 1 | 内存紧张时的降级方案 |
MEMORY_AND_DISK_DESER | ✅ | ✅ | ❌ | ❌ | 1 | Dataset 默认(Spark 2.x+) |
OFF_HEAP | ❌ | ❌ | ✅ | ✅ | 1 | Tachyon/Alluxio 场景 |
3.3 序列化 vs 反序列化:空间换时间
- 反序列化(Java Object):直接存储对象,读取快,但占空间大(通常 2~5 倍)
- 序列化(Byte Array):存储字节数组,节省空间,但每次读取需反序列化(CPU 开销)
// 反序列化存储 — 速度快但占内存大
rdd.persist(StorageLevel.MEMORY_ONLY)
// 序列化存储 — 省内存但读取有 CPU 开销
rdd.persist(StorageLevel.MEMORY_ONLY_SER)经验法则:如果内存紧张,
MEMORY_ONLY_SER是比MEMORY_AND_DISK更好的选择——因为磁盘 I/O 远慢于反序列化 CPU 开销。
四、核心算子深入解析(How)
4.1 cache() — 快速缓存
// RDD 的 cache() 实现
def cache(): this.type = persist()
// persist() 默认使用 MEMORY_AND_DISK(RDD)// Java 示例
JavaRDD<String> lines = sc.textFile("hdfs://data/logs.txt");
JavaRDD<String> filtered = lines.filter(line -> line.contains("ERROR"));
filtered.cache(); // 标记缓存
long count1 = filtered.count(); // ← 第一个 Action:触发计算 + 缓存
long count2 = filtered.count(); // ← 直接读缓存,不再计算4.2 persist() — 灵活缓存
// 各种持久化策略
rdd.persist(StorageLevel.MEMORY_ONLY) // 仅内存
rdd.persist(StorageLevel.MEMORY_ONLY_SER) // 内存序列化(省空间)
rdd.persist(StorageLevel.MEMORY_AND_DISK) // 内存+磁盘(RDD 默认)
rdd.persist(StorageLevel.MEMORY_AND_DISK_SER) // 内存+磁盘+序列化(推荐)
rdd.persist(StorageLevel.DISK_ONLY) // 仅磁盘(内存紧张时)4.3 unpersist() — 手动释放
// 立即释放(阻塞)
rdd.unpersist()
// 异步释放
rdd.unpersist(blocking = false)什么时候该调用 unpersist()?
- 确定后续不再使用该 RDD/DataFrame
- 内存压力较大时,主动释放
- 在循环/迭代中使用,避免内存泄漏
// 迭代场景:及时释放旧版本
var current = initialRdd
for (i <- 1 to 100) {
val next = current.map(compute)
next.persist()
current.unpersist() // 释放旧版本
current = next
}五、Checkpoint — 持久化的终极形态
5.1 Checkpoint 与 Cache/Persist 的本质区别
| 特性 | Cache / Persist | Checkpoint |
|---|---|---|
| 存储位置 | 各 Executor 本地 | 可靠存储(HDFS / S3) |
| 截断血缘(Lineage) | ❌ 不截断 | ✅ 截断 |
| 容错方式 | 丢失后重算 | 丢失后从检查点恢复 |
| 适用场景 | 加速重复计算 | 截断过长的血缘、迭代安全 |
| 执行时机 | 第一次 Action 时 | 额外启动一个 Job |
5.2 最佳实践:先 Cache 再 Checkpoint
// ⚠️ 错误做法:直接 Checkpoint — 计算了 2 次!
rdd.checkpoint()
rdd.count()
// ✅ 正确做法:先 Cache 再 Checkpoint
rdd.cache()
rdd.checkpoint()
rdd.count() // 第1次:计算→缓存→写入Checkpoint为什么? Checkpoint 会启动额外 Job 写入 HDFS。如果不先 Cache,这个额外 Job 会重新计算整个 RDD。
5.3 迭代计算中的 Checkpoint
// PageRank 迭代 — 定期 Checkpoint 截断血缘
sc.setCheckpointDir("hdfs://namenode:8020/spark/checkpoint")
var ranks = initialRanks
for (i <- 1 to 100) {
val contribs = links.join(ranks).values.flatMap { ... }
ranks = contribs.reduceByKey(_ + _).mapValues(0.15 + 0.85 * _)
// 每 10 次迭代做一次 Checkpoint
if (i % 10 == 0) {
ranks.cache()
ranks.checkpoint()
ranks.count()
}
}六、持久化策略选择指南(Best Practice)
6.1 决策树
数据会被多次使用吗?
├── 否 → 不需要持久化
└── 是 → 内存够吗?
├── 够 → 使用 MEMORY_ONLY(最快)
│ └── 内存不足导致溢出? → 降级为 MEMORY_ONLY_SER
└── 不够 → 数据量特别大?
├── 是 → 使用 MEMORY_AND_DISK_SER(最安全)
└── 否 → 使用 MEMORY_ONLY_SER(避免磁盘 I/O)6.2 场景化推荐
| 场景 | 推荐策略 | 原因 |
|---|---|---|
| ETL 中的中间结果复用 | MEMORY_AND_DISK_SER | 数据量大,需要稳定性 |
| 机器学习迭代计算 | MEMORY_AND_DISK + checkpoint | 迭代多次使用 + 截断血缘 |
| 仪表盘查询加速 | MEMORY_ONLY | 延迟敏感,内存充裕 |
| 探索性 Ad-hoc 分析 | cache()(默认) | 简单省心 |
| DataFrame 多表 Join | MEMORY_AND_DISK | SQL 优化器可利用缓存 |
七、源码视角:持久化机制底层实现
7.1 缓存写入流程
Task 执行完毕
→ BlockManager.doPutIterator()
→ MemoryStore.putIteratorAsValues() ← 反序列化存储
└── 内存不足? → 驱逐旧的缓存块(LRU)
→ DiskStore.put() ← 序列化存储
└── 写入本地磁盘文件(DiskBlockManager 管理)7.2 缓存读取流程
RDD.iterator() 被调用
→ BlockManager.getOrElseUpdate()
→ 查本地 MemoryStore → 命中?返回
→ 查本地 DiskStore → 命中?返回
→ 远程拉取 → 命中?返回
→ 全部未命中 → 重新计算八、常见问题与避坑(FAQ)
Q1: cache() 后数据一定在内存中吗?
不一定。MEMORY_ONLY 级别的数据放不进内存,会直接丢弃,下次重算。
Q2: persist() 可以在 Action 之后调用吗?
可以,但没有意义——它是懒执行的,后续没有 Action 触发就不会生效。
Q3: cache() 和 persist() 能重复调用改变存储级别吗?
不能。一旦标记了非 NONE 的 StorageLevel,再次调用会抛出异常。
Q4: Checkpoint 后还需要 Cache 吗?
Checkpoint 完成后,缓存可能已过期。后续读取会从 HDFS 获取,可以 unpersist() 释放内存。
九、性能调优实战
9.1 Spark UI 监控
在 Spark UI 的 Storage 标签页查看:
- Cached Partitions:已缓存的 Partition 数
- Fraction Cached:缓存比例(< 100% → 内存不足)
- Size in Memory / Disk:内存/磁盘占用
9.2 Kryo 序列化优化
val conf = new SparkConf()
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.registerKryoClasses(Array(classOf[MyDataClass]))
rdd.persist(StorageLevel.MEMORY_ONLY_SER)Kryo 比 Java 原生序列化快 10 倍且体积更小。
十、总结(Summary)
核心要点
- 持久化的本质:用存储空间换计算时间,避免 RDD 血缘的重复执行
- cache() = persist(MEMORY_AND_DISK):是简化版,适合大多数场景
- persist() 提供 12 种 StorageLevel:可根据内存和数据量灵活选择
- Checkpoint 截断血缘:适用于迭代计算和超长血缘场景
- 先 Cache 再 Checkpoint:避免重复计算
- 启用 Kryo + MEMORY_ONLY_SER:性价比最高的持久化方案
- 及时 unpersist:避免内存泄漏
选型速查
| 条件 | 推荐 |
|---|---|
| 内存充足,追求极致性能 | MEMORY_ONLY |
| 内存有限,数据量中等 | MEMORY_ONLY_SER + Kryo |
| 数据量大,需要稳定性 | MEMORY_AND_DISK_SER + Kryo |
| 迭代计算 | MEMORY_AND_DISK + 周期性 checkpoint |
| 默认选择(不确定时) | cache() |
作者:starzy
博客:blog.starzy.cn
GitHub:starzy1990.github.io
专注领域:AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践