如何在Spark SQL中对海量数据进行并行分组聚合

星敏酱_7553

星敏酱_7553

2026-09-30

645人浏览

原创

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

如何在spark sql中对海量数据进行并行分组聚合

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时尤为致命。

PHP速学视频免费教程(入门到精通)
PHP速学视频免费教程(入门到精通)

PHP怎么学习?PHP怎么入门?PHP在哪学?PHP怎么学才快?不用担心,这里为大家提供了PHP速学教程(入门到精通),有需要的小伙伴保存下载就能学习啦!

下载

相关标签:

本站声明:本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系admin@php.cn

相关专题

更多
大数据分析工具有哪四个
大数据分析工具有哪四个

大数据分析的四个工具分别是rapidminer、Hpcc、Hadoop和Pentaho bi。大数据分析用于从各种来源生成的原始数据中提取有价值的数据。这些数据帮助我们获得有意义的见解、隐藏的模式、未知的相关性、市场趋势等等,具体取决于行业。大数据分析的主要动机是提供有价值的见解,以便为未来做出更好的决策。php中文网为大家带来了大数据分析的相关教程、以及相关文章等内容,供大家免费下载使用。

2023.06.21

4316

5

Java 大数据处理基础(Hadoop 方向)
Java 大数据处理基础(Hadoop 方向)

本专题聚焦 Java 在大数据离线处理场景中的核心应用,系统讲解 Hadoop 生态的基本原理、HDFS 文件系统操作、MapReduce 编程模型、作业优化策略以及常见数据处理流程。通过实际示例(如日志分析、批处理任务),帮助学习者掌握使用 Java 构建高效大数据处理程序的完整方法。

2025.12.08

1229

12

大数据专业学习教程
大数据专业学习教程

本专题整合了大数据专业学习相关教程,阅读专题下面的文章了解更多详细内容。

2026.01.05

223

5

python处理大数据合集
python处理大数据合集

本专题整合了python处理大数据相关教程,阅读专题下面的文章了解更多详细内容。

2026.01.05

446

22

数据分析工具有哪些
数据分析工具有哪些

数据分析工具有Excel、SQL、Python、R、Tableau、Power BI、SAS、SPSS和MATLAB等。详细介绍:1、Excel,具有强大的计算和数据处理功能;2、SQL,可以进行数据查询、过滤、排序、聚合等操作;3、Python,拥有丰富的数据分析库;4、R,拥有丰富的统计分析库和图形库;5、Tableau,提供了直观易用的用户界面等等。

2023.10.12

3863

8

SQL中distinct的用法
SQL中distinct的用法

SQL中distinct的语法是“SELECT DISTINCT column1, column2,...,FROM table_name;”。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

2023.10.27

831

4

SQL中months_between使用方法
SQL中months_between使用方法

在SQL中,MONTHS_BETWEEN 是一个常见的函数,用于计算两个日期之间的月份差。想了解更多SQL的相关内容,可以阅读本专题下面的文章。

2024.02.23

1009

5

SQL出现5120错误解决方法
SQL出现5120错误解决方法

SQL Server错误5120是由于没有足够的权限来访问或操作指定的数据库或文件引起的。想了解更多sql错误的相关内容,可以阅读本专题下面的文章。

2024.03.06

5701

10

sql procedure语法错误解决方法
sql procedure语法错误解决方法

sql procedure语法错误解决办法:1、仔细检查错误消息;2、检查语法规则;3、检查括号和引号;4、检查变量和参数;5、检查关键字和函数;6、逐步调试;7、参考文档和示例。想了解更多语法错误的相关内容,可以阅读本专题下面的文章。

2024.03.06

2643

4

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
热门推荐
/
最新课程
phpStudy极速入门视频教程
phpStudy极速入门视频教程

共6课时 | 54.6万人学习

独孤九贱(4)_PHP视频教程
独孤九贱(4)_PHP视频教程

共89课时 | 133.3万人学习