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

Spark EMR作业因Shuffle Read Fetch Wait Time长达4小时失败,求优化方法

Spark作业Shuffle Read Fetch Wait Time过长的优化方案

处理65TB数据的Spark作业因Shuffle Read Fetch Wait Time过长(达4小时)失败,已启用自适应查询执行(AQE)。

当前spark-submit配置

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=36000 \
--conf spark.sql.adaptive.localShuffleReader.enabled=true \
--conf spark.network.timeout=6000s  --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.executor.cores=5 \
--conf spark.executor.instances=3600 \
--conf spark.sql.shuffle.partitions=36000 \
--conf spark.executor.memory=32g --conf spark.driver.memory=60g --conf spark.executor.memoryOverhead=8g --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

监控截图

监控截图1
监控截图2

可行优化方法

  • 调整Shuffle请求并发数:当前spark.reducer.maxReqsInFlight=1极大限制了Reducer同时发起的Shuffle请求数量,直接拖慢Fetch速度。建议逐步调高该值(从10开始测试,最高可到64),平衡并发请求与网络压力,减少等待时间。
  • 强化AQE倾斜处理:虽然已启用倾斜分区优化,但需确认spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes(默认256MB)和spark.sql.adaptive.skewJoin.skewedPartitionFactor(默认10)配置适配65TB数据场景,让AQE更精准拆分倾斜分区,减少跨节点远程数据拉取。
  • 延长Shuffle超时阈值:当前spark.files.fetchTimeout=600s对于大作业可能不足,建议调高至1800s(30分钟)或更长;同时检查spark.shuffle.io.connectionTimeout(默认120s),避免网络波动导致连接超时中断Fetch。
  • 优化Executor资源配比:当前每个Executor分配5核32G内存,若CPU使用率偏低,可适当降低spark.executor.cores(如调整为3-4),增加Executor实例数提升整体并行度;同时确保spark.executor.memoryOverhead足够应对Shuffle堆外内存需求,避免GC停顿或数据溢出。
  • 启用External Shuffle Service:开启spark.shuffle.service.enabled=true,避免Executor退出后丢失Shuffle数据,减少重复拉取;若Shuffle临时数据存储在S3,可切换到本地磁盘或EMR FS等低延迟存储,降低IO等待。
  • 缩减Shuffle数据量:确认spark.sql.parquet.filterPushdown=true已生效,在数据源层过滤无效数据;同时梳理作业逻辑,移除不必要的join、groupBy等Shuffle操作,从根源减少需要传输的数据量。
  • 调整AQE分区合并策略:当前spark.sql.adaptive.coalescePartitions.parallelismFirst=false优先保证分区大小,可能导致大作业并行度不足。可尝试设置为true,优先保证足够的并行度,再调整分区大小,提升Shuffle处理效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 11:14:54