摘要:系统拆解 Spark 全部 Action 算子——按输出类、存储类、聚合类、统计类四大分类,逐个剖析语法、底层原理、Driver/Executor 数据流向和性能陷阱。配有 2 张原创架构图、完整 Scala 代码示例和常见 OOM 排查指南。面向 Java、大数据及 AI 开发工程师。


一、Action vs Transformation — 本质区别

Action(行动算子)是 Spark 中触发计算的唯一入口。所有 Transformation 只构建 DAG,Action 才是真正按下"执行按钮"的那个操作。

维度TransformationAction
返回值新 RDD非 RDD(值/数组/写入存储)
触发计算❌ 不触发✅ 触发 DAG 执行
DAGScheduler不参与触发 Stage 划分 + Task 提交
执行时机惰性求值立即触发
代表map / filter / joincount / 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 OOMOutOfMemoryErrortake(n) 替代
countByKey OOMDriver GC overheadreduceByKey(_+_).collect()
foreach 打印不显示输出在Executor日志collect().foreach(println)
saveAsTextFile 目录已存在FileAlreadyExistsException先删除目录

写在最后

Action 算子虽然数量不多(约 20 个),但理解它们的数据流向内存压力至关重要:

  • collect/countByKey → Driver:数据向 Driver 汇聚 → 控制数据量
  • foreach/saveAsTextFile → Executor/存储:数据向外发散 → 无内存压力
  • reduce/aggregate → 聚合后返回:Executor 间聚合 → 关注 Shuffle

选对 Action,不仅决定性能,更决定你的 Driver 会不会 OOM。

更多推荐