RDD编程
Spark RDD 编程:创建 RDD → 转换(惰性)→ 行动(触发),三步走
RDD 编程的核心流程就三步:
| 步骤 | 做什么 | 关键方法 |
|---|---|---|
| 1. 创建 RDD | 从数据源生成 RDD | parallelize / textFile |
| 2. 转换 RDD | 数据处理,惰性执行 | map / filter / reduceByKey |
| 3. 行动触发 | 触发计算,返回结果或写存储 | collect / count / saveAsTextFile |
记住:转换(Transformation)懒,行动(Action)勤。转换只是记下”要做什么”,行动才真正开始跑。
创建 RDD:从集合或文件来
1.1 从内存集合创建(测试用)
val rdd1 = sc.parallelize(List(1, 2, 3, 4, 5))
val rdd2 = sc.makeRDD(List(1, 2, 3, 4, 5)) // 推荐,简洁
parallelize 和 makeRDD 功能一样。numSlices 参数控制分区数,默认是 CPU 核心数。
1.2 从外部文件创建(生产用)
// 读本地文件(加 file:// 前缀)
val lines = sc.textFile("file:///path/to/file.txt")
// 读 HDFS 文件
val lines = sc.textFile("hdfs:///user/data/input.txt")
textFile 按行读取,每行作为 RDD 的一个元素。分区数默认跟 HDFS 块数一致(128MB/块)。
转换算子:数据处理的”流水线”
转换算子返回新 RDD,惰性执行。分三类:单值型、双值型、键值型。
2.1 单值型:一个 RDD → 另一个 RDD
| 算子 | 做什么 | 什么时候用 |
|---|---|---|
map |
每个元素一对一转换 | 对每个元素做变换,比如 _ * 2 |
flatMap |
每个元素 → 多个元素(压平) | 拆分字符串、展开嵌套列表 |
filter |
保留满足条件的元素 | 过滤数据,比如去掉空行 |
mapPartitions |
按分区批量处理 | 每个分区创建一次连接(比 map 高效) |
distinct |
去重 | 去掉重复元素 |
coalesce |
减少分区(不 Shuffle) | 合并小分区,减少 Task 数 |
repartition |
调整分区(会 Shuffle) | 增加分区或均衡数据分布 |
示例:

