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
监控截图


可行优化方法
- 调整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
相关产品推荐
相关产品推荐

