Spark Core 学习:RDD 编程模型与实践
Apache Spark 核心学习项目,涵盖 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
| 对比 | RDD | DataFrame |
|---|---|---|
| 类型 | 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 的设计就会自然理解。