spark核心数据结构之RDD
RDD 深度解析:Spark 最核心的数据结构,为什么它这么设计
RDD 是 Spark 最核心的抽象——它是分布式数据集合,也是 Spark 实现”内存计算 + 容错 + 并行”的基础。
RDD 是什么?”Resilient Distributed Dataset”:
- Resilient(弹性):数据丢了能重算,不用备份副本
- Distributed(分布式):数据分散在多台机器上,并行处理
- Dataset(数据集):一个不可变的、分区的元素集合
为什么 Spark 要设计 RDD? 分布式计算面临三个难题:
- 数据怎么分散到多台机器?→ 分区(Partitions)
- 某台机器挂了数据丢了怎么办?→ 血缘关系(Dependencies)
- 怎么知道数据在哪台机器上?→ 首选位置(Preferred Locations)
RDD 的五大属性就是回答这三个问题的
RDD 的五大属性(源码定义)
// Spark 源码注释:每个 RDD 由五个主要属性描述
// - A list of partitions(分区列表)
// - A function for computing each split(分区计算函数)
// - A list of dependencies on other RDDs(依赖关系)
// - Optionally, a Partitioner for key-value RDDs(分区器)
// - Optionally, a list of preferred locations to compute each split on(首选位置)
属性1:分区列表(Partitions)——数据怎么分散
protected def getPartitions: Array[Partition]
RDD 的数据分散在多个分区里,每个分区是数据的一个子集,分布在某个 Executor 上。
| 问题 | RDD 的答案 |
|---|---|
| 数据怎么分散到多台机器? | 切成多个分区,每个分区可以放到不同的节点 |
| 并行度是多少? | 分区数 = 并行度,一个分区对应一个 Task |
| 怎么控制并行度? | repartition(n) 增加分区,coalesce(n) 减少分区 |
举例: 读取 HDFS 上一个 1GB 的文件(块大小 128MB),RDD 有 8 个分区,每个分区对应一个 HDFS 块。
属性2:分区计算函数(Compute)——数据怎么算
def compute(split: Partition, context: TaskContext): Iterator[T]
每个分区怎么计算?由 compute 函数定义。
关键特性:惰性执行。 compute 只记录”要做什么”,不真正算。等到 Action 触发才执行。
compute 返回的是迭代器(Iterator)——数据一条一条地流式处理,不会一次性把整个分区加载到内存。
举例: map(_ * 2) 的 compute 函数就是”对分区里的每个元素乘以 2”。
属性3:依赖关系(Dependencies)——数据丢了怎么办
protected def getDependencies: Seq[Dependency[_]]
每个 RDD 记录了它从哪个父 RDD 转换而来——这就是”血缘关系”(Lineage)。
依赖分两种:
| 依赖类型 | 特点 | 例子 | Stage 处理 |
|---|---|---|---|
| 窄依赖 | 父分区 → 子分区一对一 | map、filter |
可以在同一个 Stage |
| 宽依赖 | 父分区 → 子分区多对多 | reduceByKey、groupByKey |
必须切 Stage,要 Shuffle |
容错原理: 某个分区丢了,沿着血缘关系找到父 RDD,只重算丢失的那几个分区,不用重算整个 RDD。
属性4:分区器(Partitioner)——Key 怎么分布
val partitioner: Option[Partitioner] = None
只有 Key-Value 型 RDD(RDD[(K, V)])才有分区器。它决定 key 怎么分配到各个分区。
常见分区器:
| 分区器 | 规则 | 场景 |
|---|---|---|
HashPartitioner |
key.hashCode() % 分区数 | 默认,按 key 哈希分散 |
RangePartitioner |
按 key 范围划分 | 排序场景 |
| 自定义 | 自己实现 | 业务自定义(比如按地区) |
分区器的作用: 控制 Shuffle 后相同 key 的数据去同一个分区,方便后续聚合。
属性5:首选位置(Preferred Locations)——数据在哪台机器
protected def getPreferredLocations(split: Partition): Seq[String]
数据优先在哪台机器上计算?”移动计算比移动数据便宜”。
调度策略: Task 优先分配到数据所在的节点。
| 本地化级别 | 含义 | 速度 |
|---|---|---|
| PROCESS_LOCAL | 数据在 Executor 内存里 | 最快 |
| NODE_LOCAL | 数据在同一个节点上 | 快 |
| RACK_LOCAL | 数据在同一个机架 | 中等 |
| ANY | 数据在别处 | 最慢 |
举例: 从 HDFS 读数据的 RDD,首选位置就是 HDFS 块的 DataNode 列表。
RDD 的三大特性
1. 弹性(Resilient)
- 血缘关系 → 分区丢了重算,不靠副本
- 内存不够自动写磁盘
- 可以动态调整分区数
2. 分布式(Distributed)
- 数据分片存在多台机器上
- 计算并行执行
3. 不可变(Immutable)
- 一旦创建不能改
- 要改就生成新 RDD
- 简化并发,不用考虑锁
RDD 操作类型:Transformation vs Action
| 类型 | 特点 | 例子 |
|---|---|---|
| Transformation | 惰性执行,返回新 RDD | map、filter、flatMap、reduceByKey |
| Action | 触发执行,返回结果或写存储 | collect、count、saveAsTextFile |
执行时机:
val rdd = sc.textFile("data.txt") // 创建 RDD,还没读数据
val words = rdd.flatMap(_.split(" ")) // 惰性,还没执行
val counts = words.map((_, 1)).reduceByKey(_ + _) // 惰性,还没执行
counts.collect() // ← 触发执行!前面的所有操作开始跑
RDD vs DataFrame vs DataSet
| RDD | DataFrame | DataSet | |
|---|---|---|---|
| 数据类型 | 任意对象 | 结构化数据(Row) | 强类型对象 |
| 优化引擎 | 无(开发者自己优化) | Catalyst 优化器 | Catalyst 优化器 |
| 适用场景 | 复杂逻辑、非结构化数据 | SQL 分析 | 类型安全的结构化处理 |
| 性能 | 中等 | 高 | 高 |