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))   // 推荐,简洁

parallelizemakeRDD 功能一样。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 端,数据量大时网络和内存压力大。

解决:reduceByKeyaggregateByKey,Map 端预聚合。

3. 分区数不合理

分区太多 → Task 调度开销大;分区太少 → 并行度不够。

建议: 每个分区 128MB-256MB。repartition(n) 增加分区,coalesce(n) 减少分区。