group by本质是shuffle操作,需将相同分组键的数据重分区并跨节点搬运至同一executor聚合,而非单机分片计算;数据量大时易引发网络/磁盘io压力及oom,主因常为数据倾斜或高基数列直接分组。

Spark SQL的GROUP BY本质是Shuffle,不是单机分组
Spark SQL里写GROUP BY看起来像SQL,但执行时会触发全量Shuffle——所有相同分组键的记录必须被拉到同一个Executor上才能聚合。这不是“先分片再各自算完拼结果”,而是“重分区+跨节点搬运+本地聚合”。所以数据量越大,网络和磁盘IO压力越明显。
常见错误现象:java.lang.OutOfMemoryError: Java heap space出现在groupBy后,往往不是内存配少了,而是某个key倾斜(比如某地区有千万条订单),导致单个task处理数据远超其他task。
- 避免在
GROUP BY中使用高基数列(如user_id)直接分组,除非你明确需要每个用户一条结果 - 若需按
user_id聚合但担心倾斜,先加盐(salt):用concat(user_id, '_', floor(rand() * 10))构造新分组键,聚合后再二次合并 - 检查
spark.sql.adaptive.enabled是否开启(Spark 3.2+默认true),它能在运行时自动拆分长尾task
agg()里别滥用collect_list或collect_set
这两个函数会把整个分组的所有原始值存进内存,极易OOM。例如groupBy("region").agg(collect_list("order_id")),当某region有50万订单时,单个task就要加载50万个字符串对象。
使用场景:仅当业务强依赖“列出全部明细”(如审计日志归档),且已确认该分组最大规模可控(
- 替代方案优先选
count、approx_count_distinct、first、max等常数空间聚合函数 - 真要取Top N,用
array_sort(array_agg(...), ...)[0] as top1比先collect再sort更省内存 - Spark 3.4+支持
aggregate高阶函数做流式折叠,可替代部分collect_list + UDF逻辑
小表JOIN大表聚合前,先广播小表
如果聚合前要JOIN维度表(比如product_id → category),而维度表
性能影响:未广播时,JOIN + GROUP BY整体耗时可能比单纯GROUP BY高3–5倍;广播后,JOIN退化为map-side lookup,基本不增加shuffle负担。
- 显式调用
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "10485760")(10MB) - 或对DataFrame手动广播:
df_dim.broadcast(),再用join(..., broadcast=True) - 注意:广播只对
BroadcastHashJoin生效,SortMergeJoin不会降级——确保JOIN条件是等值且无复杂表达式
聚合结果写入前,用repartition控制输出文件数
直接df.groupBy(...).agg(...).write.parquet(...),输出文件数 = Shuffle后分区数,默认由spark.sql.shuffle.partitions(通常200)决定。但200个小文件对下游查询不友好,尤其Hive表统计信息收集会变慢。
容易踩的坑:用coalesce(1)强行压成1个文件——这会让所有数据涌向单个task,拖慢整体完成时间,还可能OOM。
- 合理做法:按业务主键
repartition(50, "region"),既减少文件数,又保持数据分布均匀 - 若下游是Hive,建议文件大小控制在128–256MB,可用
df.repartitionByRange配合排序避免热点 - 写Parquet前加
option("compression", "snappy"),压缩比和解压速度平衡较好
实际聚合链路中,最易被忽略的是Shuffle阶段的中间数据序列化开销。如果你用的是Kryo序列化器,记得注册自定义类;否则默认Java序列化会让Shuffle体积膨胀2–3倍——这点在GROUP BY后接UDF时尤为致命。










