Skip to content

SparkCore 之 Spark 持久化类算子详解

作者:starzy | 日期:2026-07-30
关键词:Spark、RDD 持久化、Cache、Persist、Checkpoint、StorageLevel、大数据性能优化


一、为什么需要持久化(Why)

在大数据计算中,我们常会遇到这样一个场景:某个 RDD 或 DataFrame 被多个下游 Action 算子重复使用。默认情况下,Spark 的 RDD 是惰性计算不存储中间结果的——每次 Action 触发时,都会从源头重新计算整个 DAG 链路。

Spark 持久化机制架构图

这不仅浪费计算资源,在数据量大时还会导致严重的性能问题。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 个布尔参数组合而成:

scala
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_ONLY1最快,适合内存充裕场景
MEMORY_ONLY_SER1内存有限时推荐,节省空间
MEMORY_AND_DISK1RDD 默认,平衡性能与可靠性
MEMORY_AND_DISK_SER1空间效率最高的安全方案
DISK_ONLY1内存紧张时的降级方案
MEMORY_AND_DISK_DESER1Dataset 默认(Spark 2.x+)
OFF_HEAP1Tachyon/Alluxio 场景

3.3 序列化 vs 反序列化:空间换时间

  • 反序列化(Java Object):直接存储对象,读取快,但占空间大(通常 2~5 倍)
  • 序列化(Byte Array):存储字节数组,节省空间,但每次读取需反序列化(CPU 开销)
scala
// 反序列化存储 — 速度快但占内存大
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() — 快速缓存

scala
// RDD 的 cache() 实现
def cache(): this.type = persist()
// persist() 默认使用 MEMORY_AND_DISK(RDD)
java
// 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() — 灵活缓存

scala
// 各种持久化策略
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() — 手动释放

scala
// 立即释放(阻塞)
rdd.unpersist()
// 异步释放
rdd.unpersist(blocking = false)

什么时候该调用 unpersist()

  1. 确定后续不再使用该 RDD/DataFrame
  2. 内存压力较大时,主动释放
  3. 在循环/迭代中使用,避免内存泄漏
scala
// 迭代场景:及时释放旧版本
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 / PersistCheckpoint
存储位置各 Executor 本地可靠存储(HDFS / S3)
截断血缘(Lineage)❌ 不截断✅ 截断
容错方式丢失后重算丢失后从检查点恢复
适用场景加速重复计算截断过长的血缘、迭代安全
执行时机第一次 Action 时额外启动一个 Job

5.2 最佳实践:先 Cache 再 Checkpoint

scala
// ⚠️ 错误做法:直接 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

scala
// 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 多表 JoinMEMORY_AND_DISKSQL 优化器可利用缓存

七、源码视角:持久化机制底层实现

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 序列化优化

scala
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)

核心要点

  1. 持久化的本质:用存储空间换计算时间,避免 RDD 血缘的重复执行
  2. cache() = persist(MEMORY_AND_DISK):是简化版,适合大多数场景
  3. persist() 提供 12 种 StorageLevel:可根据内存和数据量灵活选择
  4. Checkpoint 截断血缘:适用于迭代计算和超长血缘场景
  5. 先 Cache 再 Checkpoint:避免重复计算
  6. 启用 Kryo + MEMORY_ONLY_SER:性价比最高的持久化方案
  7. 及时 unpersist:避免内存泄漏

选型速查

条件推荐
内存充足,追求极致性能MEMORY_ONLY
内存有限,数据量中等MEMORY_ONLY_SER + Kryo
数据量大,需要稳定性MEMORY_AND_DISK_SER + Kryo
迭代计算MEMORY_AND_DISK + 周期性 checkpoint
默认选择(不确定时)cache()

作者:starzy
博客blog.starzy.cn
GitHubstarzy1990.github.io
专注领域:AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践