SparkCore 之 Spark Action 类算子详解
·
摘要:系统拆解 Spark 全部 Action 算子——按输出类、存储类、聚合类、统计类四大分类,逐个剖析语法、底层原理、Driver/Executor 数据流向和性能陷阱。配有 2 张原创架构图、完整 Scala 代码示例和常见 OOM 排查指南。面向 Java、大数据及 AI 开发工程师。
一、Action vs Transformation — 本质区别
Action(行动算子)是 Spark 中触发计算的唯一入口。所有 Transformation 只构建 DAG,Action 才是真正按下"执行按钮"的那个操作。
| 维度 | Transformation | Action |
|---|---|---|
| 返回值 | 新 RDD | 非 RDD(值/数组/写入存储) |
| 触发计算 | ❌ 不触发 | ✅ 触发 DAG 执行 |
| DAGScheduler | 不参与 | 触发 Stage 划分 + Task 提交 |
| 执行时机 | 惰性求值 | 立即触发 |
| 代表 | map / filter / join | count / collect / saveAsTextFile |
二、输出类 Action — 将数据拉回 Driver
2.1 collect — 收集全部数据到 Driver
// 签名:def collect(): Array[T]
// ⚠️ 所有数据通过网络传输到 Driver 内存
// 数据量大时 → Driver OOM!
val rdd = sc.parallelize(1 to 1000)
rdd.collect() // Array[Int](1,2,...,1000) — 全部在Driver内存中
// ❌ 危险用法:PB 级数据 collect
sc.textFile("hdfs://100TB-logs/").collect() // Driver OOM 必现!
// ✅ 正确用法:数据量小(<几MB)时才用 collect
rdd.filter(_.contains("specific_keyword")).collect() // 过滤后数据量小
2.2 take — 取前 N 条
// 签名:def take(num: Int): Array[T]
// 只拉指定数量到 Driver,不会 OOM
rdd.take(10) // 前10条
rdd.take(100) // 前100条
// take 的底层执行:
// 先取第一个 Partition → 不够再取第二个 → 直到凑满 num 条
// 性能优于 collect(尤其数据量大时)
2.3 foreach / foreachPartition — 遍历执行(不返回Driver)
// foreachPartition: 每个分区执行一次(推荐!)
rdd.foreachPartition { iter =>
val conn = DriverManager.getConnection(url) // 每个分区创建1次连接
iter.foreach { record =>
conn.execute(s"INSERT INTO t VALUES ($record)")
}
conn.close()
}
// mapPartitions + foreachPartition 组合是写入外部系统的最优模式
三、存储类 Action — 将数据写入外部存储
3.1 saveAsTextFile — 写入文本文件
// 签名:def saveAsTextFile(path: String)
// 每个 Partition 写入一个文件(part-00000, part-00001, ...)
rdd.saveAsTextFile("hdfs://output/result/")
// ⚠️ 目标目录不能已存在(否则抛异常)
// 解决方案:先删除目录
val path = new Path("hdfs://output/result/")
path.getFileSystem(sc.hadoopConfiguration).delete(path, true)
rdd.saveAsTextFile("hdfs://output/result/")
四、聚合类 Action — 在 Executor 聚合后返回 Driver
4.1 reduce — 全局聚合
// 签名:def reduce(f: (T, T) => T): T
// 先在每个分区内聚合 → Shuffle → 最终聚合 → 返回Driver
val rdd = sc.parallelize(1 to 100)
rdd.reduce(_ + _) // 5050
rdd.reduce(_ max _) // 100
// ⚠️ reduce 要求函数满足结合律和交换律
4.2 aggregate — 灵活的分区内/区间聚合
// 求平均值
val rdd = sc.parallelize(1 to 100)
val (sum, count) = rdd.aggregate((0, 0))(
seqOp = { case ((s, c), v) => (s + v, c + 1) },
combOp = { case ((s1, c1), (s2, c2)) => (s1 + s2, c1 + c2) }
)
val avg = sum.toDouble / count // 50.5
4.3 treeAggregate / treeReduce — 树形聚合(推荐)
// 树形聚合避免单点 OOM
rdd.treeReduce(_ + _, depth = 3) // 推荐用于大数据集
rdd.treeAggregate(zero)(seqOp, combOp, depth = 3)
五、统计类 Action — 便捷统计函数
rdd.count() // 计数
rdd.countByKey() // 按Key统计 → Map[K, Long] (⚠️ OOM)
rdd.max() // 最大值
rdd.min() // 最小值
rdd.sum() // 求和
rdd.histogram(10) // 直方图
rdd.toDebugString // 查看Lineage
六、Action 算子速查表
| 算子 | 分类 | 返回 | 数据流向 | OOM 风险 |
|---|---|---|---|---|
collect() | 输出 | Array[T] | →Driver | ⚠️ 高 |
take(n) | 输出 | Array[T] | →Driver | ✅ 低 |
foreachPartition | 输出 | Unit | →Executor | ✅ 无 |
saveAsTextFile | 存储 | Unit | →磁盘 | ✅ 无 |
reduce(f) | 聚合 | T | →Driver | ✅ 低 |
aggregate | 聚合 | U | →Driver | ✅ 低 |
treeReduce | 聚合 | T | →Driver | ✅ 极低 |
count() | 统计 | Long | →Driver | ✅ 极低 |
countByKey() | 统计 | Map | →Driver | ⚠️ 高 |
七、常见 Action 陷阱与排查
| 陷阱 | 现象 | 解决 |
|---|---|---|
| collect OOM | OutOfMemoryError | take(n) 替代 |
| countByKey OOM | Driver GC overhead | reduceByKey(_+_).collect() |
| foreach 打印不显示 | 输出在Executor日志 | collect().foreach(println) |
| saveAsTextFile 目录已存在 | FileAlreadyExistsException | 先删除目录 |
写在最后
Action 算子虽然数量不多(约 20 个),但理解它们的数据流向和内存压力至关重要:
- collect/countByKey → Driver:数据向 Driver 汇聚 → 控制数据量
- foreach/saveAsTextFile → Executor/存储:数据向外发散 → 无内存压力
- reduce/aggregate → 聚合后返回:Executor 间聚合 → 关注 Shuffle
选对 Action,不仅决定性能,更决定你的 Driver 会不会 OOM。
更多推荐
所有评论(0)