ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Modin 分布式 XGBoost:基于 Ray 的 `modin.experimental.xgboost` 模块原理与实战指南

Modin 分布式 XGBoost:基于 Ray 的 `modin.experimental.xgboost` 模块原理与实战指南 Modin 分布式 XGBoost基于 Ray 的modin.experimental.xgboost模块原理与实战指南【免费下载链接】modinModin: Scale your Pandas workflows by changing a single line of code项目地址: https://gitcode.com/gh_mirrors/mo/modinModin 在modin.experimental.xgboost模块中提供了分布式 XGBoost 实现让用户以接近原生 XGBoost 的编程体验在 Ray 集群上完成大规模训练与预测。本文以仓库内 docs/flow/modin/experimental/xgboost.rst 为主线结合modin/experimental/xgboost目录下的完整源码与测试用例从公共接口、内部执行流程、核心算法到限制与测试验证逐层展开读完即可上手使用并理解其底层运行机制。模块定位与整体结构modin.experimental.xgboost是 Modin 为 XGBoost 提供的分布式实现位于modin/experimental/xgboost/目录下由四个文件组成init.py导出公共接口DMatrix、Booster、trainxgboost.py公共接口定义向用户暴露熟悉的 XGBoost APIxgboost_ray.py基于 Ray 执行引擎的内部实现包含 Ray actor 类、数据切分与训练/预测内部函数utils.pyRabit 全归约上下文管理工具负责分布式训练过程中的状态同步。从源码结构看公共接口层通过判断Engine.get()是否为Ray来决定调用哪个后端见 xgboost.py非 Ray 引擎会抛出ValueError(Current version supports only Ray engine.)。也就是说分布式 XGBoost 目前仅支持 Ray 执行引擎这一点在使用前必须确认。公共接口与原生 XGBoost 对齐的三个入口DMatrix分布式数据的容器modin.experimental.xgboost.DMatrix继承自原生xgboost.DMatrix但重写了构造函数使其直接接受 Modin DataFrame。构造时data参数必须为modin.pandas.DataFramelabel参数可以为modin.pandas.DataFrame或modin.pandas.Series源码中以assert强制校验见 xgboost.py内部通过modin.distributed.dataframe.pandas.unwrap_partitions将 DataFrame 解包为行分区对象的延迟引用列表——data部分还会额外获取每个分区所在节点的 IP 信息get_ipTrue用于后续的本地化调度构造时会检查各列 dtype若存在object类型列则抛出ValueError记录元数据data.index、data.columns以及各行的长度row_lengths供预测阶段重建结果 DataFrame 使用。值得注意的限制当前构造函数只支持data与label两个核心参数虽然函数签名中保留了missing、silent、feature_names、feature_types、feature_weights、enable_categorical但 docstring 明确说明weight、base_margin、nthread、group、qid、label_lower_bound、label_upper_bound参数暂不支持。DMatrix还实现了__iter__迭代时依次产出解包后的data与label引用列表。这正是训练流程中X_row_parts, y_row_parts dtrain解包的底层机制。Booster重写 predict 的模型对象modin.experimental.xgboost.Booster继承原生xgboost.Booster仅重写了predict方法。与原生接口的唯一区别是data参数类型改为modin.experimental.xgboost.DMatrix。predict 内部会校验传入数据必须是DMatrix校验训练时的feature_names与预测数据的feature_names是否一致不一致时抛出包含缺失字段明细的ValueError见 xgboost.py调用内部函数_predict(self.copy(), data, **kwargs)返回一个modin.pandas.DataFrame类型的预测结果。train分布式训练入口modin.experimental.xgboost.train与原生的xgboost.train有两个核心差异见 xgboost.pydtrain参数类型必须是modin.experimental.xgboost.DMatrix新增num_actors参数用于控制参与训练的 Ray actor 数量。其余参数与原生保持一致params为 booster 参数字典evals为(DMatrix, name)组成的验证集列表evals_result为训练过程中评估指标的收集字典。训练完成后函数会取出各 worker 通过 Rabit 同步得到的同一份 booster源码注释明确说明因为 Rabit tracking所有 actor 的结果一致只取第一个即可包装为Booster返回并将history评估结果写入用户传入的evals_result字典。环境准备安装与 Ray 运行时初始化docs/usage_guide/advanced_usage/modin_xgboost.rst给出了安装与初始化的完整说明。Modin 默认自带除xgboost之外的全部依赖因此只需额外安装pip install xgboost由于目前仅支持 Ray 引擎需要先初始化 Ray 运行时。单机模式import ray ray.init() # 启动单节点 Ray 运行时若已有 Ray 集群则通过地址连接import ray ray.init(addressauto)从源码看模块运行时会通过ray.cluster_resources()与ray.nodes()探测集群 CPU 与节点拓扑见 xgboost_ray.py因此确保 Ray 集群正确初始化是分布式训练的前提。端到端示例Iris 数据集的训练与预测以下示例来自 modin_xgboost.rst在单节点模式下完整演示了建 DMatrix → 训练 → 预测的流程from sklearn import datasets import ray ray.init() # 启动单节点 Ray 运行时 import modin.pandas as pd import modin.experimental.xgboost as xgb # 加载 sklearn 自带的 iris 数据集 iris datasets.load_iris() # 构造 Modin DataFrame X pd.DataFrame(iris.data) y pd.DataFrame(iris.target) # 构造 DMatrix训练集与测试集 dtrain xgb.DMatrix(X, y) dtest xgb.DMatrix(X, y) # 训练参数 xgb_params { eta: 0.3, max_depth: 3, objective: multi:softprob, num_class: 3, eval_metric: mlogloss, } steps 20 # 收集评估结果 evals_result dict() # 分布式训练 model xgb.train( xgb_params, dtrain, steps, evals[(dtrain, train)], evals_resultevals_result, ) print(fEvals results:\n{evals_result}) # 分布式预测返回 modin.pandas.DataFrame prediction model.predict(dtest) print(fPrediction results:\n{prediction})这里y在 iris 分类场景下可作为 DataFrame 传入而在 test_xgboost.py 的测试中label同时支持pd.DataFrame与pd.Series两种形态测试覆盖了二分类breast_cancer、多分类iris/digits/wine与回归diabetes三类典型任务。Ray 引擎上的内部执行流程docs/flow/modin/experimental/xgboost.rst将内部流程分为训练与预测两大部分以下结合 xgboost_ray.py 源码逐步骤还原。训练流程7 个步骤步骤 1解包 DMatrix 获取分区引用数据以DMatrix对象传入_train通过迭代器解包得到行分区引用列表X_row_parts, y_row_parts dtrain其中X_row_parts中的每个元素是(IP_ref, partition_ref)形式的引用对y_row_parts为标签分区引用列表见 xgboost.py。步骤 2计算 actor 数量_get_num_actors见 xgboost_ray.py处理用户传入的num_actors若未提供按每个 actor 最多使用 2 个 CPU的条件计算num_actors_per_node max(1, int(min_cpus_per_node // 2))再乘以节点数。该条件是为配合多线程 XGBoost 训练每个 worker 使用 2 个线程而选择的目的是在并行度与线程资源之间取得平衡若显式传入 int则必须能被 Ray 集群节点数整除否则抛出断言错误其他类型抛出RuntimeError。num_actors之所以暴露给公共train函数是为了让用户在特定场景下能够手动调优以获得最佳性能。步骤 3创建 ModinXGBoostActor 对象create_actors见 xgboost_ray.py根据集群资源计算每个 actor 的 CPU 数_get_cpus_per_actor下限为 1并利用ModinXGBoostActor.options(resources{node_ip: 0.01})将 actor 均匀调度到各节点。ModinXGBoostActor是以ray.remote(num_cpus0)声明的 Ray actor 类每个实例持有自己的 rank 与线程数nthread。步骤 4数据在 actor 间均匀切分_split_data_across_actors调用_assign_row_partitions_to_actors见 xgboost_ray.py完成分配首先按分区 IP 与 actor IP 进行本地化分配尽量把数据留在原节点避免跨节点传输当某节点分区数不足以填满本节点 actor 时从其他节点的剩余分区中补充最终保证每个 actor 获得的分区数尽可能均匀分配结果是一个形如{actor_rank: ([part_i0, part_i3, ...], [0, 3, ...]), ...}的字典第二个元素记录分区在原始数据中的顺序对y分区的分配还支持data_for_aligning参数保证y与X在 actor 上的顺序严格对齐见_split_data_across_actors中对X_parts_by_actors的复用。步骤 5远程加载训练数据对每个 actor 远程调用set_train_data方法。数据以ray.ObjectRef列表形式传入 actor 后会自动物化ray.ObjectRef - pandas.DataFrameactor 内部用pandas.concat拼接本地的 X 与 y 分区并构造xgb.DMatrix见 xgboost_ray.py。若训练集本身也出现在evals中会通过add_as_eval_method一并登记其余验证集则通过add_eval_data分发。步骤 6远程执行本地训练并通过 Rabit 同步远程调用 actor 的train方法见 xgboost_ray.py在本地参数中强制写入nthread进入RabitContextutils.py后调用原生xgb.train通过RabitContextManagerutils.py启动xgb.RabitTracker各 actor 以DMLC_*环境变量连接 tracker实现训练状态梯度、树模型等的 all-reduce 同步。步骤 7汇总结果_train中通过RayWrapper.materialize(fut[0])取回第一个 actor 的结果字典{booster: xgb.Booster, history: dict}。由于 Rabit 保证各 actor 训练出完全一致的模型只需返回任意一个结果即可见 xgboost_ray.py。另外有两个值得注意的边界处理若num_actors超过 X 的分区数会被截断为分区数若evals中某个验证集分区数少于 actor 数num_actors会被强制下调并发出警告见 xgboost_ray.py。预测流程3 个步骤数据以DMatrix对象传入_predict见 xgboost_ray.py通过_get_num_columns用随机种子固定的小样本探测 booster 输出列数推断结果列名随后将 booster 与列名放入 Ray 对象存储并对每个数据分区远程调用_map_predict在 worker 上执行booster.predict(xgb.DMatrix(part))得到局部预测见 xgboost_ray.py利用DMatrix中记录的 index、columns 与 row_lengths 元数据通过modin.distributed.dataframe.pandas.from_partitions将各分区的预测结果重建为modin.pandas.DataFrame返回。内部 API 一览除上述核心流程外xgboost_ray.py还提供了若干供内部使用的函数与类均可通过modin.experimental.xgboost.xgboost_ray访问名称职责ModinXGBoostActorRay actor 类持有set_train_data、add_eval_data、train三个远程方法_assign_row_partitions_to_actors按 IP 将行分区分配给 actor支持 y 与 X 的顺序对齐_split_data_across_actors封装 X/y 的分区分配并触发远程数据加载_get_num_actors计算/校验 actor 数量默认 2 CPU 对应 1 actor_train/_predictRay 引擎上的分布式训练与预测入口_map_predict单个分区上的远程局部预测_get_cluster_cpus/_get_min_cpus_per_node/_get_cpus_per_actor集群资源探测与 actor CPU 配额计算正确性与一致性验证测试覆盖仓库中的 test_xgboost.py 以Modin XGBoost 与原生 XGBoost 结果对齐为核心验证手段二分类breast_cancer、多分类iris/digits/wine、回归diabetes任务中将同样的数据分别喂给原生xgboost.train与xgb.train对比evals_result中每轮评估指标logloss/error、mlogloss、rmse要求数值在容差内一致对比两者在相同数据上的预测结果与 accuracy/MSE 指标num_actors参数在[1, cpu 总数, None, NPartitions1]四种取值下分别测试test_invalid_input验证了类型约束DMatrix拒绝非 DataFrame 输入、train拒绝非 DMatrix 输入、predict拒绝非 DMatrix 输入test_default.py 则验证非 Ray 引擎Python/Dask 等下调用会抛出ValueError。这些测试从侧面印证了文档描述的接口约束与引擎限制也是理解模块行为边界的最佳参考。注意事项与限制引擎限制仅支持 Ray 引擎Engine.get()非Ray时直接报错输入限制DMatrix仅接受modin.pandas.DataFramelabel 可为 DataFrame/Seriesobject类型列会抛错weight、base_margin、nthread、group、qid等原生 DMatrix 参数暂不支持actor 数量约束显式传入的num_actors必须是 Ray 集群节点数的整数倍数据分区数不足时会自动下调实验性按 modin_xgboost.rst 的说明该功能处于 experimental 阶段接口与行为可能在未来版本中变化。综上modin.experimental.xgboost以最小化 API 差异 深度复用 Modin 分区与 Ray 资源的方式把 XGBoost 从单机扩展到了 Ray 集群数据解包、按 IP 本地化切分、actor 并行训练、Rabit 状态同步、分区级并行预测与结果重建构成了一个完整且自洽的分布式机器学习方案。【免费下载链接】modinModin: Scale your Pandas workflows by changing a single line of code项目地址: https://gitcode.com/gh_mirrors/mo/modin创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表