Spark RDD算子深度解析:Transformation与Action全掌握

Spark RDD算子深度解析:Transformation与Action全掌握

    • 一、Transformation与Action的核心区别
      • 1.1 执行机制对比
      • 1.2 核心区别
    • 二、Transformation算子详解
      • 2.1 常用Transformation算子一览
      • 2.2 map:一对一转换
      • 2.3 filter:过滤数据
      • 2.4 flatMap:一对多转换
      • 2.5 groupByKey:按Key分组
      • 2.6 reduceByKey:按Key聚合
      • 2.7 repartition:重新分区
    • 三、Action算子详解
      • 3.1 常用Action算子一览
      • 3.2 reduce:聚合操作
      • 3.3 collect:收集到Driver
      • 3.4 count:计数
      • 3.5 take:取前N个元素
      • 3.6 saveAsTextFile:保存到文件系统
    • 四、算子选择与性能优化
      • 4.1 窄依赖 vs 宽依赖
      • 4.2 性能优化建议
      • 4.3 算子选择指南
    • 五、面试高频问题
      • Q1:Transformation和Action有什么区别?
      • Q2:reduceByKey和groupByKey有什么区别?
      • Q3:map和flatMap有什么区别?
      • Q4:repartition和coalesce有什么区别?
      • Q5:collect有什么风险?什么时候用?
    • 六、总结
      • 6.1 算子速记表
      • 6.2 核心原则

🌺The Begin🌺点点关注,收藏不迷路🌺

关键词:Spark RDD、Transformation算子、Action算子、懒执行、窄依赖、宽依赖、算子选择

在Spark中,RDD(弹性分布式数据集)是所有计算的基石。而RDD算子则是对数据进行操作的核心工具,分为**Transformation(转换算子)Action(行动算子)**两大类。理解这两类算子的区别及各自的常用算子,是掌握Spark编程的第一步。

今天,我们将深入剖析常见的RDD算子,包括它们的作用、使用示例、执行原理以及性能特点。


一、Transformation与Action的核心区别

1.1 执行机制对比

Action算子

触发计算

触发计算

触发计算

reduce

提交Job

collect

count

返回结果/写入存储

Transformation算子

懒执行

懒执行

懒执行

记录血缘关系

map

构建DAG

filter

flatMap

Lineage

1.2 核心区别

特性 Transformation Action
返回值 新的RDD 非RDD(值、数组、或Unit)
执行方式 懒执行,只记录血缘 立即执行,触发作业
是否构建DAG 否(但会触发DAG执行)
是否产生结果 不产生最终结果 产生最终结果或副作用
示例 map, filter, flatMap reduce, collect, count

二、Transformation算子详解

2.1 常用Transformation算子一览

算子 作用 示例
map 对每个元素应用函数 rdd.map(x => x * 2)
filter 过滤元素 rdd.filter(x => x > 5)
flatMap 一对多映射 rdd.flatMap(line => line.split(" "))
groupByKey 按Key分组 pairRdd.groupByKey()
reduceByKey 按Key聚合 pairRdd.reduceByKey(_ + _)
sortByKey 按Key排序 pairRdd.sortByKey()
distinct 去重 rdd.distinct()
union 并集 rdd1.union(rdd2)
intersection 交集 rdd1.intersection(rdd2)
join 内连接 rdd1.join(rdd2)
repartition 重新分区 rdd.repartition(10)
coalesce 减少分区 rdd.coalesce(5)

2.2 map:一对一转换

// map:对每个元素应用函数,输入输出一对一
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5))
// 每个元素乘以2
val mapRdd = rdd.map(x => x * 2)
// 结果:2, 4, 6, 8, 10
// 类型可以改变
val stringRdd = rdd.map(x => s"Number: $x")
// 结果:"Number: 1", "Number: 2", ...

原理:每个分区内的数据逐个处理,窄依赖,不需要Shuffle。

2.3 filter:过滤数据

// filter:保留满足条件的元素
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5, 6))
// 保留大于3的元素
val filteredRdd = rdd.filter(x => x > 3)
// 结果:4, 5, 6
// 保留偶数
val evenRdd = rdd.filter(_ % 2 == 0)
// 结果:2, 4, 6

原理:也是窄依赖,每个分区独立过滤。

2.4 flatMap:一对多转换

// flatMap:每个输入产生0到多个输出
val rdd = sc.parallelize(Seq(
  "hello world",
  "spark is awesome",
  "flatMap example"
))
// 按空格切分成单词
val wordsRdd = rdd.flatMap(line => line.split(" "))
// 结果:"hello", "world", "spark", "is", "awesome", "flatMap", "example"
// 每个数字扩展成自身和自身+1
val nums = sc.parallelize(Seq(1, 2, 3))
val expanded = nums.flatMap(x => Seq(x, x + 1))
// 结果:1, 2, 2, 3, 3, 4

原理:也是窄依赖,但输出记录数可能增加。

2.5 groupByKey:按Key分组

// groupByKey:将相同Key的值分组到Iterable中
val pairRdd = sc.parallelize(Seq(
  ("a", 1), ("b", 2), ("a", 3), ("c", 4), ("b", 5)
))
val groupedRdd = pairRdd.groupByKey()
// 结果:(a, [1,3]), (b, [2,5]), (c, [4])
// 收集结果查看
groupedRdd.collect().foreach(println)
// (a,CompactBuffer(1, 3))
// (b,CompactBuffer(2, 5))
// (c,CompactBuffer(4))

原理宽依赖,需要Shuffle,相同Key的数据汇聚到同一分区。

⚠️ 注意:groupByKey在数据量大时性能较差,通常优先使用reduceByKey。

2.6 reduceByKey:按Key聚合

// reduceByKey:按Key进行聚合(比groupByKey更高效)
val pairRdd = sc.parallelize(Seq(
  ("a", 1), ("b", 2), ("a", 3), ("c", 4), ("b", 5)
))
// 按Key求和
val sumRdd = pairRdd.reduceByKey(_ + _)
// 结果:(a,4), (b,7), (c,4)
// 按Key求最大值
val maxRdd = pairRdd.reduceByKey((x, y) => if (x > y) x else y)
// 结果:(a,3), (b,5), (c,4)

原理宽依赖,但会在Map端预聚合(combine),减少Shuffle数据量。

groupByKey vs reduceByKey

reduceByKey

预聚合

Map端

Shuffle数据少

Reduce端最终聚合

groupByKey

全部数据

Map端

Shuffle

Reduce端聚合

2.7 repartition:重新分区

// repartition:增加或减少分区数(会进行Shuffle)
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5), 2)  // 初始2个分区
println(rdd.getNumPartitions)  // 2
// 增加到5个分区
val repartitionedRdd = rdd.repartition(5)
println(repartitionedRdd.getNumPartitions)  // 5
// 减少到1个分区(但repartition会Shuffle,用coalesce更好)
val coalescedRdd = rdd.coalesce(1)  // 不Shuffle,只合并分区

repartition vs coalesce

算子 作用 是否Shuffle 适用场景
repartition 增/减分区 需要均匀分布数据
coalesce 减分区 否(默认) 减少分区,避免Shuffle

三、Action算子详解

3.1 常用Action算子一览

算子 作用 返回值
reduce 聚合所有元素 单个值
collect 收集所有元素到Driver Array
count 统计元素个数 Long
take 取前N个元素 Array
first 取第一个元素 单个值
foreach 遍历每个元素 Unit
saveAsTextFile 保存为文本文件 Unit
countByKey 统计每个Key的出现次数 Map

3.2 reduce:聚合操作

// reduce:使用函数聚合所有元素
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5))
// 求和
val sum = rdd.reduce(_ + _)
println(sum)  // 15
// 求最大值
val max = rdd.reduce((x, y) => if (x > y) x else y)
println(max)  // 5
// 字符串连接
val strRdd = sc.parallelize(Seq("a", "b", "c", "d"))
val combined = strRdd.reduce(_ + _)
println(combined)  // "abcd"

原理:先在各分区内聚合,再聚合各分区结果。

3.3 collect:收集到Driver

// collect:将所有数据收集到Driver端(小心OOM!)
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5))
val data = rdd.collect()
println(data.mkString(", "))  // 1, 2, 3, 4, 5
// collectAsMap:对PairRDD收集成Map
val pairRdd = sc.parallelize(Seq(("a", 1), ("b", 2), ("a", 3)))
val map = pairRdd.collectAsMap()
// 结果:Map("a" -> 3, "b" -> 2)  // 注意相同Key会覆盖

⚠️ 警告:collect会把所有数据拉取到Driver,数据量大时会导致Driver OOM!

3.4 count:计数

// count:统计元素个数
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5))
val cnt = rdd.count()
println(cnt)  // 5
// countApprox:近似计数(大数据量时)
val approxCnt = rdd.countApprox(1000)  // 超时1000ms
println(approxCnt)
// countByValue:统计每个值出现的次数
val items = sc.parallelize(Seq("a", "b", "a", "c", "b", "a"))
val valueCounts = items.countByValue()
// 结果:Map("a" -> 3, "b" -> 2, "c" -> 1)

3.5 take:取前N个元素

// take:取前N个元素(按分区顺序,不一定是全局有序)
val rdd = sc.parallelize(Seq(5, 3, 1, 4, 2))
val first3 = rdd.take(3)
println(first3.mkString(", "))  // 可能是 5, 3, 1(取决于分区)
// takeOrdered:取排序后的前N个
val sortedFirst3 = rdd.takeOrdered(3)
println(sortedFirst3.mkString(", "))  // 1, 2, 3
// top:取最大的N个
val top3 = rdd.top(3)
println(top3.mkString(", "))  // 5, 4, 3

3.6 saveAsTextFile:保存到文件系统

// saveAsTextFile:保存RDD为文本文件
val rdd = sc.parallelize(Seq(
  "hello world",
  "spark is awesome",
  "save to file"
))
// 保存到HDFS或本地文件系统
rdd.saveAsTextFile("/output/spark/rdd-output")
// 会生成 part-00000, part-00001 等文件
// 保存为SequenceFile(只适用于PairRDD)
val pairRdd = sc.parallelize(Seq(("a", 1), ("b", 2), ("c", 3)))
pairRdd.saveAsSequenceFile("/output/spark/seq-output")

四、算子选择与性能优化

4.1 窄依赖 vs 宽依赖

宽依赖

Shuffle

Shuffle

Shuffle

Shuffle

父RDD分区1

子RDD分区1

子RDD分区2

父RDD分区2

窄依赖

一对一

一对一

一对一

父RDD分区1

子RDD分区1

父RDD分区2

子RDD分区2

父RDD分区3

子RDD分区3

依赖类型 特点 算子示例 性能
窄依赖 无Shuffle,流水线执行 map, filter, flatMap
宽依赖 需要Shuffle,数据落盘 groupByKey, reduceByKey, join

4.2 性能优化建议

// 1. 优先使用reduceByKey而不是groupByKey
// 差:groupByKey + map
rdd.groupByKey().mapValues(_.sum)
// 好:reduceByKey
rdd.reduceByKey(_ + _)
// 2. 尽早filter,减少数据量
// 差:先map后filter
rdd.map(expensiveFunc).filter(isGood)
// 好:先filter后map
rdd.filter(isGood).map(expensiveFunc)
// 3. 使用mapPartitions替代map(批量操作)
// 当有批量初始化开销时
rdd.mapPartitions { iter =>
  val conn = createConnection()  // 每个分区创建一次
  iter.map(processWithConn(conn))
}
// 4. 合理设置分区数
rdd.repartition(200)  // 根据集群资源调整

4.3 算子选择指南

需求 推荐算子 原因
一对一转换 map 简单高效
一对多转换 flatMap 灵活
过滤数据 filter 尽早过滤
按Key聚合 reduceByKey Map端预聚合
分组但不聚合 groupByKey 必要时用
取前N个 take/takeOrdered 避免collect
统计个数 count 高效
保存结果 saveAsTextFile 分布式保存

五、面试高频问题

Q1:Transformation和Action有什么区别?

  • Transformation:懒执行,返回新RDD,构建DAG血缘关系
  • Action:触发作业执行,返回结果或产生副作用
  • 关键:没有Action,Transformation不会真正执行

Q2:reduceByKey和groupByKey有什么区别?

  • reduceByKey:先在Map端预聚合(combine),减少Shuffle数据量
  • groupByKey:不预聚合,所有数据Shuffle后再处理
  • 性能:reduceByKey通常比groupByKey更高效

Q3:map和flatMap有什么区别?

  • map:一对一,每个输入产生一个输出
  • flatMap:一对多,每个输入产生0到多个输出
  • 使用场景:flatMap常用于分词、展开等操作

Q4:repartition和coalesce有什么区别?

  • repartition:可以增/减分区,一定会Shuffle
  • coalesce:只能减分区,默认不Shuffle(只合并分区)
  • 建议:减少分区用coalesce,增加分区用repartition

Q5:collect有什么风险?什么时候用?

  • 风险:将全量数据拉取到Driver,大数据量会导致OOM
  • 适用场景:结果集较小,或调试开发时
  • 替代方案:用take抽样、saveAsTextFile保存

六、总结

6.1 算子速记表

分类 算子 作用 依赖类型
Transformation map 一对一转换
filter 过滤
flatMap 一对多
groupByKey 分组
reduceByKey 聚合
repartition 重分区
Action reduce 聚合
collect 收集
count 计数
take 取前N
saveAsTextFile 保存

6.2 核心原则

Transformation构建蓝图,Action让蓝图变为现实

掌握这些常用算子,你就能在Spark开发中得心应手,写出高效、优雅的分布式计算代码!


思考题:mapPartitions和map有什么区别?在什么场景下应该用mapPartitions?欢迎在评论区讨论!

在这里插入图片描述

🌺The End🌺点点关注,收藏不迷路🌺
© 版权声明

相关文章