XGBoost4J-Spark-GPU 实践指南:在 Apache Spark 集群上端到端 GPU 加速分布式 XGBoost 训练
2026/9/19 12:49:47 网站建设 项目流程

XGBoost4J-Spark-GPU 实践指南:在 Apache Spark 集群上端到端 GPU 加速分布式 XGBoost 训练

【免费下载链接】xgboostScalable, Portable and Distributed Gradient Boosting (GBDT, GBRT or GBM) Library, for Python, R, Java, Scala, C and more. Runs on single machine, Hadoop, Spark, Dask, Flink and DataFlow项目地址: https://gitcode.com/gh_mirrors/xg/xgboost

本文以官方教程 doc/jvm/xgboost4j_spark_gpu_tutorial.rst 为主体,讲解如何使用 XGBoost4J-Spark-GPU 结合 RAPIDS Accelerator for Apache Spark,在 Spark 集群上完成“数据加载 → 特征/标签预处理 → GPU 训练 → GPU 推理”的完整流程,并覆盖spark-submit提交配置、stage 级 GPU 资源调度以及 RMM 内存池等进阶配置。读完本文,你可以独立搭建并跑通一个 GPU 加速的分布式 XGBoost 机器学习应用。

一、XGBoost4J-Spark-GPU 是什么:定位与整体架构

XGBoost4J-Spark-GPU是一个开源库,目标是借助 RAPIDS Accelerator for Apache Spark 产品,把 Apache Spark 集群上分布式 XGBoost 训练**从数据准备到模型训练(end to end)**整体迁移到 GPU 上加速。它建立在xgboost4j-spark(CPU 版)之上,额外引入 cuDF 列式数据(ai.rapids.cudf.Table)作为 executor 侧的数据交换格式。

从源码结构看,该能力以Spark 插件形式实现:

  • 插件类为 GpuXGBoostPlugin,实现了 XGBoostPlugin trait(提供isEnabledbuildRddWatchestransform三个扩展点);
  • 通过 Java SPI 文件 META-INF/services/ml.dmlc.xgboost4j.scala.spark.XGBoostPlugin 注册,XGBoostEstimator.fit时会被自动发现并接管数据管道;
  • 模块的 Maven 坐标为ml.dmlc:xgboost4j-spark-gpu_2.12(当前仓库快照版本 3.5.0-SNAPSHOT),见 xgboost4j-spark-gpu/pom.xml,它依赖xgboost4jxgboost4j-spark,并以provided作用域声明 Spark 与rapids-4-spark依赖——这意味着运行时由 Spark 集群(通过--packages)提供这些依赖,而不是打进用户应用 fat jar。

插件的启用条件可以直接从 GpuXGBoostPlugin.isEnabled 读出:必须满足spark.plugins中包含com.nvidia.spark.SQLPlugin,且spark.rapids.sql.enabled为 true(默认 true)。不满足时回退到常规 CPU 管道。另外 validate 会强制要求训练参数device=cuda,否则会抛出 “Using Spark-Rapids to accelerate XGBoost must set device=cuda” 错误——这就是后文训练参数中必须写"device" -> "cuda"的底层原因。

二、添加 XGBoost 依赖到项目

在开始教程代码之前,先参考 JVM 包安装说明(install_jvm_packages小节)把 XGBoost4J-Spark-GPU 作为项目依赖加入。官方同时提供**稳定版(release)与快照版(snapshot)**两种 Maven 构件供选择。

由 pom.xml 可确认构件坐标:

<groupId>ml.dmlc</groupId> <artifactId>xgboost4j-spark-gpu_2.12</artifactId> <version>3.5.0-SNAPSHOT</version>

其中_2.12后缀对应 Scala 2.12 二进制版本。构建产物还会通过maven-shade-pluginxgboost4jxgboost4j-spark两个模块 shade 进来(见 pom.xml 构建配置),因此用户侧只需引用这一个 GPU 构件即可获得完整功能。

三、数据准备:用 Spark 把原始数据整形为 XGBoost 数据接口

本节沿用官方教程,以Iris 数据集为例演示如何用 Apache Spark 转换原始数据集、使其符合 XGBoost 的数据接口。

Iris 数据集以 CSV 格式提供。每条记录包含 4 个特征列:“sepal length”(花萼长)、“sepal width”(花萼宽)、“petal length”(花瓣长)、“petal width”(花瓣宽),外加一个 “class” 列,它是标签,取值为 “Iris Setosa”、“Iris Versicolour” 和 “Iris Virginica” 三类。

3.1 用 Spark 内置 CSV Reader 读取数据集

import org.apache.spark.sql.SparkSession import org.apache.spark.sql.types.{DoubleType, StringType, StructField, StructType} val spark = SparkSession.builder().getOrCreate() val labelName = "class" val schema = new StructType(Array( StructField("sepal length", DoubleType, true), StructField("sepal width", DoubleType, true), StructField("petal length", DoubleType, true), StructField("petal width", DoubleType, true), StructField(labelName, StringType, true))) val xgbInput = spark.read.option("header", "false") .schema(schema) .csv(dataPath)

要点解读:

  • SparkSession是所有基于 DataFrame 的 Spark 应用的统一入口;
  • schema变量显式定义了包裹 Iris 数据的 DataFrame schema。显式指定 schema 后,你可以自定义列名与类型;否则列名将退化为 Spark 默认推导的_col0_col1之类;
  • 最后用 Spark 内置 CSV reader 把 Iris CSV 文件读入 DataFramexgbInput。Spark 还内置了 ORC、Parquet、Avro、JSON 等格式的 reader,可按需替换。

注意:在 GPU 场景下,CSV 读取要落到 GPU 上,需要 RAPIDS Accelerator 支持该类型的 CSV 读取(提交参数中的spark.rapids.sql.csv.read.double.enabled=true即为此而设,见第六节)。

3.2 转换原始 Iris 数据:字符串标签编码为数值标签

为了让 XGBoost 认识 Iris 数据集,必须把 String 类型的标签列 “class” 编码为 Double 类型标签。

一种直接的方式是使用 Spark 内置的StringIndexer特征转换器,但该算子并未被 RAPIDS Accelerator 加速,使用它会回退到 CPU。因此官方教程采用等价的纯列操作方案:

import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ val spec = Window.orderBy(labelName) val Array(train, test) = xgbInput .withColumn("tmpClassName", dense_rank().over(spec) - 1) .drop(labelName) .withColumnRenamed("tmpClassName", labelName) .randomSplit(Array(0.7, 0.3), seed = 1) train.show(5)

输出示例:

+------------+-----------+------------+-----------+-----+ |sepal length|sepal width|petal length|petal width|class| +------------+-----------+------------+-----------+-----+ | 4.3| 3.0| 1.1| 0.1| 0| | 4.4| 2.9| 1.4| 0.2| 0| | 4.4| 3.0| 1.3| 0.2| 0| | 4.4| 3.2| 1.3| 0.2| 0| | 4.6| 3.2| 1.4| 0.2| 0| +------------+-----------+------------+-----------+-----+

原理说明:dense_rank().over(Windows.orderBy(labelName))按标签字典序给不同类别分配 1、2、3 的稠密排名,- 1后得到 0、1、2 的标签索引;随后丢弃原字符串标签列并把临时列重命名回labelName,同时用randomSplit(Array(0.7, 0.3), seed = 1)按 7:3 划分训练/测试集。整条链路(窗口函数、列重命名、randomSplit)都是 RAPIDS 可加速的 Spark SQL 操作,从而保证 ETL 阶段也能跑在 GPU 上。

四、训练:定义并 fit 一个 GPU XGBoostClassifier

XGBoost4J-Spark-GPU 支持**回归、分类和排序(ranking)**三类模型。本教程以 Iris 演示多分类问题;回归与排序的用法与分类非常相似(对应XGBoostRegressorXGBoostRanker)。

4.1 构造 XGBoostClassifier 与关键参数

import ml.dmlc.xgboost4j.scala.spark.XGBoostClassifier val xgbParam = Map( "objective" -> "multi:softprob", "num_class" -> 3, "num_round" -> 100, "device" -> "cuda", "num_workers" -> 1) val featuresNames = schema.fieldNames.filter(name => name != labelName) val xgbClassifier = new XGBoostClassifier(xgbParam) .setFeaturesCol(featuresNames) .setLabelCol(labelName)

参数说明:

  • "objective" -> "multi:softprob":多分类 softprob 目标,输出每个类别的概率;
  • "num_class" -> 3:类别数,须与标签编码一致;
  • "num_round" -> 100:boosting 迭代轮数;
  • "device" -> "cuda":告知 XGBoost 使用 CUDA 设备而非 CPU。与单机模式不同,Spark 分布式模式下 GPU 由 Spark 资源管理器分配,而不是 XGBoost 自己管理,因此不支持cuda:1这类显式指定设备序号的写法——executor 实际使用哪张卡由 Spark 按 GPU 资源调度决定。
  • "num_workers" -> 1:XGBoost 分布式训练的 worker 数。

训练参数的完整清单可参考 参数文档。与 XGBoost4J-Spark 包一致,除默认的下划线命名参数外,XGBoost4J-Spark-GPU 也支持这些参数的camel-case 变体,以对齐 Spark MLlib 的命名习惯。例如设置max_depth可以像上面一样放进Map,也可以通过 setter:

val xgbClassifier = new XGBoostClassifier(xgbParam) .setFeaturesCol(featuresNames) .setLabelCol(labelName) xgbClassifier.setMaxDepth(2)

与 CPU 版 XGBoost4J-Spark 的一个重要差异:CPU 版既接受VectorUDT类型的单列特征,也接受特征列名数组;而XGBoost4J-Spark-GPU 只接受特征列名数组,即setFeaturesCol(value: Array[String])。这与源码一致——GpuXGBoostPlugin.preprocess 只遍历getFeaturesCols(列名数组)做类型转换与列选择。

4.2 fit:训练过程

设置好参数与特征/标签列之后,用输入 DataFrame 调用fit即可构建转换器XGBoostClassificationModelfit本质上就是训练过程,产出的模型可用于预测等后续任务:

val xgbClassificationModel = xgbClassifier.fit(train)

源码层面,fit触发的 GPU 数据管道值得展开(对应 GpuXGBoostPlugin.buildRddWatches):

  1. preprocess:选出 label/weight/baseMargin/group/特征列并做必要类型转换,然后按需重分区(repartitionIfNeeded)、按需排序(sortPartitionIfNeeded,排序场景需要组内顺序);
  2. 列式化:通过ColumnarRdd(train.toDF())把 DataFrame 转成 cuDFTable迭代器;
  3. 构建 QuantileDMatrix:在 executor 侧把 cuDF Table 包装为CudfColumnBatch(含特征/标签/权重/margin/group 五组列索引),再交给 QuantileDMatrix 做 GPU 直方图分箱。若配置了评估集,验证集 DMatrix 会以训练集 DMatrix 为ref创建,保证训练与验证使用同一套分箱边界;
  4. 外部内存:若启用外部内存,则走 ExtMemQuantileDMatrix 路径,把中间页缓存到spark.local.dir指向的本地盘(默认/tmp)。

这也解释了 QuantileDMatrix 上setLabel/setWeight/setBaseMargin/setQueryId等 setter 一律抛XGBoostError——分箱矩阵的数据在构造时即已给定,不再支持事后追加元数据。相关行为有专门的测试套件验证,见 GpuXGBoostPluginSuite 与 GpuTestSuite。

五、预测:GPU 上的 transform

得到XGBoostClassificationModel(或XGBoostRegressionModelXGBoostRankerModel)后,模型以 DataFrame 为输入,读取特征列、逐行预测,默认输出一个新的 DataFrame:

  • XGBoostClassificationModel输出 margin(rawPredictionCol)、每个类别的概率(probabilityCol)以及最终预测标签(predictionCol);
  • XGBoostRegressionModel输出预测标签(predictionCol);
  • XGBoostRankerModel输出预测标签(predictionCol)。
val xgbClassificationModel = xgbClassifier.fit(train) val results = xgbClassificationModel.transform(test) results.show()

结果示例(截取自教程输出):

+------------+-----------+------------------+-------------------+-----+--------------------+--------------------+----------+ |sepal length|sepal width| petal length| petal width|class| rawPrediction| probability|prediction| +------------+-----------+------------------+-------------------+-----+--------------------+--------------------+----------+ | 4.5| 2.3| 1.3|0.30000000000000004| 0|[3.16666603088378...|[0.98853939771652...| 0.0| | 4.6| 3.1| 1.5| 0.2| 0|[3.25857257843017...|[0.98969423770904...| 0.0| | 4.9| 2.4| 3.3| 1.0| 1|[-2.1498908996582...|[0.00596602633595...| 1.0| | 5.7| 2.5| 5.0| 2.0| 2|[-2.1498908996582...|[0.00280966912396...| 2.0| +------------+-----------+------------------+-------------------+-----+--------------------+--------------------+----------+

从 GpuXGBoostPlugin.transform 的源码可以看到 GPU 推理的执行方式:

  • 模型 Booster 通过sc.broadcast(model.nativeBooster)广播到各 executor,同一 executor 内所有 Spark task 共享同一个 booster 实例;
  • 首次预测时,executor 依据 Spark 资源管理器分配的 GPU 地址(XGBoost.getGPUAddrFromResources)调用booster.setParam("device", "cuda:$gpuId"),即训练参数里不写cuda:1、而由 Spark 运行时决定具体设备序号的原因;
  • 每个分区按 cuDF 列式Table批量读取,包装成CudfColumnBatchDMatrix后调用predictInternal批量推理,再转回行式Row输出;
  • 通过TaskContext.addTaskCompletionListener确保每个 task 结束时关闭持有的 GPUColumnarBatch,避免显存泄漏。

六、提交应用:spark-submit 配置与 stage 级 GPU 调度

前提:你已经配置好支持 GPU 的 Spark standalone 集群(配置方式参见 NVIDIA Spark-RAPIDS 官方指南)。

自 XGBoost 2.1.0 起,stage 级调度(stage-level scheduling)被自动启用。因此如果你使用的是 Spark standalone 3.4.0 及以上版本,强烈建议把spark.task.resource.gpu.amount配置为分数值——这样 ETL 阶段可以有多个 task 并行运行。示例配置:"spark.task.resource.gpu.amount=1/spark.executor.cores"。反之,如果你使用的是早于 2.1.0 的 XGBoost 版本或低于 3.4.0 的 Spark standalone 集群,则仍需让spark.task.resource.gpu.amount等于spark.executor.resource.gpu.amount

该调度逻辑在源码中位于 XGBoost.scala 的 stage-level scheduling 管理段:它读取spark.executor.resource.gpu.amountspark.task.resource.gpu.amount两个配置,并结合 Spark 版本判断是否跳过 stage 级调度——当 task 级 GPU 量为分数时,ETL task 不会占用 GPU,只有训练 task 才真正申请 GPU,从而让 ETL 与训练解耦。

假设应用主类为Iris、应用 jar 为iris-1.0.0.jar,提交 XGBoost 应用到 Apache Spark Standalone 集群的示例如下:

rapids_version=24.08.0 xgboost_version=$LATEST_VERSION main_class=Iris app_jar=iris-1.0.0.jar spark-submit \ --master $master \ --packages com.nvidia:rapids-4-spark_2.12:${rapids_version},ml.dmlc:xgboost4j-spark-gpu_2.12:${xgboost_version} \ --conf spark.executor.cores=12 \ --conf spark.task.cpus=1 \ --conf spark.executor.resource.gpu.amount=1 \ --conf spark.task.resource.gpu.amount=0.08 \ --conf spark.rapids.sql.csv.read.double.enabled=true \ --conf spark.rapids.sql.hasNans=false \ --conf spark.plugins=com.nvidia.spark.SQLPlugin \ --class ${main_class} \ ${app_jar}

关键配置解读:

配置作用
--packages com.nvidia:rapids-4-spark_2.12:...,ml.dmlc:xgboost4j-spark-gpu_2.12:...一次性拉取 RAPIDS Accelerator 与 XGBoost GPU 包(对应 pom.xml 中provided依赖的运行时来源)
spark.executor.resource.gpu.amount=1每个 executor 独占 1 张 GPU
spark.task.resource.gpu.amount=0.08task 级 GPU 分数(12 核 executor 下 1/12≈0.08),使 ETL task 并行而不独占 GPU
spark.task.cpus=1每个 task 1 核
spark.rapids.sql.csv.read.double.enabled=true允许 RAPIDS 用 GPU 读双精度 CSV 列(Iris 特征均为 double)
spark.rapids.sql.hasNans=false告知 RAPIDS 数据无 NaN,避免额外的空值检查开销
spark.plugins=com.nvidia.spark.SQLPluginRAPIDS Accelerator 以 Spark 插件形式加载——这也是 GpuXGBoostPlugin.isEnabled 判定 GPU 管道生效的必要条件

RAPIDS Accelerator 的更多配置项与其 FAQ 可查阅 NVIDIA 官方文档(Spark-RAPIDS 配置页与 FAQ 页)。

七、RMM 内存池支持(3.5.0 起已弃用,改用 CUDA 异步内存池)

版本提示:RMM 插件自3.5.0起被弃用(deprecated),官方建议改用 CUDA async pool;RMM 支持自 3.0 版本加入。当前仓库 pom.xml 的版本正是 3.5.0-SNAPSHOT,与该弃用说明一致。

当 XGBoost 以 RMM 插件编译(构建方式见 构建文档)时,XGBoost Spark 包可以根据spark.rapids.memory.gpu.pooling.enabledspark.rapids.memory.gpu.pool自动复用 RMM 内存池,两个提交参数需同时设置。此外,XGBoost 使用NCCL做 GPU 间通信,NCCL 需要一部分显存作为通信缓冲区,因此不要让 RMM 占满全部可用显存。内存池相关配置示例:

spark-submit \ --master $master \ --conf spark.rapids.memory.gpu.allocFraction=0.5 \ --conf spark.rapids.memory.gpu.maxAllocFraction=0.8 \ --conf spark.rapids.memory.gpu.pool=ARENA \ --conf spark.rapids.memory.gpu.pooling.enabled=true \ ...

插件侧的内存池联动逻辑可以在 GpuXGBoostPlugin.buildRddWatches 尾部 看到:它读取spark.rapids.memory.gpu.pool(默认async),据此为底层训练追加参数——

  • 值为async:追加use_cuda_async_pool=true(CUDA 异步内存池,当前推荐路径);
  • 值为none:不附加任何内存池参数;
  • 其他值(如ARENA):追加use_rmm=true,走 RMM 池。

这解释了上文spark.rapids.memory.gpu.pool=ARENA与“3.5.0 起改用 async pool”两者在实现上的对应关系:切换池策略不需要改代码,只改 Spark 配置即可。

八、小结

本文围绕 XGBoost4J-Spark-GPU 官方教程 完整走通了一条 GPU 加速的分布式训练链路:

  1. 依赖:Maven 引入ml.dmlc:xgboost4j-spark-gpu_2.12,Spark/RAPIDS 依赖运行时由--packages提供;
  2. 数据:Spark 内置 reader 读入 CSV → 显式 schema → 窗口函数dense_rank()把字符串标签编码为 0/1/2,规避未被 RAPIDS 加速的StringIndexer
  3. 训练XGBoostClassifier+device=cuda(GPU 由 Spark 调度,不支持cuda:N),fit(train)触发“列式化 → QuantileDMatrix 分箱 → 分布式 boosting”管道;
  4. 预测transform(test)输出rawPrediction/probability/prediction三列,executor 侧按 cuDF 批次做 GPU 推理;
  5. 提交spark-submit挂载 RAPIDS 插件并配置 GPU 资源;2.1.0+ 下用分数型spark.task.resource.gpu.amount获得 ETL 阶段的多 task 并行;
  6. 内存:RMM 池已弃用(3.5.0),默认走 CUDA async pool,spark.rapids.memory.gpu.pool一键切换。

继续深入时可查阅的仓库材料:JVM 包总览与安装、XGBoost 参数参考、构建文档、CPU 版 Spark 教程、GPU 插件实现 及其测试套件。

【免费下载链接】xgboostScalable, Portable and Distributed Gradient Boosting (GBDT, GBRT or GBM) Library, for Python, R, Java, Scala, C and more. Runs on single machine, Hadoop, Spark, Dask, Flink and DataFlow项目地址: https://gitcode.com/gh_mirrors/xg/xgboost

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询