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
相关产品推荐
相关产品推荐

