You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark 2.2.0分组聚合多Sum耗时翻倍的原因及优化咨询

我之前在处理Spark聚合任务时碰到过完全一样的情况,当时看到多一个sum耗时直接翻倍也懵了,后来翻了Spark的执行日志和底层逻辑才搞清楚来龙去脉,给你详细说说:

为什么双Sum的耗时几乎翻倍?

虽然你提到DataFrame是行存储,但真正影响耗时的是Spark执行聚合时的底层流程:

  1. Shuffle数据量的隐性增长:单sum时,每个Mapper输出的是(c1, sum_n1)的键值对;加了sum("n2")后,输出变成(c1, (sum_n1, sum_n2))。看起来每条记录只多了一个字段,但如果你的数据量级很大,整体Shuffle的数据量会接近翻倍——毕竟每个分组要携带两个聚合中间值。而Shuffle本身就是Spark里最吃IO和网络资源的环节,这部分开销直接拉满了耗时。
  2. Spark 2.2.0的聚合逻辑局限:这个版本的Spark在多聚合函数的处理上还比较“朴素”,相当于把两个sum的聚合逻辑拆成了近似独立的流程,只是最后合并了分组键。不像后续3.x版本有AQE(自适应查询执行)这类优化,能把多聚合的操作合并得更高效。
  3. 计算与IO的叠加效应:单sum时每个分区只需要对n1做累加,双sum要同时处理两个字段,虽然单条记录的计算量翻倍,但这部分其实占比不大,主要还是Shuffle的IO开销放大了整体耗时。
优化思路:让双Sum耗时接近单Sum

根据我当时的实践,这几个方法能有效缩小耗时差距:

  • 先过滤再聚合:如果n1、n2存在大量空值或者不需要计算的记录,先通过filter把这些数据筛掉,减少后续处理的数据量。比如:
    df.filter(col("n1").isNotNull && col("n2").isNotNull)
      .groupBy("c1")
      .agg(sum("n1"), sum("n2"))
    
  • 优化Shuffle参数:
    • 调整spark.sql.shuffle.partitions:默认是200,如果你的数据量很大,适当增大这个值(比如调到500-1000),让Shuffle任务更并行化,但别设得太大,否则会产生过多小任务的调度开销。
    • 开启Shuffle压缩:把spark.sql.shuffle.spill.compress设为true,减少磁盘IO的开销,这个参数对大Shuffle场景提升很明显。
  • 用RDD的AggregateByKey自定义聚合:DataFrame的agg在多字段聚合时会生成较多中间数据,换成RDD的aggregateByKey可以一次性对两个字段做累加,减少Shuffle的数据量。示例代码:
    // 转成RDD,提取分组键和需要聚合的两个字段
    val rdd = df.rdd.map(row => (row.getAs[String]("c1"), (row.getAs[Long]("n1"), row.getAs[Long]("n2"))))
    // 自定义聚合,一次性累加n1和n2
    val aggRdd = rdd.aggregateByKey((0L, 0L))(
      (acc, value) => (acc._1 + value._1, acc._2 + value._2), // 分区内累加
      (acc1, acc2) => (acc1._1 + acc2._1, acc1._2 + acc2._2)  // 分区间合并
    )
    // 转回DataFrame
    val resultDf = aggRdd.toDF("c1", "sum_pair")
      .select("c1", "sum_pair._1 as sum_n1", "sum_pair._2 as sum_n2")
    
    这种方式能在Mapper端就完成两个字段的同时累加,Shuffle时每个分组只传输一组中间结果,而不是两组,能大幅降低Shuffle的开销。
  • 升级Spark版本:如果条件允许,直接升到Spark 3.x以上版本。新版本的AQE能动态调整Shuffle分区数,还有更高效的聚合算子实现,多聚合函数的耗时会接近单聚合的水平。
  • 切换列存储格式:如果你的数据是存在磁盘上的行存储格式(比如CSV、JSON),转成Parquet或ORC这类列存储格式。这样Spark读取时只会加载n1、n2这两个需要的列,减少磁盘IO的开销,尤其是当你的DataFrame有很多其他无关列时。

内容的提问来源于stack exchange,提问作者devopslife

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 04:14:48