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

EMR处理400TB Parquet数据集时Spark Shuffle FetchFailedException问题

Spark Shuffle FetchFailedException 排查与解决

问题场景

从S3读取400TB Parquet数据集,运行在250台r7.16xlarge实例(每台64 vCore、488 GiB内存)上,执行聚合作业时抛出org.apache.spark.shuffle.FetchFailedException错误。

环境配置

Spark提交命令

spark-submit --conf spark.sql.files.maxPartitionBytes=268435456 \
--master yarn --deploy-mode cluster  --conf spark.yarn.maxAppAttempts=1 \
--conf spark.sql.adaptive.enabled=true --conf spark.dynamicAllocation.enabled=false \
--conf spark.sql.parquet.filterPushdown=true \
--conf spark.sql.adaptive.coalescePartitions.enabled=true \
--conf spark.sql.adaptive.advisoryPartitionSizeInBytes=268435456 \
--conf spark.sql.adaptive.optimizeSkewsInRebalancePartitions.enabled=true \
--conf spark.sql.adaptive.rebalancePartitionsSmallPartitionFactor=.5 \
--conf spark.sql.adaptive.coalescePartitions.parallelismFirst=false \
--conf spark.sql.adaptive.coalescePartitions.initialPartitionNum=30000 \
--conf spark.sql.adaptive.localShuffleReader.enabled=true \
--conf spark.shuffle.io.connectionTimeout=8000 \
--conf spark.network.timeout=50000s  --conf spark.files.fetchTimeout=600s \
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.kryoserializer.buffer.max=1g \
--conf spark.memory.storageFraction=0.05 --conf spark.memory.fraction=.8 \
--conf spark.shuffle.compress=true --conf spark.shuffle.spill.compress=true \
--conf spark.hadoop.fs.s3.multipart.th.fraction.parts.completed=0.99 \
--conf spark.sql.objectHashAggregate.sortBased.fallbackThreshold=4000000 \
--conf spark.reducer.maxReqsInFlight=1 \
--conf spark.shuffle.service.enabled=true \
--conf spark.shuffle.io.maxRetries=10 \
--conf spark.executor.cores=5 \
--conf spark.executor.instances=3000 \
--conf spark.executor.memory=36g --conf spark.driver.memory=60g --conf spark.executor.memoryOverhead=4g --conf spark.driver.memoryOverhead=4g \
--conf spark.hadoop.fs.s3a.fast.output.enabled=true \
--conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p' -Djavax.net.ssl.trustStore=/home/hadoop/.config/certs/InternalAndExternalTrustStore.jks" --conf spark.driver.extraJavaOptions="-XX:+UseG1GC " \
test.py

作业代码

spark.sql(''' create or replace temporary view vwds_cnt_output as
    select
      k
    , m
    , count(distinct(psen, psed)) as kaw
    , count(distinct (psen, psed, pate)) as kt
      from vt
      group by k,m ''')

finaldf000=spark.sql(''' select * from vwds_cnt_output ''')
test_data00='s3://sea/seed/ktr/'
finaldf000.write.mode("overwrite").parquet(test_data00)

排查与解决办法

1. 网络与Shuffle服务优化

  • 调整Shuffle网络超时:将spark.shuffle.io.connectionTimeout从8000ms(8秒)调至60000ms(60秒),避免大Shuffle数据拉取超时;spark.network.timeout保持50000s的同时,确认YARN NodeManager/ResourceManager的超时参数(如yarn.nodemanager.container-executor.task.timeout)与Spark配置匹配,防止容器被误杀。
  • 优化Shuffle拉取并发:spark.reducer.maxReqsInFlight=1过于保守,可调整至5-10提升数据拉取效率;spark.shuffle.io.maxRetries保持10或调至15,增加重试次数应对临时网络波动。
  • 确认Shuffle服务状态:检查所有节点的YARN Shuffle Service是否正常运行,确保Executor退出后Shuffle数据不会丢失。

2. 数据倾斜与分区调整

  • 检测并解决数据倾斜:通过Spark UI查看Stage任务执行时间分布,若存在少数任务耗时远超平均,说明group by k,m存在倾斜。可采用以下方案:
    • 对倾斜的key添加随机前缀,先做局部聚合,再去掉前缀做全局聚合;
    • 新增配置spark.sql.adaptive.skewJoin.enabled=true,并设置spark.sql.adaptive.skewedPartitionThresholdInBytes=1073741824(1GB),让Spark自动识别并处理倾斜分区。
  • 优化分区大小:当前spark.sql.adaptive.advisoryPartitionSizeInBytes=256MB,对于大聚合场景可调至512MB或1GB,减少总分区数降低Shuffle开销;spark.sql.adaptive.coalescePartitions.initialPartitionNum=30000设置过小,建议调整为与Executor总核数匹配的倍数(如3000*5=15000)。

3. 内存与GC优化

  • 调整Executor内存分配:每台r7.16xlarge实例跑12个Executor(64核/5核 per Executor),当前每个Executor内存+Overhead=40g,总占用480g接近实例内存上限,易引发GC或OOM。建议改为spark.executor.memory=32g、spark.executor.memoryOverhead=8g,预留更多内存给系统;同时将spark.memory.fraction从0.8调至0.7,减少执行内存占比,降低GC频率。
  • 优化G1GC参数:将-XX:InitiatingHeapOccupancyPercent=35调至40-45,减少GC触发频率;添加-XX:MaxGCPauseMillis=200控制GC停顿时间,避免因长时间GC导致网络超时。

4. S3读写优化

  • 提升S3读取效率:确认spark.sql.parquet.filterPushdown=true生效,若vt视图有过滤条件需确保下推至S3;设置spark.hadoop.fs.s3a.block.size=134217728(128MB)匹配Parquet文件块大小;检查S3网络带宽是否充足,避免因读取过慢导致Executor空闲超时。
  • 调整S3多部分上传参数:spark.hadoop.fs.s3.multipart.th.fraction.parts.completed=0.99可改为0.95,减少多部分上传的等待时间,优化写入阶段性能。

5. 其他配置调整

  • 关闭spark.sql.adaptive.localShuffleReader.enabled:大Shuffle场景下,本地ShuffleReader可能导致数据分布不均,增加拉取压力,建议关闭该配置。
  • 优化Shuffle缓存:设置spark.shuffle.file.buffer=64k(默认32k)提升Shuffle写缓存效率;开启spark.shuffle.io.preferDirectBufs=true,使用直接内存减少内存拷贝开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 12:47:07