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

Spark EMR任务疑问:AQE是否决定任务数?是否需调整Executor?

问题解答

问题1:这1000个任务是否是Spark在shuffle阶段后由AQE定义的?

是的,这1000个任务是AQE(自适应查询执行)在shuffle阶段后根据实际数据处理情况生成的,具体原因如下:

  • 你已启用spark.sql.adaptive.coalescePartitions.enabled=true,并配置了spark.sql.adaptive.advisoryPartitionSizeInBytes=256m、spark.sql.adaptive.coalescePartitions.minPartitionSize=256m,AQE会在shuffle阶段后根据实际输出的数据量,将过小的分区合并为符合目标大小的分区。
  • 你的任务包含SQL聚合操作,这类操作会大幅缩减数据量:初始读取的141000个输入分区经过聚合后总数据量骤降,按照256MB的分区大小计算,最终仅能生成约1000个分区(对应1000个任务)。
  • 你设置的spark.sql.adaptive.coalescePartitions.minPartitionNum=27000未生效,是因为该参数优先级低于实际数据量限制——AQE不会强制生成超出数据量支撑的分区数,当总数据量除以目标分区大小的结果远小于27000时,会以实际计算的分区数为准。

问题2:若集群最多仅使用1000个任务(即最多占用1000核),是否应减少executor的数量?

必须减少executor数量,否则会造成严重资源浪费:

  • 当前配置1800个executors、每个5核,总计9000个可用核心,但任务仅需1000个核心,意味着约8/9的executors处于空闲状态,这些空闲实例会占用内存、CPU资源却不产生实际计算价值。
  • 基于当前任务的并行度需求,建议将executor数量调整为200-300个(200*5=1000核,刚好匹配任务所需核心数,预留少量冗余应对波动)。
  • 结合你使用的r5.16xlarge实例(64vCore),每个实例可部署12个5核executors(12*5=60vCore,剩余4vCore留给系统进程),200个executors仅需约17台实例,远少于当前的150台,可进一步减少集群实例数量以降低成本。

当前Spark提交配置

--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=256m \
--conf spark.sql.adaptive.coalescePartitions.minPartitionSize=256m \
--conf spark.sql.adaptive.coalescePartitions.parallelismFirst=false \
--conf spark.sql.adaptive.coalescePartitions.minPartitionNum=27000 \
--conf spark.network.timeout=5400s  --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.network.timeout=1200s \
--conf spark.executor.cores=5 \
--conf spark.executor.instances=1800 \
--conf spark.executor.memory=32g --conf spark.driver.memory=60g --conf spark.executor.memoryOverhead=4g --conf spark.driver.memoryOverhead=4g \

任务监控截图

Spark任务阶段监控
Spark任务分区详情

内容的提问来源于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 14:40:59