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) | 增加分区或均衡数据分布 |
示例:
// map:每个元素翻倍
val rdd = sc.makeRDD(List(1, 2, 3, 4))
rdd.map(_ * 2).collect() // [2, 4, 6, 8]
// flatMap:拆分句子成单词
val lines = sc.makeRDD(List("hello world", "hello spark"))
lines.flatMap(_.split(" ")).collect() // [hello, world, hello, spark]
// mapPartitions:每个分区创建一次连接
rdd.mapPartitions(iter => {
val conn = getConnection() // 一个分区一次
iter.map(record => process(record, conn))
})
// coalesce:合并分区(3个→2个,不 Shuffle)
rdd.coalesce(2)
// repartition:重分区(3个→5个,会 Shuffle)
rdd.repartition(5)
2.2 双值型:两个 RDD → 一个新 RDD
| 算子 | 做什么 | 什么时候用 |
|---|---|---|
union |
并集(不去重) | 合并两个数据集 |
intersection |
交集(去重) | 找两个数据集的共同元素 |
subtract |
差集 | 找 A 有 B 没有的元素 |
示例:
val rdd1 = sc.makeRDD(List(1, 2, 3))
val rdd2 = sc.makeRDD(List(3, 4, 5))
rdd1.union(rdd2).collect() // [1,2,3,3,4,5]
rdd1.intersection(rdd2).collect() // [3]
rdd1.subtract(rdd2).collect() // [1,2]
2.3 键值型:Pair RDD 专用
针对 RDD[(K, V)],Spark 提供了丰富的聚合和连接算子。
| 算子 | 做什么 | 什么时候用 |
|---|---|---|
reduceByKey |
按 key 聚合(Map 端预聚合) | 推荐,比 groupByKey 快 |
groupByKey |
按 key 分组(无预聚合) | 需要保留所有 value 时用 |
aggregateByKey |
自定义分区内和分区间聚合 | 逻辑复杂时用 |
join |
内连接 | 两个表按 key 关联 |
leftOuterJoin |
左外连接 | 保留左表所有 key |
示例:
val rdd = sc.makeRDD(List(("a", 1), ("a", 2), ("b", 3), ("b", 4)))
// reduceByKey:按 key 求和(Map 端预聚合)
rdd.reduceByKey(_ + _).collect() // [(a,3), (b,7)]
// groupByKey:按 key 分组
rdd.groupByKey().mapValues(_.toList).collect() // [(a,[1,2]), (b,[3,4])]
// join:两个 RDD 按 key 连接
val rdd2 = sc.makeRDD(List(("a", "x"), ("b", "y")))
rdd.join(rdd2).collect() // [(a,(1,x)), (a,(2,x)), (b,(3,y)), (b,(4,y))]
行动算子:触发计算
行动算子触发整个作业执行。没有行动,转换永远不会真正计算。
| 算子 | 做什么 | 什么时候用 |
|---|---|---|
collect |
把数据拉到 Driver(数组) | 小数据集调试用 |
count |
返回元素总数 | 计数 |
take(n) |
取前 n 个元素 | 采样查看 |
reduce |
全局聚合 | 求和、求最大值 |
foreach |
遍历每个元素(分布式) | 每个元素执行操作(写外部系统) |
saveAsTextFile |
保存到文件 | 输出结果 |
示例:
val rdd = sc.makeRDD(List(1, 2, 3, 4, 5))
rdd.collect() // [1,2,3,4,5]
rdd.count() // 5
rdd.take(3) // [1,2,3]
rdd.reduce(_ + _) // 15
rdd.foreach(println) // 分布式打印(各 Executor 打印,Driver 看不到)
rdd.saveAsTextFile("hdfs:///output")
两个重要对比
map vs flatMap
| map | flatMap | |
|---|---|---|
| 输入 | 1 个元素 | 1 个元素 |
| 输出 | 1 个元素 | 0 个或多个元素 |
| 结果维度 | 不变 | 可能变(压平) |
val rdd = sc.makeRDD(List("hello world", "hello spark"))
rdd.map(_.split(" ")).collect() // [Array(hello, world), Array(hello, spark)]
rdd.flatMap(_.split(" ")).collect() // [hello, world, hello, spark]
reduceByKey vs groupByKey
| reduceByKey | groupByKey | |
|---|---|---|
| Map 端预聚合 | 有 | 无 |
| Shuffle 数据量 | 小 | 大 |
| 性能 | 快 | 慢 |
| 适用场景 | 聚合(求和、计数) | 需要保留所有 value |
// 推荐用 reduceByKey
rdd.reduceByKey(_ + _)
// groupByKey 会 Shuffle 所有数据,尽量少用
rdd.groupByKey()
常见问题
1. collect() 报 OOM
collect() 把全量数据拉回 Driver,数据量大就爆内存。
解决: 用 take(n) 采样,或者 saveAsTextFile 直接写 HDFS。
2. groupByKey 慢
groupByKey 把所有数据都 Shuffle 到 Reduce 端,数据量大时网络和内存压力大。
解决: 换 reduceByKey 或 aggregateByKey,Map 端预聚合。
3. 分区数不合理
分区太多 → Task 调度开销大;分区太少 → 并行度不够。
建议: 每个分区 128MB-256MB。repartition(n) 增加分区,coalesce(n) 减少分区。