ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

SparkCore 之 Spark Action 类算子详解

SparkCore 之 Spark Action 类算子详解 摘要系统拆解 Spark 全部 Action 算子——按输出类、存储类、聚合类、统计类四大分类逐个剖析语法、底层原理、Driver/Executor 数据流向和性能陷阱。配有 2 张原创架构图、完整 Scala 代码示例和常见 OOM 排查指南。面向 Java、大数据及 AI 开发工程师。一、Action vs Transformation — 本质区别Action行动算子是 Spark 中触发计算的唯一入口。所有 Transformation 只构建 DAGAction 才是真正按下执行按钮的那个操作。维度TransformationAction返回值新 RDD非 RDD值/数组/写入存储触发计算❌ 不触发✅ 触发 DAG 执行DAGScheduler不参与触发 Stage 划分 Task 提交执行时机惰性求值立即触发代表map / filter / joincount / collect / saveAsTextFile二、输出类 Action — 将数据拉回 Driver2.1 collect — 收集全部数据到 Driver// 签名def collect(): Array[T]// ⚠️ 所有数据通过网络传输到 Driver 内存// 数据量大时 → Driver OOMvalrddsc.parallelize(1to1000)rdd.collect()// Array[Int](1,2,...,1000) — 全部在Driver内存中// ❌ 危险用法PB 级数据 collectsc.textFile(hdfs://100TB-logs/).collect()// Driver OOM 必现// ✅ 正确用法数据量小几MB时才用 collectrdd.filter(_.contains(specific_keyword)).collect()// 过滤后数据量小2.2 take — 取前 N 条// 签名def take(num: Int): Array[T]// 只拉指定数量到 Driver不会 OOMrdd.take(10)// 前10条rdd.take(100)// 前100条// take 的底层执行// 先取第一个 Partition → 不够再取第二个 → 直到凑满 num 条// 性能优于 collect尤其数据量大时2.3 foreach / foreachPartition — 遍历执行不返回Driver// foreachPartition: 每个分区执行一次推荐rdd.foreachPartition{itervalconnDriverManager.getConnection(url)// 每个分区创建1次连接iter.foreach{recordconn.execute(sINSERT 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/)// ⚠️ 目标目录不能已存在否则抛异常// 解决方案先删除目录valpathnewPath(hdfs://output/result/)path.getFileSystem(sc.hadoopConfiguration).delete(path,true)rdd.saveAsTextFile(hdfs://output/result/)四、聚合类 Action — 在 Executor 聚合后返回 Driver4.1 reduce — 全局聚合// 签名def reduce(f: (T, T) T): T// 先在每个分区内聚合 → Shuffle → 最终聚合 → 返回Drivervalrddsc.parallelize(1to100)rdd.reduce(__)// 5050rdd.reduce(_ max _)// 100// ⚠️ reduce 要求函数满足结合律和交换律4.2 aggregate — 灵活的分区内/区间聚合// 求平均值valrddsc.parallelize(1to100)val(sum,count)rdd.aggregate((0,0))(seqOp{case((s,c),v)(sv,c1)},combOp{case((s1,c1),(s2,c2))(s1s2,c1c2)})valavgsum.toDouble/count// 50.54.3 treeAggregate / treeReduce — 树形聚合推荐// 树形聚合避免单点 OOMrdd.treeReduce(__,depth3)// 推荐用于大数据集rdd.treeAggregate(zero)(seqOp,combOp,depth3)五、统计类 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。
返回列表