Spark集群执行GroupBy操作触发java.lang.OutOfMemoryError错误求助
你碰到的这个内存溢出问题,在处理大分组数据时很常见,结合你的7节点512GB总内存的Cloudera Spark 2.1.0集群配置和代码逻辑,我给你几个针对性的解决方案:
1. 先排查是否存在超大分组
这个错误最核心的原因大概率是单个userid对应的filename数量过多,导致collect_list生成的数组大小超过了JVM的最大数组限制(通常是Integer.MAX_VALUE,约20亿元素,但实际内存不足时会更早触发)。先跑个查询确认是否有极端数据:
from pyspark.sql.functions import desc # 查看每个userid对应的文件数,按数量降序排列 ndf.groupBy("name").count().orderBy(desc("count")).show(10)
如果发现某个userid对应几十万甚至上百万个文件,那直接用collect_list合并肯定会炸内存,这种情况得先评估业务需求:是否真的需要把所有文件名合并成一个列表?如果可以接受拆分或者其他存储方式,优先调整业务逻辑。
2. 调整分区数,避免过度 repartition
你代码里两次设置repartition(20000)完全没必要,甚至会加剧内存压力:7节点集群,就算每个节点跑8个executor,总executor数也就56个,分区数设置为executor数的2-4倍(比如100-200)就足够并行处理了。20000个分区会导致每个分区的数据量极小,任务调度开销暴增,同时每个分区的元数据(比如分区内的统计信息、对象引用)会占用大量内存,反而容易触发OOM。
修改后的代码可以去掉不合理的repartition,或者换成更合理的数值:
# 替换成合理的分区数,比如150(根据你的数据量调整) ndf = ndf.repartition(150) by_user_df = ndf.groupBy(ndf.name) \ .agg(collect_list("file_name")) \ .withColumnRenamed('collect_list(file_name)', 'file_names') # 如果后续需要按userid分区,推荐用repartitionByRange,避免哈希分区不均匀 by_user_df = by_user_df.repartitionByRange("name") by_user_df.count()
3. 调整集群内存参数(Cloudera环境)
通过Cloudera Manager调整Spark的内存配置,针对你的场景重点调这几个参数:
- Executor内存:每个Executor分配16GB左右(7节点总内存512GB,去掉系统和其他服务占用,每个节点大概能分配60-70GB,按4个executor/节点算,每个16GB是合理的)
- Executor堆外内存:设置
spark.executor.memoryOverhead为2-4GB,防止堆外内存不足导致的OOM - Driver内存:如果执行
count()时Driver端OOM,把spark.driver.memory调到8-16GB,因为count()会把结果汇总到Driver端
4. 替代collect_list的方案(针对超大分组)
如果确实有超大分组无法避免,考虑用以下方式替代collect_list:
- 字符串拼接:用
concat_ws把文件名用分隔符(比如逗号、竖线)拼接成一个字符串,但要注意JVM字符串长度限制,不过比数组的内存开销小一些:from pyspark.sql.functions import concat_ws by_user_df = ndf.groupBy("name").agg(concat_ws(",", collect_list("file_name")).alias("file_names")) - 保留多行结构:不合并成列表,后续处理时按userid批量读取,避免一次性加载所有数据到内存
- 分桶表:如果是经常需要按userid分组的场景,可以把原始表做成分桶表,按
name分桶,减少分组时的数据 shuffle 开销
5. 考虑升级Spark版本
Spark 2.1.0是比较老的版本(2017年发布),后续版本(比如2.3+)对collect_list的内存管理、分组聚合的性能都有优化,Cloudera对应的CDH版本也有支持,升级后可能会缓解这类内存问题。
内容的提问来源于stack exchange,提问作者kupe

