Dataproc集群下PySpark代码未在所有executor上并行是什么原因?
问题根因
Spark调度器不会启动多余executor执行任务,仅用2个executor运行是以下多个常见配置共同作用的结果:
- 分区数不足:
sc.parallelize()生成RDD的默认分区数由spark.default.parallelism参数控制,Dataproc在小型集群场景下默认该值为2,RDD仅有2个分区对应2个并行task,最多只会占用2个executor资源,没有更多task需要调度到剩余executor上运行。 - 动态资源分配阈值未触发:Dataproc默认开启Spark动态资源分配,初始仅启动少量executor,仅当待调度task队列堆积超过阈值时才会扩容更多executor。你运行的测试任务计算量过轻,在调度器触发扩容逻辑前就已经执行完成,不会触发新executor的调度。
- 抢占式节点调度优先级限制:Dataproc默认优先将任务调度到普通非抢占式executor上,仅当普通节点资源不足时才会将任务分发到抢占式executor运行,你的测试任务资源消耗极低,3个主executor的空闲资源完全可以覆盖任务需求,调度器不会触发抢占式节点的任务调度。
- executor资源配置冗余:如果单个executor配置的CPU核心、内存额度较高,2个executor的空闲资源已经可以容纳当前任务的所有并行task,调度器为了避免跨节点数据传输开销,会优先将所有task调度到已启动的executor上,不会调度到其他空闲节点。
验证与修复方案
- 首先确认RDD分区数:运行
sc.parallelize(range(1,1000000)).getNumPartitions(),如果返回值为2即可确认是分区数不足问题,测试时手动指定更多分区即可触发多executor调度,示例代码如下:
import numpy as np import socket # 手动指定30个分区,数量超过集群总CPU核心数即可保证调度所有可用executor set(sc.parallelize(range(1,1000000), 30).map(lambda x : socket.gethostname()).collect())
- 调整动态资源分配配置:如果需要执行轻量任务也能用到所有executor,可以修改Spark配置降低扩容阈值:
- 将
spark.dynamicAllocation.schedulerBacklogTimeout从默认1s调整为100ms - 将
spark.dynamicAllocation.minExecutors设置为6,强制启动所有6个executor
- 将
- 调整抢占式节点调度比例:修改
spark.yarn.preemptibleAllocationRatio参数,该参数控制抢占式executor承载的task比例,默认值为0.5,调高后可以让更多任务调度到抢占式节点上。 - 固定executor数量测试:提交任务时添加参数
--num-executors 6关闭动态分配的自动扩缩容逻辑,强制Spark启动所有6个executor执行任务。
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

