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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 22:40:29