基于joblibspark的Databricks集群并行任务均匀分配问题问询
根因
- Spark 默认优先遵循任务本地化调度规则,会优先将任务分配给有数据缓存、或被判定为当前负载更低的 executor,不会主动强制跨节点均匀分发
- Databricks 默认开启 executor 动态分配,任务量较小时会自动缩减实际启用的 executor 数量,哪怕你集群配置了4个 executor,实际调度时可能只有2个处于可用状态
- joblibspark 的
batch_size参数仅控制单批次提交的任务数量,不会干预 Spark 底层的任务分配逻辑,所以单独设置为1无法解决分布不均问题
需要补充的配置要点
1. 关闭Spark本地化等待
强制Spark不等待本地化资源就绪,只要有空闲executor就直接分发任务,在代码开头添加配置:
spark.conf.set("spark.locality.wait", "0")
也可以直接在Databricks集群的Spark配置页添加该参数,全局生效。
2. 固定executor数量,关闭动态分配
避免Databricks自动缩减executor导致可用节点不足,添加以下配置:
spark.conf.set("spark.dynamicAllocation.enabled", "false") # 固定启用4个executor,和你的集群配置匹配 spark.conf.set("spark.executor.instances", "4")
3. 调整joblib并行参数
新增 pre_dispatch="all" 参数,强制将所有任务一次性提交到Spark调度队列,避免分批提交导致的调度不均,修改后的代码如下:
from joblib import Parallel, delayed from joblibspark import register_spark register_spark() output = Parallel(backend="spark", n_jobs=4, # 与实际executor数量匹配,可根据核数总和向上调整 verbose=config.JOBLIB_VERBOSE, batch_size=1, pre_dispatch="all")( delayed(fit_one) (x, model_data=model_data, dlmodel=dlmodel, outdir=outdir, frac=sample_p, score_type=score_type, save=save, verbose=verbose) for x in ZZ)
4. 强制单任务单核心分配(可选,适用于CPU密集型任务)
如果你的fit_one是单线程CPU密集型任务,可以添加以下配置,强制每个任务独占1个核心,避免单个executor同时运行多个任务:
spark.task.cpus=1 # 如有需要可以将单个executor核数设为1,强制每个executor只能跑1个任务,天然均匀分布 spark.executor.cores=1
内容的提问来源于stack exchange,提问作者pauljohn32
相关产品推荐
相关产品推荐

