Spark任务SortMergeJoin与磁盘溢出问题排查求助
问题排查求助
环境信息
- 运行环境:50节点AWS EMR集群(实例类型r5.24xlarge,单节点配置:96 vCore、768 GiB内存、1000 GiB EBS存储)
- Spark版本:3.2.1
故障现象
- Stage 20因磁盘溢出失败,日志显示:
INFO UnsafeExternalSorter: Thread 163 spilling sort data of 20.4 GiB to disk (107 times so far) - 故障节点磁盘使用率超90%,存在两个任务耗时极长
- 主表原始1500万行,经代码扩展至70倍后,已广播所有关联小表,但仍触发SortMergeJoin
- 某段代码执行后数据全部集中到单个分区,怀疑数据倾斜但无法定位
附Spark配置
spark_conf.set("spark.executor.cores", 5) spark_conf.set("spark.executor.memory", "35g") spark_conf.set("spark.executor.instances", 479) spark_conf.set("spark.yarn.executor.memoryOverhead", "5g") spark_conf.set("spark.driver.memory", "35g") spark_conf.set("spark.network.timeout", 1000) spark_conf.set("spark.sql.autoBroadcastJoinThreshold", "500m") spark_conf.set("spark.default.parallelism", 21000) spark_conf.set("spark.memory.storageFraction", 0.2) spark_conf.set("spark.driver.cores", 5) spark_conf.set("spark.sql.shuffle.partitions", 21000) spark_conf.set("spark.dynamicAllocation.enabled", "false") spark_conf.set("spark.sql.execution.arrow.pyspark.enabled", "true") spark_conf.set("spark.cleaner.referenceTracking.cleanCheckpoints", "true") spark_conf.set("spark.checkpoint.compress", "true") spark_conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", -1)
相关代码片段
asm = self.spark_session.createDataFrame(asm_pd, schema=asm_schema).drop("2009") asm_b = f.broadcast(asm) df.groupBy(f.spark_partition_id()).count().show(200) print("before asm join") # assumption name is under index df = df.join(asm_b, df["join_asm"] == asm_b["index"]).drop("join_asm", "index") df.groupBy(f.spark_partition_id()).count().show(200) print("after asm join")
排查方向及建议
1. 确认数据倾斜是否为核心原因
从现象看,单个分区数据集中、部分任务耗时极长、磁盘溢出频繁,高度符合数据倾斜特征,可通过以下方式定位:
- 对比
join前后df.groupBy(f.spark_partition_id()).count()的输出,确认是否有分区数据量远超其他分区 - 统计
join_asm字段的频次:df.groupBy("join_asm").count().orderBy(f.desc("count")).show(50),查看是否存在某个/某些key的行数占比极高 - 查看Spark UI的Stage详情,找到耗时极长的任务,对比其处理数据量与其他任务的平均数据量
2. 广播小表后仍触发SortMergeJoin的原因排查
- 检查
asm表实际大小:若表大小超过spark.sql.autoBroadcastJoinThreshold=500m,或Spark缺失表的统计信息,会自动 fallback 到SortMergeJoin。可通过asm.count()、Spark UI的Storage页面确认表大小 - 验证手动广播是否生效:查看Spark UI的Jobs页面,对应join阶段是否有Broadcast Hash Join标识。若未生效,需检查join条件的字段类型是否匹配(如
df["join_asm"]为字符串、asm_b["index"]为整数),类型不匹配会导致Spark无法识别等值join的广播条件
3. 单分区数据集中的原因排查
- 若
join后出现单分区数据集中,大概率是join的key存在极端值(如null、空字符串或高频key),导致该key对应的所有数据被shuffle到同一分区 - 检查
join_asm字段的null值数量:df.filter(f.col("join_asm").isNull()).count(),Spark默认会将所有null值分配到同一个分区 - 排查
join前是否有repartition/coalesce或groupBy等shuffle操作,这类操作也可能导致分区分布不均
4. 资源配置优化建议
- 当前
spark.executor.cores=5,单节点r5.24xlarge有96vCore,集群总核数4800,但executor总核数仅2395,资源利用率不足。可调整为spark.executor.cores=8或16,减少executor实例数,提升单executor的资源配比,降低调度开销 spark.memory.storageFraction=0.2设置过低,存储内存占比小,缓存数据易被驱逐,增加磁盘IO。可调整为0.3或0.4,提升存储内存占比,减少溢出- 确认
spark.shuffle.spill.compress=true(默认开启),压缩溢出数据,降低磁盘占用压力
内容的提问来源于stack exchange,提问作者user19627118
相关产品推荐
相关产品推荐

