PySpark:如何在Databricks集群节点并行运行独立XGBoost模型?
问题:Databricks并行训练独立XGBoost模型效率低于串行
我需要在Databricks集群的多个节点上分布式运行独立XGBoost模型(每个节点训练一个独立模型,而非将单个模型拆分到多节点)。编写了接收整数参数的训练函数,代码可正常运行,但并行执行效率远低于串行:用sc.parallelize运行2个模型耗时20-30分钟,而串行执行仅需8-10分钟。
训练函数代码如下:
import xgboost as xgb from scipy.sparse import csr_matrix, save_npz, load_npz def train_xgboost_bootstrap_models(n: int): train_id = pd.read_pickle(f'/dbfs/Data/01/train_id_{n}.pickle' ).sort_values(by='row_index') sp_train = load_npz(f'/dbfs/Data/01/sp_train_{n}.npz') id_oob = pd.read_pickle(f'/dbfs/Data/01/id_oob_{n}.pickle' ).sort_values(by='row_index') sp_oob = load_npz(f'/dbfs/Data/01/sp_oob_{n}.npz') dmat_train = xgb.DMatrix(sp_train, label=train_id.y.values) dmat_oob = xgb.DMatrix(sp_oob, label=id_oob.y.values) param = { 'objective' : "reg:tweedie", 'tweedie_variance_power' : 1.5, 'eval_metric' : "rmse", 'tree_method' : 'hist', 'booster' : 'gbtree' } nrounds = 100 model = xgb.train(params=param, dtrain = dmat_train, num_boost_round=nrounds, evals=[(dmat_train, 'train'), (dmat_oob, 'eval')]) return model # 并行执行 models = sc.parallelize(list(range(2))).map(lambda x: train_xgboost_bootstrap_models(x)).collect() # 串行执行对比 train_xgboost_bootstrap_models(n=0) train_xgboost_bootstrap_models(n=1)
问题原因分析
- Spark调度开销过高:用Spark RDD的
map执行独立训练任务,会引入任务调度、数据序列化/反序列化、节点通信等额外开销。当单模型训练耗时较短时,这些开销会抵消并行收益,甚至导致总耗时增加。 - 数据本地性缺失:直接从DBFS加载数据时,Spark任务可能被调度到未缓存该数据的节点,引发远程IO,拖慢训练速度。
- 节点内资源竞争:XGBoost的
hist树方法默认使用多核心训练,并行运行多个模型时,节点内CPU资源被抢占,单个模型的训练效率下降。
优化方案
1. 改用Databricks任务级并行(替代Spark RDD)
使用Databricks的Job API提交独立训练任务,调度开销更低且资源隔离性更好:
from databricks.sdk import WorkspaceClient import concurrent.futures w = WorkspaceClient() def run_single_model(n): # 提前将单模型训练逻辑封装为独立Notebook并创建Job run_result = w.jobs.run_now( job_id="YOUR_TRAIN_JOB_ID", notebook_params={"model_index": str(n)} ) # 等待任务完成,返回模型存储路径或结果 return run_result # 用多线程并行提交任务 with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor: futures = [executor.submit(run_single_model, idx) for idx in range(2)] results = [future.result() for future in futures]
2. 优化数据本地性
将训练数据提前复制到节点本地磁盘,避免远程加载:
# 在训练函数内添加数据本地化逻辑 def train_xgboost_bootstrap_models(n: int): # 复制DBFS数据到本地磁盘 dbutils.fs.cp(f"/dbfs/Data/01/train_id_{n}.pickle", f"file:/local_disk0/train_id_{n}.pickle") dbutils.fs.cp(f"/dbfs/Data/01/sp_train_{n}.npz", f"file:/local_disk0/sp_train_{n}.npz") dbutils.fs.cp(f"/dbfs/Data/01/id_oob_{n}.pickle", f"file:/local_disk0/id_oob_{n}.pickle") dbutils.fs.cp(f"/dbfs/Data/01/sp_oob_{n}.npz", f"file:/local_disk0/sp_oob_{n}.npz") # 从本地加载数据 train_id = pd.read_pickle(f'/local_disk0/train_id_{n}.pickle').sort_values(by='row_index') sp_train = load_npz(f'/local_disk0/sp_train_{n}.npz') id_oob = pd.read_pickle(f'/local_disk0/id_oob_{n}.pickle').sort_values(by='row_index') sp_oob = load_npz(f'/local_disk0/sp_oob_{n}.npz') # 后续训练逻辑不变...
3. 限制单个XGBoost模型的资源占用
调整XGBoost参数,限制每个模型使用的CPU核心数,避免节点内资源竞争:
param = { 'objective': "reg:tweedie", 'tweedie_variance_power': 1.5, 'eval_metric': "rmse", 'tree_method': 'hist', 'booster': 'gbtree', 'nthread': 4 # 根据节点总核心数调整,比如节点有16核则设置为4,同时并行4个模型 }
4. 优化Spark任务粒度
如果坚持使用Spark,改用mapPartitions减少任务调度次数:
# 设置partition数等于模型数量,每个partition跑一个模型 models = sc.parallelize(list(range(2)), numSlices=2).mapPartitions( lambda num_list: [train_xgboost_bootstrap_models(n) for n in num_list] ).collect()
内容的提问来源于stack exchange,提问作者gabagool
相关产品推荐
相关产品推荐

