You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 12:17:41