Spark Action 算子详解:四大分类、数据流向与 DAG 触发执行原理

📅 2026/7/30 15:01:55 👁️ 阅读次数 📝 编程学习
Spark Action 算子详解:四大分类、数据流向与 DAG 触发执行原理
# SparkCore 之 Spark Action 类算子详解> **摘要**:系统拆解 Spark 全部 Action 算子——按输出类、存储类、聚合类、统计类四大分类,逐个剖析语法、底层原理、Driver/Executor 数据流向和性能陷阱。配有 2 张原创架构图、完整 Scala 代码示例和常见 OOM 排查指南。面向 Java、大数据及 AI 开发工程师。---## 一、Action vs Transformation — 本质区别Action(行动算子)是 Spark 中**触发计算的唯一入口**。所有 Transformation 只构建 DAG,Action 才是真正按下"执行按钮"的那个操作。![Action 算子全景分类](DIAGRAM:d1)| 维度 | 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 触发与数据流向![Action 触发 DAG 执行流程](DIAGRAM:d2)### 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 真正落地,让数据创造价值