RDD持久化
RDD 持久化:缓存(快,作业内复用)vs 检查点(稳,跨作业复用)
Spark 默认不存数据,每次 Action 都从头算一遍。如果同一个 RDD 要被用两次,默认行为就是算了两次。
持久化就是”把中间结果存下来,下次直接用”。
两种方式:
| 缓存(Cache / Persist) | 检查点(Checkpoint) | |
|---|---|---|
| 存哪 | 内存 / 本地磁盘 | HDFS(分布式存储) |
| 多快 | 快(内存级) | 慢(要写 HDFS) |
| 血缘关系 | 保留 | 切断 |
| 什么时候用 | 同一个 Job 里多次使用 | 跨 Job 复用,或者要切断长血缘 |
| 最常用方式 | .cache() |
sc.setCheckpointDir(...) + .checkpoint() |
记住:缓存是”快但临时”,检查点是”慢但持久”。两个一起用是黄金组合。
缓存(Cache / Persist):内存里存一份,下次直接用
val rdd = sc.textFile("data.txt")
.flatMap(_.split(" "))
.map((_, 1))
.reduceByKey(_ + _)
rdd.cache() // 标记:需要缓存
rdd.count() // 第一次 Action:算 + 缓存
rdd.collect() // 第二次 Action:直接从缓存读,不重算
.cache() 默认存内存(MEMORY_ONLY)。如果内存不够,数据会丢失,然后重算。
.persist(StorageLevel.XXX) 可以选存储级别:
| 存储级别 | 说明 | 什么时候用 |
|---|---|---|
MEMORY_ONLY |
纯内存,不序列化 | 默认,数据量小用 |
MEMORY_ONLY_SER |
内存,序列化 | 数据量大,用序列化省空间 |
MEMORY_AND_DISK |
内存优先,不够写磁盘 | 数据不确定能不能装下内存 |
MEMORY_AND_DISK_SER |
内存优先,序列化,不够写磁盘 | 数据量大、内存紧张 |
DISK_ONLY |
纯磁盘 | 内存完全不够,或数据不用频繁读 |
缓存的时机: .cache() 只是”标记”,真正缓存发生在第一次 Action 的时候。
缓存的删除:
- 内存不够时 LRU 自动淘汰
- 手动:
.unpersist()
检查点(Checkpoint):写入 HDFS,切断血缘
检查点把 RDD 数据写到 HDFS,同时切断血缘关系——后续 RDD 直接从检查点读数据,不追溯前面的依赖链。
// 先设检查点目录
sc.setCheckpointDir("hdfs://namenode:8020/checkpoint")
val rdd = sc.textFile("data.txt")
.flatMap(_.split(" "))
.map((_, 1))
.reduceByKey(_ + _)
rdd.checkpoint() // 标记检查点
rdd.count() // 触发:写检查点
检查点的特点:
| 特点 | 说明 |
|---|---|
| 存 HDFS | 分布式存储,节点挂了数据还在 |
| 切断血缘 | 依赖链变成从检查点开始,不用追溯前面的长链 |
| 额外算一次 | 检查点会触发一次独立计算(除非配合缓存) |
| 跨 Job 可用 | 同一个检查点文件可以被不同 Spark 应用读 |
为什么不只靠 HDFS 副本? HDFS 是”存”,检查点是”把计算过的结果存下来”。如果只靠血缘,每次都要从头算;靠检查点,直接从存好的结果读。
缓存 + 检查点 = 黄金组合
检查点的问题是:它会额外算一次。如果 RDD 还没算过,checkpoint() 会触发一次完整计算,然后才写 HDFS。
最佳实践:先缓存,再检查点
rdd.cache() // 缓存到内存
rdd.checkpoint() // 标记检查点
rdd.count() // 第一次 Action:算 → 缓存 → 写检查点(从缓存读,不重算)
// 后续 Action:优先从缓存读(快),缓存丢了就从检查点读(稳)
为什么好?
- 缓存给后续操作提供快速访问
- 检查点给缓存丢失提供”后路”
- 检查点不额外重算(从缓存读数据)
缓存 vs 检查点:一张图看懂
| 对比项 | 缓存 | 检查点 |
|---|---|---|
| 存储位置 | 内存 / 本地磁盘 | HDFS |
| 存储持久性 | 临时(Job 结束释放) | 永久(文件存在 HDFS) |
| 血缘关系 | 保留 | 切断 |
| 读取速度 | 快(内存) | 慢(读 HDFS) |
| 是否额外计算 | 否(首次 Action 时顺便缓存) | 是(会触发一次独立计算) |
| 跨 Job 复用 | 否 | 是 |
| 适用场景 | 同一个 Job 内多次使用 | 长血缘链、跨 Job 复用、容灾 |
怎么看缓存和检查点生效了
val rdd = sc.textFile("data.txt").flatMap(_.split(" ")).map((_, 1))
// 查看血缘关系
println(rdd.toDebugString)
// 有 CachedPartitions 说明缓存生效
// 检查点后看血缘
rdd.checkpoint()
rdd.count()
println(rdd.toDebugString)
// 出现 ReliableCheckpointRDD 说明检查点生效
Spark UI 的 Storage 页面会显示缓存了哪些 RDD,占用了多少内存。