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,占用了多少内存。