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


内容的提问来源于stack exchange,提问作者user3858193
相关产品推荐
相关产品推荐

