RDD 深入
概述
RDD(Resilient Distributed Dataset,弹性分布式数据集)是 Spark 1.x 的基石抽象,理解 RDD 就理解了 Spark 的容错与执行模型。本文讲透 RDD 五大特性、依赖关系、血统容错与缓存/Checkpoint 机制。
一、什么是 RDD
RDD 是只读、分区、可并行操作的数据集合抽象,代表一个中间计算状态。关键特征:
| 特征 | 说明 |
|---|---|
| 只读 | 不可修改,只能生成新 RDD |
| 分区 | 数据分布在集群多个分区 |
| 弹性 | 分区丢失可从血统重算 |
| 惰性 | 只有 Action 才触发计算 |
二、RDD 五大特性
| 特性 | 方法 | 含义 |
|---|---|---|
| 分区列表 | getPartitions | 数据切分为 N 个分区 |
| 分区计算函数 | compute | 对每个分区执行的计算 |
| 依赖关系 | getDependencies | 父 RDD 依赖(窄/宽) |
| 分区器 | partitioner | 键值 RDD 的分区规则(Hash/Range) |
| 首选位置 | getPreferredLocations | 数据本地性提示(哪个节点有数据) |
五大特性决定了 Spark 如何划分任务、如何容错、如何调度。
2.1 分区与并行度
RDD 分区数 → Stage 的 Task 数 → 并行度
数据源分区:HDFS 块数、文件切片、集合分区数
算子可重分区:repartition / coalesce| 分区太多 | 分区太少 |
|---|---|
| 任务调度开销大 | 并行度不足,浪费资源 |
2.2 数据本地性(首选位置)
getPreferredLocations 返回分区数据所在节点,TaskScheduler 优先把 Task 派到这些节点,减少网络传输。
三、依赖关系
3.1 窄依赖 vs 宽依赖
| 依赖 | 含义 | 算子示例 |
|---|---|---|
| 窄依赖 | 每个父分区只被一个子分区使用,可管道式计算 | map、filter、union |
| 宽依赖 | 每个父分区被多个子分区使用,需 Shuffle | groupByKey、reduceByKey、join |
窄依赖:父1 → 子1 (一个分区算完即走,无需等待)
父2 → 子2
宽依赖:父1 → 子1、子2、子3(父分区数据要分发到多个子分区,Shuffle)
父2 → 子1、子2、子33.2 依赖与 Stage 划分
| 依赖 | 对 Stage 的影响 |
|---|---|
| 窄依赖 | 不切 Stage,父 Stage 内管道式计算 |
| 宽依赖 | 产生 Shuffle 边界,切出新 Stage |
map → filter → groupByKey → map → reduce
Stage1(map/filter) Stage2(map/reduce)
└──── 宽依赖 Shuffle ────┘四、血统 Lineage 与容错
4.1 血统机制
RDD 记录自己的父 RDD 依赖链(Lineage)。分区数据丢失时,沿血统从源头重算:
rdd3 = rdd2.map(...) ← 血统记录 rdd3 → rdd2 → rdd1 → 数据源
rdd2 = rdd1.filter(...)
rdd1 = sc.textFile(...)| 优点 | 缺点 |
|---|---|
| 无需复制数据,容错成本低 | 依赖链过长时重算代价大 |
4.2 血统 vs 复制
| 容错方式 | 机制 | 代价 |
|---|---|---|
| 血统重算 | 沿依赖链重算丢失分区 | CPU 重算 |
| 复制 | 每分区多份(如 MR 3 副本) | 存储 3 倍 |
| 缓存/Checkpoint | 保存中间结果,断链 | 存储或序列化成本 |
Spark 默认用血统 + 缓存,Checkpoint 用于切断长链。
五、缓存与 Checkpoint
5.1 缓存(Cache/Persist)
把 RDD 数据保存在内存/磁盘,后续 Action 复用,避免重复计算:
scala
rdd.cache() // 默认 MEMORY_ONLY
rdd.persist(StorageLevel.MEMORY_AND_DISK)| StorageLevel | 说明 |
|---|---|
| MEMORY_ONLY | 只存内存,放不下则丢弃重算 |
| MEMORY_AND_DISK | 内存放不下落盘,防重算 |
| MEMORY_AND_DISK_SER | 序列化存,省内存费 CPU |
| DISK_ONLY | 只落盘,适合超大数据 |
| OFF_HEAP | 堆外内存(Tachyon) |
5.2 缓存与血统的关系
缓存是断点保护:被缓存的分区丢失时直接从缓存恢复,不必沿血统重算全部。
5.3 Checkpoint
Checkpoint 把 RDD 数据写到可靠存储(HDFS),并切断血统:
| 对比 | 缓存 | Checkpoint |
|---|---|---|
| 存储 | 内存/本地磁盘 | HDFS(跨节点可靠) |
| 血统 | 保留 | 切断 |
| 生命周期 | Executor 结束即失效 | 跨作业持久 |
| 使用场景 | 多次复用 | 长依赖链、迭代计算 |
scala
sc.setCheckpointDir("hdfs:///spark-checkpoint")
rdd.checkpoint()5.4 使用建议
| 场景 | 手段 |
|---|---|
| 一个 RDD 被多个 Action 复用 | cache |
| 血统链特别长 | checkpoint |
| 迭代计算(如机器学习) | checkpoint + cache |
| 中间结果要跨作业复用 | 直接写表/文件更稳 |
六、算子分类
6.1 转换算子(Transformation,惰性)
| 类型 | 示例 |
|---|---|
| 单值 | map、flatMap、filter、sample |
| 键值 | mapValues、reduceByKey、groupByKey、sortByKey |
| 分区 | coalesce、repartition |
| 集合 | union、intersection、distinct |
6.2 行动算子(Action,触发计算)
| 示例 | 作用 |
|---|---|
| count、collect、take | 统计与取样 |
| reduce、fold、aggregate | 聚合 |
| saveAsTextFile、saveAsParquetFile | 输出 |
| foreach | 遍历 |
惰性 + 血统:转换只记录依赖,Action 才真正计算,这是 Spark 能整体优化 DAG 的基础。
七、常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 缓存未生效 | 缓存是惰性的,需 Action 触发后才真正缓存 |
| 内存不足频繁重算 | 用 MEMORY_AND_DISK 或调大内存 |
| 血统链过长 OOM | 中间 checkpoint 断链 |
| repartition 造成 Shuffle | 缩小分区用 coalesce(无 Shuffle) |
| collect 大数据 OOM | 改 take/save 到文件 |