跳到正文
Joeplover
学习笔记·2025-10-13·约 3 分钟阅读

Spark Core 学习:RDD 编程模型与实践

Apache Spark 核心学习项目,涵盖 Spark 环境搭建、RDD 创建方式(内存/磁盘)、分区控制等基础操作。

Spark RDD

Spark Core:RDD 编程模型与实践

RDD(弹性分布式数据集)

RDD 是 Spark 的核心抽象。虽然现在大部分场景直接用 DataFrame(API 更友好、性能更好),但理解 RDD 是深入理解 Spark 的基础。

RDD 三大特性

1. 不可变

RDD 一旦创建就不能修改。对 RDD 的所有操作(map、filter、flatMap)都生成新的 RDD,原始 RDD 不变。

rdd = sc.parallelize([1, 2, 3, 4, 5])
doubled = rdd.map(lambda x: x * 2)   # 生成新 RDD,原 RDD 不变

不可变性带来的好处:缓存安全、故障恢复简单(重新计算即可)。

2. 可分区

RDD 的数据分布在多个分区(Partition)上,每个分区可以在不同节点上并行处理。

# 设置并行度
rdd = sc.parallelize(data, numSlices=4)  # 4 个分区
rdd.getNumPartitions()  # 返回 4

分区数决定了并行度。太少浪费资源,太多调度开销大。通常建议每个分区 128MB-256MB 数据。

3. 弹性(可恢复)

Spark 通过 Lineage(血缘关系) 实现容错。每个 RDD 记录了自己怎么来的:

HDFS 文件 → textFile() → 扁平化后 → map → filter
                                      ↓
                     记录转换操作的 DAG 图

如果一个分区丢失,Spark 根据 Lineage 重新计算这个分区,不需要重算整个数据集。

Transformation vs Action

这是 Spark 最重要的概念区分:

# Transformation(惰性执行):构建 DAG,不触发计算
rdd1 = rdd.filter(lambda x: x > 0)     # 还没做
rdd2 = rdd1.map(lambda x: (x, x * 2))  # 还没做
# 此时 Spark 只是记录了要做什么

# Action:触发 DAG 的实际执行
rdd2.collect()      # 触发计算,收集结果回 driver
rdd2.count()        # 触发计算,返回元素数量
rdd2.saveAsTextFile("output")  # 触发计算,写入 HDFS

Transformation 构建执行计划(DAG),Action 触发执行。这个设计让 Spark 可以对多个操作进行优化(如合并流水线)。

cache() 缓存

# 标记 RDD 为持久化
cached = rdd.cache()  # 等价于 persist(StorageLevel.MEMORY_ONLY)

# 第一次计算时存入内存
cached.count()  # 慢:从源头计算并缓存

# 第二次直接使用缓存
cached.count()  # 快:从内存读取

缓存适用于多次使用同一个 RDD 的场景。如果只用一次,缓存反而浪费资源。

RDD vs DataFrame

对比RDDDataFrame
类型JVM 对象,运行时类型有 Schema,编译时类型
API函数式(map/filter)SQL 风格 + 函数式
优化依赖开发者优化Catalyst 优化器自动优化
性能基准线通常快 2-5 倍
适用场景自定义复杂逻辑数据分析(90% 场景)

DataFrame 的性能优势来自 Catalyst 优化器——列裁剪(只读需要的列)、谓词下推(过滤在数据源执行)、自动选择 Join 策略等。相同逻辑,DataFrame 通常比 RDD 快 2-5 倍。

总结

RDD 是理解 Spark 一切的基础。即使现在 DataFrame 是主流,RDD 的核心概念(不可变、分区、惰性计算、血缘)贯穿整个 Spark 生态。理解了 RDD,DataFrame 和 Structured Streaming 的设计就会自然理解。