# 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```scala
// 签名: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 条```scala
// 签名:def take(num: Int): Array[T]
// 只拉指定数量到 Driver,不会 OOMrdd.take(10) // 前10条
rdd.take(100) // 前100条// take 的底层执行:
// 先取第一个 Partition → 不够再取第二个 → 直到凑满 num 条
// 性能优于 collect(尤其数据量大时)
```### 2.3 first — 取第一条```scala
// first() = take(1).head
val first = rdd.first() // 快速查看数据样例
```### 2.4 top / takeOrdered — 取最大/最小 N 条```scala
// top: 降序取前N条(大的在前)
rdd.top(5) // 默认排序
rdd.top(5)(Ordering.by(_.length)) // 按长度排序// takeOrdered: 升序取前N条(小的在前)
rdd.takeOrdered(5)// 底层:每个Partition取TopN → Shuffle到单个Partition → 最终TopN
// 数据量大时注意单个Partition的OOM风险
```### 2.5 foreach / foreachPartition — 遍历执行(不返回Driver)```scala
// foreach: 每个元素执行副作用操作(写入数据库/打印等)
rdd.foreach(println) // ⚠️ 打印在Executor日志,不在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 — 写入文本文件```scala
// 签名:def saveAsTextFile(path: String)
// 每个 Partition 写入一个文件(part-00000, part-00001, ...)rdd.saveAsTextFile("hdfs://output/result/")
// 生成文件:part-00000, part-00001, ..., part-00xxx// ⚠️ 目标目录不能已存在(否则抛异常)
// 解决方案:先删除目录
val path = new Path("hdfs://output/result/")
path.getFileSystem(sc.hadoopConfiguration).delete(path, true)
rdd.saveAsTextFile("hdfs://output/result/")
```### 3.2 saveAsObjectFile / saveAsSequenceFile```scala
// saveAsObjectFile: Java 序列化写入(不推荐,兼容性差)
rdd.saveAsObjectFile("hdfs://output/obj/")// saveAsSequenceFile: Hadoop SequenceFile 格式(仅 K-V RDD 可用)
val kvRDD = sc.parallelize(Seq(("a", 1), ("b", 2)))
kvRDD.saveAsSequenceFile("hdfs://output/seq/")
```---## 四、聚合类 Action — 在 Executor 聚合后返回 Driver### 4.1 reduce — 全局聚合```scala
// 签名:def reduce(f: (T, T) => T): T
// 先在每个分区内聚合 → Shuffle 到一个分区 → 最终聚合 → 返回Driverval rdd = sc.parallelize(1 to 100)
rdd.reduce(_ + _) // 5050
rdd.reduce(_ max _) // 100// ⚠️ reduce 要求函数满足结合律和交换律(因分区聚合顺序不确定)
// ✅ sum/max/min/count 满足 → 安全
// ❌ (a-b) — 不满足结合律 → 结果不确定
```### 4.2 fold — 带初始值的 reduce```scala
// 签名:def fold(zeroValue: T)(op: (T, T) => T): T
// 每个分区先用 zeroValue 初始化,再聚合rdd.fold(0)(_ + _) // 5050
rdd.fold(1)(_ * _) // ⚠️ 结果不确定!每个分区都乘以1,最终多乘了1// fold 的 zeroValue 在每个分区都参与一次,最终结果多一次
// → 使用 aggregate 替代 fold
```### 4.3 aggregate — 灵活的分区内/区间聚合```scala
// 签名:def aggregate[U](zeroValue: U)(seqOp: (U,T)=>U, combOp: (U,U)=>U): U
// seqOp: 分区内聚合 · combOp: 分区间聚合// 求平均值
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.4 treeAggregate / treeReduce — 树形聚合(推荐)```scala
// 普通 reduce:所有分区→Shuffle→1个分区→OOM风险
// treeReduce:多轮局部聚合→Shuffle→最终聚合→避免OOMrdd.treeReduce(_ + _, depth = 3) // 推荐用于大数据集聚合
rdd.treeAggregate(zero)(seqOp, combOp, depth = 3)// depth 越大 = 多轮聚合 = 更省内存但多轮Shuffle
// 建议:depth=2或3 即可,过大会增加延迟
```---## 五、统计类 Action — 便捷统计函数### 5.1 count / countByKey / countByValue```scala
rdd.count() // 计数(触发DAG)
rdd.countByKey() // 按Key统计 → Map[K, Long]
rdd.countByValue() // 按值统计 → Map[T, Long]// countByKey 的问题:所有数据传到 Driver → 大数据量OOM
// 替代方案:rdd.mapValues(_ => 1L).reduceByKey(_ + _).collect()
```### 5.2 collectAsMap / lookup```scala
// collectAsMap: 收集为 Map[K, V](仅 K-V RDD,重复 Key 只保留最后一个)
val rdd = sc.parallelize(Seq(("a",1),("b",2),("a",3)))
rdd.collectAsMap() // Map(a -> 3, b -> 2) — "a"取了最后一个// lookup: 通过 Key 查找所有 Value
rdd.lookup("a") // Seq(1, 3) — 可能触发全表扫描
```### 5.3 isEmpty / max / min / sum```scala
rdd.isEmpty() // 是否为空
rdd.max() // 最大值(需要 Ordering)
rdd.min() // 最小值
rdd.sum() // 求和(需要 Numeric)
rdd.mean() // 平均值
rdd.variance() // 方差
rdd.stdev() // 标准差
rdd.histogram(10) // 直方图(10个桶)
```---## 六、Action 执行原理 — DAG 触发与数据流向### 6.1 Action 触发执行的全流程```
① 调用 Action 算子(如 count())↓
② SparkContext.runJob() 被调用↓
③ DAGScheduler.handleJobSubmitted()├── 反向回溯 Lineage 构建完整 DAG├── 遇到 ShuffleDependency → 切 Stage└── 生成 TaskSet 提交给 TaskScheduler↓
④ TaskScheduler 按数据本地性分配 Task 到 Executor↓
⑤ Executor 执行 Task,结果返回 Driver(或写入存储)↓
⑥ Action 返回结果给用户代码
```### 6.2 数据流向差异| Action 类型 | 数据流向 | Driver 内存压力 | 适用场景 |
|------------|----------|---------------|----------|
| `collect()` | Executor → **Driver** | ⚠️ **高**(全部数据) | 小数据集 |
| `take(n)` | Executor → **Driver** | ✅ 低(仅n条) | 预览数据 |
| `foreach` | Driver → **Executor** | ✅ 无 | 写入外部存储 |
| `saveAsTextFile` | Executor → **磁盘** | ✅ 无 | 持久化输出 |
| `reduce` | Executor → Executor → **Driver** | ✅ 低(聚合后) | 全局聚合 |
| `countByKey` | Executor → **Driver** | ⚠️ 高(所有Key) | Key种类少时 |---## 七、Action 算子速查表| 算子 | 分类 | 返回类型 | 数据流向 | OOM 风险 |
|------|------|----------|----------|----------|
| `collect()` | 输出 | Array[T] | →Driver | ⚠️ 高 |
| `take(n)` | 输出 | Array[T] | →Driver | ✅ 低 |
| `first()` | 输出 | T | →Driver | ✅ 低 |
| `top(n)` | 输出 | Array[T] | →Driver | ✅ 低 |
| `takeOrdered(n)` | 输出 | Array[T] | →Driver | ✅ 低 |
| `foreach(f)` | 输出 | Unit | →Executor | ✅ 无 |
| `foreachPartition(f)` | 输出 | Unit | →Executor | ✅ 无 |
| `saveAsTextFile` | 存储 | Unit | →磁盘 | ✅ 无 |
| `saveAsObjectFile` | 存储 | Unit | →磁盘 | ✅ 无 |
| `saveAsSequenceFile` | 存储 | Unit | →磁盘 | ✅ 无 |
| `reduce(f)` | 聚合 | T | →Driver | ✅ 低 |
| `fold(z)(op)` | 聚合 | T | →Driver | ✅ 低 |
| `aggregate(z)(s,c)` | 聚合 | U | →Driver | ✅ 低 |
| `treeReduce(f,d)` | 聚合 | T | →Driver | ✅ 极低 |
| `treeAggregate(z)(s,c,d)` | 聚合 | U | →Driver | ✅ 极低 |
| `count()` | 统计 | Long | →Driver | ✅ 极低 |
| `countByKey()` | 统计 | Map[K,Long] | →Driver | ⚠️ 高 |
| `countByValue()` | 统计 | Map[T,Long] | →Driver | ⚠️ 高 |
| `collectAsMap()` | 统计 | Map[K,V] | →Driver | ⚠️ 高 |
| `lookup(key)` | 统计 | Seq[V] | →Driver | ✅ 低 |
| `sum/max/min/mean` | 统计 | 数值 | →Driver | ✅ 极低 |
| `histogram(b)` | 统计 | (Array,Array) | →Driver | ✅ 低 |
| `toDebugString` | 调试 | String | →Driver | ✅ 极低 |---## 八、常见 Action 陷阱与排查| 陷阱 | 现象 | 根因 | 解决 |
|------|------|------|------|
| **collect OOM** | `java.lang.OutOfMemoryError: Java heap space` | 全量数据拉回Driver | `take(n)` 替代;增加 `--driver-memory` |
| **countByKey OOM** | Driver GC overhead | Key种类过多 | `mapValues(_=>1L).reduceByKey(_+_).collect()` |
| **foreach 打印不显示** | 调用 `foreach(println)` 控制台无输出 | 打印在Executor日志中 | 用 `collect().foreach(println)` 或 `take(10).foreach(println)` |
| **saveAsTextFile 目录已存在** | `FileAlreadyExistsException` | 同名目录已存在 | 先删除目标目录 |
| **reduce 结果不确定** | 每次运行结果不同 | 算子不满足结合律 | 确认函数满足结合律和交换律 |---## 写在最后Action 算子虽然数量不多(约 20 个),但理解它们的**数据流向**和**内存压力**至关重要:- **collect/countByKey → Driver**:数据向 Driver 汇聚 → 控制数据量
- **foreach/saveAsTextFile → Executor/存储**:数据向外发散 → 无内存压力
- **reduce/aggregate → 聚合后返回**:Executor 间聚合 → 关注 Shuffle**选对 Action,不仅决定性能,更决定你的 Driver 会不会 OOM。**---> **starzy** | AI Data Engineer / 大数据技术实践者
> 专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
> 技术博客:blog.starzy.cn | GitHub:starzy1990.github.io
> 让 AI 真正落地,让数据创造价值
相关新闻
柳州2026家里房子漏水怎么办?市面上多种方案可选择,哪种最适合自己?专业防水公司免费上门为您评估,家里漏水不再愁! - 吉林同城获客
2026/7/30 15:01:55
查看详情
2026钦州24小时宠物医院优质机构实用指南 - 谁都没有我好看
2026/7/30 15:00:22
查看详情
技能md文件对 LLM提问怎么解析文本并执行名命令:“生成命令- bash 块 “和“(function calling) 返回 JSON“ 两种方式
2026/7/30 15:00:22
查看详情
在Obsidian中一键导出PDF、Word和ePub:终极Pandoc插件完整指南
2026/7/30 16:14:08
查看详情
长期更新软件:注册无广告WinRaR.v7.23压缩解压工具_二合一美化版
2026/7/30 16:14:08
查看详情
bit-bang时序被编译器优化破坏了怎么办
2026/7/30 16:12:44
查看详情
LangChain 1.x升级指南:init_chat_model新特性解析
2026/7/30 16:12:44
查看详情
基因载体全解析:从染色体结构到人工载体的生物技术应用
2026/7/30 16:12:44
查看详情
CTF应急响应实战:从流量分析到内存取证的综合安全能力训练
2026/7/30 16:12:44
查看详情
nfsserve:用Rust构建跨平台NFSv3服务器的终极指南
2026/7/30 0:00:28
查看详情
如何5分钟搭建TeamSpeak3音乐机器人:TS3AudioBot完全指南
2026/7/30 0:01:54
查看详情
从wireshark抓包开始,学习车载以太网(连载六:报文解析4)
2026/7/30 0:01:54
查看详情
SmartSentinel调试日志:实时追踪JSON解析错误的终极工具
2026/7/29 9:52:24
查看详情
数字身份克隆技术:Second Me开源项目解析与应用
2026/7/29 3:46:43
查看详情
remix-i18next TypeScript类型安全实践:确保翻译键与类型定义同步
2026/7/30 1:10:13
查看详情
企业 GEO 优化完整应用场景
2026/7/29 9:52:24
查看详情
ShaderGlass:如何在Windows桌面上实时运行GPU着色器的完整指南
2026/7/29 9:52:24
查看详情
ai agent框架spring ai/alibaba 源码原理分析(六) agent和组件
2026/7/29 5:31:19
查看详情