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

Spark广播Join过慢且Executor利用率低问题排查求助

Spark广播Join中Executor利用率不足的原因分析

问题场景

Spark Standalone集群配置:56个Executor,每个7核、44GiB内存。执行广播Join代码如下:

res = df_chunk.join(
    broadcast(df_subset),
    on=(
        (df_chunk["__file_id__"] == df_subset["file_id"])
        & (df_chunk["__row_id__"] == df_subset["row_id"])
    ),
    how="inner",
)

异常现象:

  • Join任务共3920个,3910个1秒内完成,10个耗时约10分钟
  • Spark UI显示仅约8个Executor参与Join执行
  • 预期广播Join无数据倾斜,但实际出现任务耗时不均

补充数据集拆分逻辑

为适配基础设施,采用以下代码拆分大数据集:

def merge_dataframes(data_frames: List['DataFrame']) -> 'DataFrame':
    """
    纵向合并DataFrame列表
    """
    if len(data_frames) < 2:
        return data_frames[0]

    _df = data_frames.pop(0)
    while data_frames:
        _df = _df.union(data_frames.pop(0))

    return _df

to_chunk = []
for file in chunk:
     file_id = FIS.get_file_index(file)

     df_file = spark.read.parquet(file)
     df_file = df_file.withColumn('__file_id__', lit(file_id))
     to_chunk.append(df_file)
     file_ids.append(file_id)

     df_chunk = merge_dataframes(to_chunk).repartition(NUM_CORES)
     df_subset = df.filter(df.file_id.isin(file_ids)).repartition(NUM_CORES)

发现实际执行任务数等于chunk_size,单个Parquet文件约3GB。

核心原因分析

1. 广播Join并未消除大表的数据倾斜

广播Join仅将小表(df_subset)广播至所有Executor,但大表(df_chunk)的分区数据分布不均依然会导致数据倾斜:

  • 单个Parquet文件3GB,通过循环union合并后直接调用repartition(NUM_CORES),未指定分区键,Spark默认的Hash分区器无法保证数据均匀分配,若__file_id__或__row_id__存在热点值,会导致部分分区数据量远超其他分区,形成耗时极长的任务。
  • 3920个任务对应56×7的核心数,但多数任务处理小分区快速完成,仅10个任务处理超大分区,占用所在Executor的核心资源。

2. Executor利用率不足的直接诱因

  • Spark的任务调度机制中,长任务会持续占用所在Executor的核心,其他Executor完成所有小任务后处于空闲状态,仅剩余处理大分区的Executor在工作,导致整体利用率仅约8个Executor。
  • 循环union操作导致DataFrame的 lineage 过长,可能触发不必要的重计算,影响任务调度效率。
  • 对df_subset手动repartition属于多余操作,广播前Spark会自动优化小表的分区,手动 repartition 反而可能打乱数据分布,增加调度开销。

3. 任务数等于chunk_size的根源

循环中每次处理一个文件就合并并 repartition,但repartition(NUM_CORES)未指定与Join键相关的分区键,无法真正打散Parquet文件的原始数据分布,导致每个原始文件对应一个或多个大分区,实际有效任务数等同于chunk_size。

优化建议

  • 指定Join相关的分区键进行 repartition:对df_chunk按__file_id__和__row_id__组合键分区,确保数据均匀分布:
    df_chunk = merge_dataframes(to_chunk).repartition(NUM_CORES, "__file_id__", "__row_id__")
    
  • 简化DataFrame合并逻辑:使用spark.unionByName替代循环union,减少lineage复杂度:
    df_chunk = spark.unionByName(to_chunk) if len(to_chunk) > 1 else to_chunk[0]
    
  • 移除df_subset的 repartition:广播前无需手动调整小表分区:
    df_subset = df.filter(df.file_id.isin(file_ids))
    
  • 定位并处理热点键:通过统计分组数据量找到热点键,对热点键单独拆分:
    df_chunk.groupBy("__file_id__", "__row_id__").count().orderBy(F.desc("count")).show(20)
    

内容的提问来源于stack exchange,提问作者Nipun Wijerathne

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 18:58:15