XGBoost4J-Spark-GPU 实践指南:在 Apache Spark 集群上端到端 GPU 加速分布式 XGBoost 训练
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-sparkCPU 版之上额外引入 cuDF 列式数据ai.rapids.cudf.Table作为 executor 侧的数据交换格式。从源码结构看该能力以Spark 插件形式实现插件类为 GpuXGBoostPlugin实现了 XGBoostPlugin trait提供isEnabled、buildRddWatches、transform三个扩展点通过 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它依赖xgboost4j、xgboost4j-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 会强制要求训练参数devicecuda否则会抛出 “Using Spark-Rapids to accelerate XGBoost must set devicecuda” 错误——这就是后文训练参数中必须写device - cuda的底层原因。二、添加 XGBoost 依赖到项目在开始教程代码之前先参考 JVM 包安装说明install_jvm_packages小节把 XGBoost4J-Spark-GPU 作为项目依赖加入。官方同时提供**稳定版release与快照版snapshot**两种 Maven 构件供选择。由 pom.xml 可确认构件坐标groupIdml.dmlc/groupId artifactIdxgboost4j-spark-gpu_2.12/artifactId version3.5.0-SNAPSHOT/version其中_2.12后缀对应 Scala 2.12 二进制版本。构建产物还会通过maven-shade-plugin把xgboost4j与xgboost4j-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.enabledtrue即为此而设见第六节。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 XGBoostClassifierXGBoost4J-Spark-GPU 支持**回归、分类和排序ranking**三类模型。本教程以 Iris 演示多分类问题回归与排序的用法与分类非常相似对应XGBoostRegressor、XGBoostRanker。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 - 100boosting 迭代轮数device - cuda告知 XGBoost 使用 CUDA 设备而非 CPU。与单机模式不同Spark 分布式模式下 GPU 由 Spark 资源管理器分配而不是 XGBoost 自己管理因此不支持cuda:1这类显式指定设备序号的写法——executor 实际使用哪张卡由 Spark 按 GPU 资源调度决定。num_workers - 1XGBoost 分布式训练的 worker 数。训练参数的完整清单可参考 参数文档。与 XGBoost4J-Spark 包一致除默认的下划线命名参数外XGBoost4J-Spark-GPU 也支持这些参数的camel-case 变体以对齐 Spark MLlib 的命名习惯。例如设置max_depth可以像上面一样放进Map也可以通过 setterval 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即可构建转换器XGBoostClassificationModel。fit本质上就是训练过程产出的模型可用于预测等后续任务val xgbClassificationModel xgbClassifier.fit(train)源码层面fit触发的 GPU 数据管道值得展开对应 GpuXGBoostPlugin.buildRddWatchespreprocess选出 label/weight/baseMargin/group/特征列并做必要类型转换然后按需重分区repartitionIfNeeded、按需排序sortPartitionIfNeeded排序场景需要组内顺序列式化通过ColumnarRdd(train.toDF())把 DataFrame 转成 cuDFTable迭代器构建 QuantileDMatrix在 executor 侧把 cuDF Table 包装为CudfColumnBatch含特征/标签/权重/margin/group 五组列索引再交给 QuantileDMatrix 做 GPU 直方图分箱。若配置了评估集验证集 DMatrix 会以训练集 DMatrix 为ref创建保证训练与验证使用同一套分箱边界外部内存若启用外部内存则走 ExtMemQuantileDMatrix 路径把中间页缓存到spark.local.dir指向的本地盘默认/tmp。这也解释了 QuantileDMatrix 上setLabel/setWeight/setBaseMargin/setQueryId等 setter 一律抛XGBoostError——分箱矩阵的数据在构造时即已给定不再支持事后追加元数据。相关行为有专门的测试套件验证见 GpuXGBoostPluginSuite 与 GpuTestSuite。五、预测GPU 上的 transform得到XGBoostClassificationModel或XGBoostRegressionModel、XGBoostRankerModel后模型以 DataFrame 为输入读取特征列、逐行预测默认输出一个新的 DataFrameXGBoostClassificationModel输出 marginrawPredictionCol、每个类别的概率probabilityCol以及最终预测标签predictionColXGBoostRegressionModel输出预测标签predictionColXGBoostRankerModel输出预测标签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批量读取包装成CudfColumnBatch→DMatrix后调用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.amount1/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.amount与spark.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_version24.08.0 xgboost_version$LATEST_VERSION main_classIris app_jariris-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.cores12 \ --conf spark.task.cpus1 \ --conf spark.executor.resource.gpu.amount1 \ --conf spark.task.resource.gpu.amount0.08 \ --conf spark.rapids.sql.csv.read.double.enabledtrue \ --conf spark.rapids.sql.hasNansfalse \ --conf spark.pluginscom.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.amount1每个 executor 独占 1 张 GPUspark.task.resource.gpu.amount0.08task 级 GPU 分数12 核 executor 下 1/12≈0.08使 ETL task 并行而不独占 GPUspark.task.cpus1每个 task 1 核spark.rapids.sql.csv.read.double.enabledtrue允许 RAPIDS 用 GPU 读双精度 CSV 列Iris 特征均为 doublespark.rapids.sql.hasNansfalse告知 RAPIDS 数据无 NaN避免额外的空值检查开销spark.pluginscom.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 poolRMM 支持自 3.0 版本加入。当前仓库 pom.xml 的版本正是 3.5.0-SNAPSHOT与该弃用说明一致。当 XGBoost 以 RMM 插件编译构建方式见 构建文档时XGBoost Spark 包可以根据spark.rapids.memory.gpu.pooling.enabled与spark.rapids.memory.gpu.pool自动复用 RMM 内存池两个提交参数需同时设置。此外XGBoost 使用NCCL做 GPU 间通信NCCL 需要一部分显存作为通信缓冲区因此不要让 RMM 占满全部可用显存。内存池相关配置示例spark-submit \ --master $master \ --conf spark.rapids.memory.gpu.allocFraction0.5 \ --conf spark.rapids.memory.gpu.maxAllocFraction0.8 \ --conf spark.rapids.memory.gpu.poolARENA \ --conf spark.rapids.memory.gpu.pooling.enabledtrue \ ...插件侧的内存池联动逻辑可以在 GpuXGBoostPlugin.buildRddWatches 尾部 看到它读取spark.rapids.memory.gpu.pool默认async据此为底层训练追加参数——值为async追加use_cuda_async_pooltrueCUDA 异步内存池当前推荐路径值为none不附加任何内存池参数其他值如ARENA追加use_rmmtrue走 RMM 池。这解释了上文spark.rapids.memory.gpu.poolARENA与“3.5.0 起改用 async pool”两者在实现上的对应关系切换池策略不需要改代码只改 Spark 配置即可。八、小结本文围绕 XGBoost4J-Spark-GPU 官方教程 完整走通了一条 GPU 加速的分布式训练链路依赖Maven 引入ml.dmlc:xgboost4j-spark-gpu_2.12Spark/RAPIDS 依赖运行时由--packages提供数据Spark 内置 reader 读入 CSV → 显式 schema → 窗口函数dense_rank()把字符串标签编码为 0/1/2规避未被 RAPIDS 加速的StringIndexer训练XGBoostClassifierdevicecudaGPU 由 Spark 调度不支持cuda:Nfit(train)触发“列式化 → QuantileDMatrix 分箱 → 分布式 boosting”管道预测transform(test)输出rawPrediction/probability/prediction三列executor 侧按 cuDF 批次做 GPU 推理提交spark-submit挂载 RAPIDS 插件并配置 GPU 资源2.1.0 下用分数型spark.task.resource.gpu.amount获得 ETL 阶段的多 task 并行内存RMM 池已弃用3.5.0默认走 CUDA async poolspark.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),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →