Spark 转换算子(Transformation)全面解析
Spark 转换算子(Transformation)全面解析
摘要:本文系统梳理 Spark 中的转换算子,涵盖惰性求值机制、常用 API 分类、宽窄依赖原理与典型使用场景。读完本篇,再看行动算子,整个 Spark 编程模型就通了。
一、什么是转换算子?
Spark 的核心编程模型基于 RDD(弹性分布式数据集),对 RDD 的操作分两类:
| 类型 | 返回值 | 是否触发计算 | 特点 |
|---|---|---|---|
| 转换算子(Transformation) | 新的 RDD | ❌ 不触发 | 惰性求值,只记录血缘 |
| 行动算子(Action) | 具体值或写出 | ✅ 触发 | 驱动 DAG 真正执行 |
惰性求值(Lazy Evaluation) 是理解转换算子的关键:调用转换算子时,Spark 不会立即执行任何计算,而是将操作记录在 DAG(有向无环图)中,等到行动算子被调用时才统一执行。
sc.textFile(...) // 读取数据(也是惰性的)
.flatMap(_.split(" ")) // 记录到 DAG,不执行
.map((_, 1)) // 记录到 DAG,不执行
.filter(_._2 > 0) // 记录到 DAG,不执行
.collect() // ← Action!此刻 DAG 才真正执行
二、宽依赖 vs 窄依赖
理解转换算子,绕不开依赖关系,它直接决定 Stage 如何划分。
窄依赖(Narrow Dependency)
父 RDD 的每个分区最多被一个子 RDD 分区使用。无需 Shuffle,可在同一 Stage 内流水线执行。

典型:map、filter、union、flatMap
宽依赖(Wide Dependency)
父 RDD 的一个分区被多个子 RDD 分区使用,需要 Shuffle,是 Stage 的切分点。

典型:groupByKey、reduceByKey、join、repartition
性能核心:Shuffle 是 Spark 性能的主要瓶颈,涉及磁盘 I/O 和网络传输。设计程序时应尽量减少不必要的宽依赖。
三、转换算子分类总览
Transformation 算子
├── 基础转换类 map / flatMap / filter / mapPartitions / mapPartitionsWithIndex
├── 去重排序类 distinct / sortBy / sortByKey
├── 聚合类 groupBy / groupByKey / reduceByKey / foldByKey / aggregateByKey / combineByKey
├── 集合操作类 union / intersection / subtract / zip / cartesian
├── 重分区类 repartition / coalesce
├── Key-Value 类 mapValues / flatMapValues / keys / values / partitionBy
└── 其他 glom / sample / cogroup
四、常用转换算子详解
4.1 map(func)
对 RDD 中每个元素应用函数,一对一映射,返回新 RDD。
val rdd = sc.parallelize(List(1, 2, 3, 4, 5))
val result = rdd.map(_ * 2)
result.collect() // Array(2, 4, 6, 8, 10)
4.2 flatMap(func)
对 RDD 中的每个元素应用函数,函数返回一个集合,flatMap 会将所有集合"拆开"合并成一个新的 RDD。适合处理"一行变多行"的场景,比如把一句话拆成单词。
val rdd = sc.parallelize(List("hello world", "spark is fun"))
val words = rdd.flatMap(_.split(" "))
words.collect() // Array(hello, world, spark, is, fun)
map和flatMap的区别:map返回的是每个元素映射后的结果(可能是集合),flatMap会将结果中的集合展开一层。
rdd.map(_.split(" ")).collect()
// Array(Array(hello, world), Array(spark, is, fun)) ← 嵌套数组
rdd.flatMap(_.split(" ")).collect()
// Array(hello, world, spark, is, fun) ← 已展平
4.3 filter(func)
保留使函数返回 true 的元素。
val rdd = sc.parallelize(1 to 10)
val evens = rdd.filter(_ % 2 == 0)
evens.collect() // Array(2, 4, 6, 8, 10)
4.4 mapPartitions(func)
以分区为单位执行函数,函数参数为 Iterator[T],返回 Iterator[U]。
val rdd = sc.parallelize(1 to 6, 2) // 2 个分区
val result = rdd.mapPartitions { iter =>
iter.map(_ * 10)
}
result.collect() // Array(10, 20, 30, 40, 50, 60)
与 map 的区别:
mapPartitions每个分区只调用一次函数,适合需要初始化资源(如数据库连接)的场景,避免每条数据都重复创建资源。
4.5 distinct(numPartitions)
对 RDD 元素去重,底层触发 Shuffle。
val rdd = sc.parallelize(List(1, 2, 2, 3, 3, 3))
rdd.distinct().collect() // Array(1, 2, 3) 顺序不保证
4.6 sortBy(func, ascending, numPartitions)
按指定键对 RDD 排序。
val rdd = sc.parallelize(List(3, 1, 4, 1, 5, 9, 2, 6))
rdd.sortBy(x => x).collect() // 升序: Array(1, 1, 2, 3, 4, 5, 6, 9)
rdd.sortBy(x => x, false).collect() // 降序: Array(9, 6, 5, 4, 3, 2, 1, 1)
// 对元组按第二个字段排序
val pairRdd = sc.parallelize(List(("a", 3), ("b", 1), ("c", 2)))
pairRdd.sortBy(_._2).collect()
// Array((b,1), (c,2), (a,3))
4.7 groupByKey()
将 PairRDD (K, V) 中相同 Key 的所有 Value 聚合成一个迭代器。
val rdd = sc.parallelize(List(("a", 1), ("b", 2), ("a", 3), ("b", 4)))
rdd.groupByKey().collect()
// Array((a, [1,3]), (b, [2,4]))
⚠️ 性能警告:
groupByKey会将所有数据 Shuffle 到网络,然后在内存中聚合,极易造成数据倾斜和 OOM。如果目的是聚合(求和、计数等),优先使用reduceByKey或aggregateByKey。
4.8 reduceByKey(func)
在 Shuffle 之前先在每个分区内做局部聚合(Map 端预聚合),大幅减少网络传输量。
val rdd = sc.parallelize(List(("a", 1), ("b", 1), ("a", 1), ("b", 1)))
rdd.reduceByKey(_ + _).collect()
// Array((a,2), (b,2))
执行过程对比:
groupByKey:
Executor A: (a,1),(a,1) ──Shuffle──► (a, [1,1,1,1]) ──聚合──► (a,4)
Executor B: (a,1),(a,1) ─────────────────────────────────────────▲
reduceByKey:
Executor A: (a,1),(a,1) ──局部合并──► (a,2) ──Shuffle──► (a,4)
Executor B: (a,1),(a,1) ──局部合并──► (a,2) ──────────────▲
网络传输量显著减少!
4.9 aggregateByKey(zeroValue)(seqOp, combOp)
最灵活的按 Key 聚合算子,允许 Value 类型在聚合过程中发生变化。
zeroValue:每个 Key 的初始值(分区级别)seqOp:分区内聚合(U, V) => UcombOp:分区间合并(U, U) => U
// 计算每个 Key 对应的平均值
val rdd = sc.parallelize(List(("a", 1), ("a", 3), ("b", 2), ("b", 4)), 2)
val result = rdd.aggregateByKey((0, 0))(
(acc, v) => (acc._1 + v, acc._2 + 1), // seqOp: 累加值和计数
(a, b) => (a._1 + b._1, a._2 + b._2) // combOp: 合并两个分区的结果
)
result.mapValues(v => v._1.toDouble / v._2).collect()
// Array((a,2.0), (b,3.0))
4.10 join / leftOuterJoin / rightOuterJoin / fullOuterJoin
对两个 PairRDD 按 Key 进行关联,底层触发 Shuffle。
val rdd1 = sc.parallelize(List(("a", 1), ("b", 2), ("c", 3)))
val rdd2 = sc.parallelize(List(("a", "x"), ("b", "y"), ("d", "z")))
rdd1.join(rdd2).collect()
// Array((a,(1,x)), (b,(2,y))) — 内连接,只保留两边都有的 Key
rdd1.leftOuterJoin(rdd2).collect()
// Array((a,(1,Some(x))), (b,(2,Some(y))), (c,(3,None)))
rdd1.rightOuterJoin(rdd2).collect()
// Array((a,(Some(1),x)), (b,(Some(2),y)), (d,(None,z)))
4.11 union(otherRDD)
合并两个相同类型的 RDD,不去重,窄依赖,不触发 Shuffle。
val rdd1 = sc.parallelize(List(1, 2, 3))
val rdd2 = sc.parallelize(List(3, 4, 5))
rdd1.union(rdd2).collect() // Array(1, 2, 3, 3, 4, 5)
4.12 repartition(numPartitions) / coalesce(numPartitions)
| 算子 | 说明 | Shuffle |
|---|---|---|
repartition(n) |
重新分为 n 个分区(可增可减) | 有 |
coalesce(n) |
合并分区(默认仅减少) | 默认无 |
val rdd = sc.parallelize(1 to 100, 10) // 10 个分区
rdd.repartition(5).getNumPartitions // 5
rdd.coalesce(5).getNumPartitions // 5,但无 Shuffle,效率更高
// coalesce 增加分区数时需显式开启 shuffle
rdd.coalesce(20, shuffle = true).getNumPartitions // 20
最佳实践:需要减少分区时优先用
coalesce,避免不必要的 Shuffle;需要增加分区或数据已严重倾斜时用repartition。
五、聚合算子横向对比
开发中经常纠结该用哪个聚合算子,一表看清:
| 算子 | 初始值 | 输入/输出类型 | 预聚合 | 适用场景 |
|---|---|---|---|---|
groupByKey |
无 | V → Iterable[V] | ❌ | 需要拿到所有原始值(不推荐聚合用) |
reduceByKey |
无 | V + V → V | ✅ | 简单同类型聚合(求和、求积等) |
aggregateByKey |
有 | (U, V) → U | ✅ | 类型转换(如求平均、收集集合) |
选择建议:简单求和/计数 → reduceByKey;需要类型转换 → aggregateByKey;避免在聚合场景使用 groupByKey。
六、常见误区与最佳实践
❌ 误区 1:groupByKey 代替 reduceByKey
// 低效:所有数据先 Shuffle 再聚合
rdd.groupByKey().map { case (k, vs) => (k, vs.sum) }
// 高效:先局部聚合再 Shuffle
rdd.reduceByKey(_ + _)
❌ 误区 2:频繁 repartition 增加无效 Shuffle
// 每次 repartition 都是一次全量 Shuffle
rdd.repartition(100).map(...).repartition(50)
// 尽量在处理前一次性调整分区数
rdd.repartition(100).map(...)
✅ 最佳实践总结
| 场景 | 推荐算子 |
|---|---|
| 元素级别变换 | map / flatMap / filter |
| 资源初始化开销大 | mapPartitions |
| 按 Key 简单聚合 | reduceByKey |
| 按 Key 复杂聚合 | aggregateByKey |
| 合并多个 RDD | union(无 Shuffle) |
| 减少分区 | coalesce(无 Shuffle) |
| 重新均衡分区 | repartition |
七、总结
转换算子是 Spark 程序的"骨架",核心理念是:
- 惰性求值:转换算子只构建 DAG,不执行任何计算
- 宽窄依赖:窄依赖无需 Shuffle 可流水线执行,宽依赖是 Stage 切分点
- 预聚合意识:优先使用带 Map 端预聚合的算子(
reduceByKey>groupByKey) - 分区感知:合理控制分区数,避免数据倾斜和不必要的 Shuffle
下一篇我们将介绍行动算子(Action),看看 Action 如何触发整个 DAG 真正执行
如有问题欢迎在评论区交流,也欢迎关注后续 Spark 系列文章。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)