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

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 04:15:49