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)

mapflatMap 的区别: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。如果目的是聚合(求和、计数等),优先使用 reduceByKeyaggregateByKey


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) => U
  • combOp:分区间合并 (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 程序的"骨架",核心理念是:

  1. 惰性求值:转换算子只构建 DAG,不执行任何计算
  2. 宽窄依赖:窄依赖无需 Shuffle 可流水线执行,宽依赖是 Stage 切分点
  3. 预聚合意识:优先使用带 Map 端预聚合的算子(reduceByKey > groupByKey
  4. 分区感知:合理控制分区数,避免数据倾斜和不必要的 Shuffle

下一篇我们将介绍行动算子(Action),看看 Action 如何触发整个 DAG 真正执行


如有问题欢迎在评论区交流,也欢迎关注后续 Spark 系列文章。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐