PySpark内连接异常:Executor分配不均且OOM故障排查
问题分析与解决方案
问题背景
我有两个DataFrame:
df1
id items 1 [[c,b,..], [z,x,..],..] 2 [[d,t,..], [q,a,..],..] . . . .
- Schema:
[('id', 'string'), ('items', 'array<array<string>>')] - 总记录数:157
- 保存时分区数:20
- 数据总大小:700MB
df2
id rand 1 a 2 b 4 h . . . .
- 总记录数:150万
- 保存时分区数:10
- 数据总大小:200MB
基于id列执行inner join后保存DataFrame时,任务耗时极久,Spark UI显示:
- 任务执行到固定数量后卡住(比如shuffle分区设为100时卡在76,500时卡在476,50时卡在28);
- 任务未合理分配到Executor,多数shuffle读取量/记录集中在10个Executor,尽管有更多Executor可用;
- 有Shuffle Read/记录的任务均失败,无Shuffle Read/记录的任务可成功;
- 工作仅分配给8-10个Executor,部分任务执行中失败,其余Executor无负载。
失败任务的Executor输出:
# # java.lang.OutOfMemoryError: Java heap space # -XX:OnOutOfMemoryError="kill %p" # Executing /bin/sh -c "kill 12155"...
Spark配置:
spark_config["spark.executor.memory"] = "24G" spark_config["spark.executor.memoryOverhead"] = "8G" spark_config["spark.executor.cores"] = "24" spark_config["spark.driver.memory"] = "10G" spark_config["spark.sql.shuffle.partitions"] = "50" spark_config["spark.default.parallelism"] = "50" spark_config["spark.dynamicAllocation.enabled"] = "true" spark_config["spark.sql.execution.arrow.pyspark.enabled"] = "true" spark_config["spark.shuffle.service.enabled"] = "true" spark_config["spark.dynamicAllocation.minExecutors"] = "100" spark_config["spark.dynamicAllocation.maxExecutors"] = "500" spark_config["spark.submit.deployMode"] = "client" spark_config["spark.yarn.queue"] = "default"
已尝试调整shuffle分区数、默认并行度、增加Executor内存,均无效。
补充信息:
- PySpark版本:2.4.8
- 连接代码:
df1 = df1.join(df2, how = 'inner', on = 'id') - 可用Executor数量:100,尝试过增减数量。
解决方案
1. 解决数据倾斜核心问题
从现象判断,核心是数据倾斜+单条记录过大:df1仅157条记录,但每条的items是大嵌套数组,join后相同id的df2大量记录会和该大数组关联,导致部分shuffle分区数据量暴增,触发OOM。
- 拆分df1的大数组:
将嵌套数组拆分为单条记录,降低单条数据体积:from pyspark.sql.functions import explode # 拆分外层数组,每个子数组单独成一行 df1_exploded = df1.withColumn("item", explode("items")).drop("items") - 使用广播join避免shuffle:
df1数据量极小,直接广播到所有Executor,无需对大表df2做shuffle:from pyspark.sql.functions import broadcast df_joined = broadcast(df1).join(df2, how="inner", on="id")
2. 优化Spark配置
- 调整shuffle分区数匹配资源:
当前分区数远低于总核数,导致任务无法充分并行,建议设置为总核数的1-2倍:spark_config["spark.sql.shuffle.partitions"] = "2400" # 100 Executor * 24核 - 临时关闭动态分配:
动态分配可能在任务异常时无法正确扩容,临时固定Executor数量:spark_config["spark.dynamicAllocation.enabled"] = "false" spark_config["spark.executor.instances"] = "100" - 降低Executor核数:
24核单Executor内存压力过大,降低核数让内存分配更均匀:spark_config["spark.executor.cores"] = "8" spark_config["spark.executor.memory"] = "24G" # 保持总内存不变,提升单核可用内存
3. 验证排查步骤
- 查看df1单条记录大小:
from pyspark.sql.functions import length, to_json df1.select("id", length(to_json("items")).alias("item_size")).show() - 查看Spark UI的Shuffle页面,确认分区数据分布,验证倾斜情况。
内容的提问来源于stack exchange,提问作者Chris_007
相关产品推荐
相关产品推荐

