spark核心数据结构之RDD

RDD 深度解析:Spark 最核心的数据结构,为什么它这么设计

RDD 是 Spark 最核心的抽象——它是分布式数据集合,也是 Spark 实现”内存计算 + 容错 + 并行”的基础。

RDD 是什么?”Resilient Distributed Dataset”:

  • Resilient(弹性):数据丢了能重算,不用备份副本
  • Distributed(分布式):数据分散在多台机器上,并行处理
  • Dataset(数据集):一个不可变的、分区的元素集合

为什么 Spark 要设计 RDD? 分布式计算面临三个难题:

  1. 数据怎么分散到多台机器?→ 分区(Partitions)
  2. 某台机器挂了数据丢了怎么办?→ 血缘关系(Dependencies)
  3. 怎么知道数据在哪台机器上?→ 首选位置(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 处理
窄依赖 父分区 → 子分区一对一 mapfilter 可以在同一个 Stage
宽依赖 父分区 → 子分区多对多 reduceByKeygroupByKey 必须切 Stage,要 Shuffle

容错原理: 某个分区丢了,沿着血缘关系找到父 RDD,只重算丢失的那几个分区,不用重算整个 RDD。

属性4:分区器(Partitioner)——Key 怎么分布

@transient 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 mapfilterflatMapreduceByKey
Action 触发执行,返回结果或写存储 collectcountsaveAsTextFile

执行时机:

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 分析 类型安全的结构化处理
性能 中等