Spark调优:如何提升并行度解决OOM(137错误)问题?
解决Spark DataFrame OOM(137错误)的实用方案
除了升级硬件,你可以从以下几个方向尝试解决问题:
1. 精准定位OOM发生阶段
先通过Spark UI查看失败的Stage详情:
- 确认是Shuffle阶段(比如Shuffle Read/Write)还是普通计算阶段OOM
- 查看每个Task处理的数据量,如果单个Task数据量过大,说明分区数调整未生效或仍不足
2. 确保分区数配置真正生效
你设置的spark.sql.shuffle.partitions仅对后续的Shuffle操作生效,若之前的操作已经生成固定分区的DataFrame,需要:
- 在会话启动时就设置该参数(而非中途),比如在
spark-submit中添加:--conf spark.sql.shuffle.partitions=1000 - 对已有DataFrame重分区后,执行Action操作触发分区变更,比如:
df = df.repartition(1000) df.count() # 触发分区生效,再执行后续操作 - 检查是否有操作强制覆盖了分区数:比如某些Join/GroupBy操作若依赖小表,可能会自动调整分区,此时需要手动指定分区数
3. 优化Shuffle内存配置
针对Shuffle阶段OOM,调整以下参数:
- 提高Shuffle内存占比:将Executor内存中分配给Shuffle的比例从默认20%提高到30%-40%:
spark.conf.set("spark.shuffle.memoryFraction", "0.3") - 增大堆外内存:若OOM是堆外内存不足(比如序列化、网络缓存),调大内存预留:
spark.conf.set("spark.executor.memoryOverhead", "4096") # 单位MB,根据集群情况调整 - 开启Shuffle数据压缩:确保压缩已启用(默认开启),减少内存和磁盘IO压力:
spark.conf.set("spark.shuffle.spill.compress", "true")
4. 优化Checkpoint操作
Checkpoint若未正确落地,仍会占用内存:
- 确保Checkpoint路径指向分布式存储(如HDFS),而非本地磁盘
- Checkpoint后必须重新赋值并触发Action,保证数据真正写入磁盘:
df = df.checkpoint() df.count() # 触发Checkpoint执行 - 若不需要容错,可尝试
localCheckpoint()(仅适合测试环境,不保证数据安全)
5. 精简数据与优化操作
- 过滤冗余列:仅保留后续操作必需的列,减少单条数据的内存占用:
df = df.select("col1", "col2", "required_col") - 优化数据类型:将占用内存大的类型替换为紧凑类型,比如:
- 用
IntegerType替代StringType存储数字 - 用
DateType/TimestampType替代StringType存储日期
- 用
- 使用广播Join:若Join操作中有小表,将小表广播,避免大表Shuffle:
from pyspark.sql.functions import broadcast df_join = df_large.join(broadcast(df_small), on="join_key") - 清理无用缓存:若之前使用了
cache()/persist(),及时释放内存:df_old.unpersist()
6. 调整Task并行度
- 增加Executor数量:若集群还有资源,提高
spark.executor.instances,让更多Executor分担任务 - 减少单个Executor的核心数:比如从4核改为2核,降低每个Executor的任务负载,减少内存竞争
内容的提问来源于stack exchange,提问作者John Stud
相关产品推荐
相关产品推荐

