NVIDIA cuML 分布式机器学习模块cuml.dask全解析:多节点多 GPU 算法 API 指南
【免费下载链接】cumlNVIDIA cuML: GPU-Accelerated Machine Learning项目地址: https://gitcode.com/GitHub_Trending/cu/cuml
cuml.dask是 NVIDIA cuML 面向多节点多 GPU(MNMG)场景的分布式算法子包,基于 Dask 生态把 cuDF、CuPy 与 RAFT 通信层串联起来,让 KMeans、PCA、随机森林、线性模型等经典算法可以跨 GPU 集群并行训练与推理。本文以仓库 API 文档页 docs/source/api/cuml.dask.rst 为骨架,结合python/cuml/cuml/dask/下各模块源码,完整梳理cuml.dask的模块结构、每个公开类的核心参数、基类设计原理与一套可落地的分布式训练工作流,帮助读者快速上手 cuML 的多 GPU 机器学习开发。
一、cuml.dask是什么
cuml.dask是 cuML 中"使用 Dask 实现的多节点、多 GPU 算法"的统一入口,其官方定位在 API 文档页中一句话概括为"Multi-node, multi-GPU algorithms using Dask."它并非一套独立的算法实现,而是由以下几层拼装而成:
- 数据层:以
dask_cudf.DataFrame或 CuPy 后端的dask.array承载分布式数据,每个 Dask worker 持有若干分区; - 调度层:借助
dask.distributed的Client与 Future 机制在集群上分发任务; - 通信层:通过 raft_dask 的
Comms在 worker 之间建立 NCCL 通信组,供 MNMG 算法做全局归约(如质心同步、特征协方差聚合); - 算法层:每个 worker 上实际调用的是 cuML 的
*_mg(multi-GPU)单机实现,例如cuml.cluster.kmeans_mg.KMeansMG、cuml.linear_model.linear_regression_mg.LinearRegressionMG。
从仓库结构看,python/cuml/cuml/dask/init.py 在包导入时即完成三件事:检查 Dask 依赖是否齐全、注册全部 12 个子模块、并强制关闭 Dask 的 p2p shuffle(dask.config.set({"dataframe.shuffle.method": "tasks"})),以确保分布式 DataFrame 混洗行为在 cuML 场景下稳定可控。
二、安装与依赖:缺什么会怎样
cuml.dask是 cuML 的可选扩展,不是cuml本体。其依赖检查逻辑位于 python/cuml/cuml/dask/init.py:导入时会尝试加载dask、dask.distributed、dask_cudf与raft_dask,任缺其一即抛出ModuleNotFoundError,并给出修复指引。
# 包导入时的依赖检查(源码节选) try: import dask import dask.distributed as _ # noqa import dask_cudf as _ # noqa import raft_dask as _ # noqa except ModuleNotFoundError as exc: raise ModuleNotFoundError( "Not all requirements for using `cuml.dask` are installed.\n\n" "# Install with Conda:\n" " conda install rapids-dask-dependency dask-cudf raft-dask\n\n" "# Or install with pip:\n" " pip install cuml-{cu_version}[dask]" )对应的两种安装方式为:
# Conda 方式 conda install rapids-dask-dependency dask-cudf raft-dask # pip 方式({ver} 为当前 CUDA 版本对应的标签,如 cu12) pip install cuml-{ver}[dask]除了 Python 侧依赖,部分 MNMG 算法(如 KMeans、PCA)还依赖 cuML C++ 层的多 GPU 编译产物。相关代码通过mnmg_import装饰器(见 python/cuml/cuml/dask/common/base.py)包裹*_mg模块的惰性导入:若当前构建未启用多 GPU 支持,会抛出带--multigpu构建提示的RuntimeError。
三、API 全景:模块与公开类总览
API 文档页 docs/source/api/cuml.dask.rst 按功能域划分了 12 个命名空间,与 python/cuml/cuml/dask/ 目录一一对应。完整映射如下:
| 功能域 | 命名空间 | 公开 API |
|---|---|---|
| Cluster | cuml.dask.cluster | DBSCAN、KMeans |
| Decomposition | cuml.dask.decomposition | PCA、TruncatedSVD |
| Ensemble | cuml.dask.ensemble | RandomForestClassifier、RandomForestRegressor |
| Linear Models | cuml.dask.linear_model | LinearRegression、Ridge、Lasso、ElasticNet |
| Manifold | cuml.dask.manifold | UMAP |
| Naive Bayes | cuml.dask.naive_bayes | MultinomialNB |
| Neighbors | cuml.dask.neighbors | NearestNeighbors、KNeighborsClassifier、KNeighborsRegressor |
| Preprocessing | cuml.dask.preprocessing | LabelBinarizer、OneHotEncoder |
| Feature Extraction | cuml.dask.feature_extraction.text | TfidfTransformer |
| Datasets | cuml.dask.datasets | make_blobs、make_classification、make_regression |
| Solvers | cuml.dask.solvers | CD |
| Base Classes and Mixins | cuml.dask.common.base | BaseEstimator、DelayedParallelFunc、DelayedPredictionMixin、DelayedTransformMixin、DelayedInverseTransformMixin |
注意:cuml.dask下还有metrics(如confusion_matrix)、common(分布式数据/工具类)等目录,但文档页列出的 12 个命名空间是公开 API 的主要入口。下面按域逐一展开。
四、集群:DBSCAN 与 KMeans
4.1 KMeans:多轮通信的最小化数据搬运
cuml.dask.cluster.KMeans(实现见 python/cuml/cuml/dask/cluster/kmeans.py)是文档中最典型的 MNMG 算法。源码 docstring 明确指出其设计目标:每个迭代只在 worker 之间共享质心,从而把数据搬运量降到最低;而 predict 阶段则退化为"纯并行"(embarrassingly parallel),直接调用单 GPU KMeans。
核心参数(与单 GPU 版本对齐,默认值取自源码):
| 参数 | 默认值 | 说明 |
|---|---|---|
n_clusters | 8 | 质心(簇)数量,必须是正整数 |
max_iter | 300 | EM 迭代上限,越大越准但越慢 |
tol | 1e-4 | 质心变化小于该阈值时提前收敛 |
init | 'scalable-k-means++' | 初始化策略:'scalable-k-means++'/'k-means||'/'random',或传入(n_clusters, n_features)的 ndarray 作为初始质心 |
oversampling_factor | 2 | scalable k-means++ 采样倍数,越大初始质心越好但内存开销越大,总采样数为oversampling_factor * n_clusters * 8 |
max_samples_per_batch | 32768 | 分批计算两两距离时的样本批大小,影响显存占用(每批约max_samples_per_batch * n_clusters个元素),n_clusters很大时可调小 |
random_state | None | 随机种子,保证可复现 |
verbose | False | 日志级别 |
fit的分布式流程(kmeans.py)可以拆成四步,值得作为理解 MNMG 的样板:
- 数据预处理:
DistributedDataHandler.create把 Dask 集合按 worker 组织成分区;支持sample_weight(会自动归一化); - 全局预检(preflight):
_validate_n_clusters在客户端校验n_clusters,_fetch_worker_sizes汇总各 worker 行数,若n_samples < n_clusters直接报错——把可预测的全局错误在提交分布式任务前暴露出来; - 建通信组:
Comms(comms_p2p=False)初始化 RAFT/NCCL 通信会话,随后向每个 worker 提交_func_preflight_fit(参数校验)与_func_fit(真实训练,内部实例化cuml.cluster.kmeans_mg.KMeansMG); - 结果汇总:质心经 NCCL 同步后各 worker 副本一致,因此只从第一个 worker 拉取完整模型(省内存);
inertia_由各 worker 的(inertia, n_samples)汇总求和;labels_则保留为分布式 Dask 集合(dask_cudf 或 dask.array),避免大数组回传客户端。
训练完成后可直接访问cluster_centers_、labels_、inertia_等属性,它们通过基类的属性代理机制从远端 Future 或本地模型中取回。
4.2 DBSCAN
cuml.dask.cluster.DBSCAN(实现见 python/cuml/cuml/dask/cluster/dbscan.py)是密度聚类算法的分布式版本,同样遵循"每 worker 本地计算 + 全局通信合并标签"的模式。与 KMeans 的连续迭代不同,DBSCAN 的核心开销在于跨分区的邻域查询与标签传播,适合样本量大、簇形状不规则的数据。
五、分解:PCA 与 TruncatedSVD
cuml.dask.decomposition.PCA(实现见 python/cuml/cuml/dask/decomposition/pca.py)的源码 docstring 明确了两点:
- 输入约定:MNMG PCA 期望 Dask cuDF 对象作为输入;
- 两种算法:
Full(默认)对数据做完整特征分解后取前 K 个特征向量;Jacobi通过迭代修正前 K 个特征向量,速度更快但精度可能略低。
其公开参数与单 GPU 版本一致,主要包括n_components(默认1,主成分个数)、whiten(是否白化)以及svd_solver相关选项;类继承自BaseDecomposition、DecompositionSyncFitMixin与DelayedTransformMixin/DelayedInverseTransformMixin,意味着它属于"同步拟合 + 延迟变换"型算法:fit时通过DecompositionSyncFitMixin在集群上做全局协方差聚合(见 python/cuml/cuml/dask/decomposition/base.py),transform/inverse_transform则按分区并行执行。
TruncatedSVD面向稀疏/截断场景,适合降维后仍需保留数据局部结构的任务,接口风格与 PCA 一致。
六、集成学习:随机森林分类与回归
cuml.dask.ensemble.RandomForestClassifier(实现见 python/cuml/cuml/dask/ensemble/randomforestclassifier.py)采用纯并行(embarrassingly-parallel)策略:设森林共N棵树、集群有w个 worker,则每个 worker 只在本地数据上构建N/w棵树,互不通信。源码还给出两条实用经验:
- 若每个 worker 只持有数据子集,通常要求数据预先充分洗牌,效果才接近全量训练;
- 若把全部数据复制到每个 worker(
fit收到w个内容相同的分区),结果将近似单 GPU 拟合。
核心参数(默认值取自 randomforestclassifier.py):
| 参数 | 默认值 | 说明 |
|---|---|---|
n_estimators | 100 | 森林总树数(不是每 worker 树数),会被均分到各 worker |
split_criterion | 0('gini') | 分裂准则:0/'gini'、1/'entropy'、2/'mse'、4/'poisson'、5/'gamma'、6/'inverse_gaussian'(分类任务仅 gini/entropy 有效) |
bootstrap | True | 是否对每棵树做有放回采样;False则每棵树用全量数据 |
max_samples | 1.0 | 每棵树使用的样本行比例 |
max_depth | None | 最大深度,None表示不设限(直到叶子纯净) |
max_leaves | -1 | 每棵树最大叶子数,-1表示不设限(软约束) |
max_features | 'auto' | 每个节点分裂时考虑的特征数策略 |
分类器额外支持predict_proba(通过DelayedPredictionProbaMixin),RandomForestRegressor实现位于 randomforestregressor.py,参数体系一致但准则以'mse'等回归型为主。
七、线性模型家族:LinearRegression、Ridge、Lasso、ElasticNet
cuml.dask.linear_model提供四件套,均继承自BaseEstimator、SyncFitMixinLinearModel与DelayedPredictionMixin。以 LinearRegression 为例:
- 拟合算法:
algorithm='eig'时基于协方差矩阵的特征分解,速度快,适合"高瘦"(tall and skinny)数据;文档同时提示 SVD 更慢但数值稳定性有保证,特征数增大时 eig 的精度可能下降; - 截距:
fit_intercept=True(默认)会额外拟合常数项c,把模型表达为y = x·β + c;置False要求数据已中心化; - 属性:拟合后提供
coef_(n_features长度的 cuDF Series)与intercept_。
Ridge(ridge.py)加入 L2 正则(alpha参数控制强度);Lasso(lasso.py)与ElasticNet(elastic_net.py)则通过坐标下降求解稀疏解。这些线性模型的fit走SyncFitMixinLinearModel._fit(base.py):先在集群上初始化Comms,把每个 worker 的分区大小通过parts_to_ranks广播给各 rank,再提交model_func创建LinearRegressionMG等 worker 端模型,最后统一等待 Future 完成。
八、流形学习:UMAP
cuml.dask.manifold.UMAP(实现见 python/cuml/cuml/dask/manifold/umap.py)是单 GPU UMAP 的分布式版本,面向超大样本量的降维与可视化。其核心流程是先构建分布式 kNN 图(复用cuml.dask.neighbors的索引能力),再在图上做嵌入优化;主要参数与单 GPU 版本一致,包括n_neighbors、n_components、min_dist、spread、learning_rate、random_state等,其中 kNN 图构建阶段的分布式查询由NearestNeighbors支撑。
九、朴素贝叶斯:MultinomialNB
cuml.dask.naive_bayes.MultinomialNB(实现见 python/cuml/cuml/dask/naive_bayes/naive_bayes.py)提供多项式朴素贝叶斯分类器的分布式实现,适合文本分类等计数特征场景。核心参数包括平滑系数alpha(默认1.0,即 Laplace 平滑)与fit_prior(是否学习类别先验)。fit阶段在各 worker 上统计类别计数与特征计数,再跨 worker 归约得到全局先验与条件概率;predict/predict_proba走延迟并行通道。
十、近邻:NearestNeighbors 与 KNN 分类/回归
cuml.dask.neighbors包含三个类:
NearestNeighbors(nearest_neighbors.py):在分布式数据上构建近邻索引。关键参数n_neighbors(默认5)与batch_size(默认2_000_000)——batch_size决定每次查询处理的最大行数,直接影响吞吐,且每个承载索引分区的 worker 需要约batch_size * n_features * 4字节的额外显存,需按显存预算调整;KNeighborsClassifier(kneighbors_classifier.py):在分布式索引上做 kNN 投票分类,支持predict与predict_proba;KNeighborsRegressor(kneighbors_regressor.py):kNN 回归,预测值为邻居标签的加权/平均聚合。
三者共享BaseEstimator与DistributedDataHandler基础设施,fit建立分布式索引,查询阶段按分区并行。
十一、预处理与特征提取
cuml.dask.preprocessing提供:
LabelBinarizer(preprocessing/_label.py):把类别标签二值化为 one-hot 风格的分布式表示;OneHotEncoder(preprocessing/encoders.py):类别特征独热编码,支持dtype、handle_unknown、sparse_output等参数,transform/inverse_transform均走DelayedTransformMixin/DelayedInverseTransformMixin的延迟并行通道。
cuml.dask.feature_extraction.text.TfidfTransformer(feature_extraction/text/)把词频矩阵转换为 TF-IDF 权重,用于文本向量化流水线,接口与 scikit-learn 对应组件对齐。
十二、分布式数据集生成:make_blobs 等
cuml.dask.datasets的三个生成器是快速验证 MNMG 工作流的最佳工具,它们的工作方式一致:在每个 Dask worker 上调用单 GPU 版本的生成器,再把结果拼接成分布式集合。
以make_blobs(datasets/blobs.py)为例,签名与关键参数:
from cuml.dask.datasets import make_blobs X, y = make_blobs( n_samples=100, # 总行数(跨 worker 分配) n_features=2, # 特征数 centers=None, # 簇数或固定质心数组;None 时生成 3 个中心 cluster_std=1.0, # 簇内点标准差 n_parts=None, # 分区数,可大于 worker 数;None 时等于 worker 数 center_box=(-10, 10), # 质心所在边界框 shuffle=True, # 是否在 worker 内打乱样本 random_state=None, # 随机种子 return_centers=False, # 是否同时返回质心 dtype="float32", # 数据类型 client=None, # 显式传入 Dask Client workers=None, # 限定使用的 worker 地址列表;None 表示全部 )返回值为 CuPy 后端的dask.array(X形状(n_samples, n_features),y形状(n_samples,)),return_centers=True时额外返回质心数组。源码把n_samples均分到n_parts个分区并分派到各 worker,用独立的随机种子保证可复现。make_classification与make_regression位于 classification.py 和 regression.py,接口风格一致。
十三、求解器:CD(坐标下降)
cuml.dask.solvers.CD(python/cuml/cuml/dask/solvers/)是坐标下降(Coordinate Descent)求解器的分布式版本,对应单 GPU 的cuml.solvers.CD,是 Lasso/ElasticNet 内部使用的优化器,也可独立使用。核心参数包括alpha(正则强度)、max_iter、tol、selection(坐标选择策略)等。由于坐标下降天然适合行分区并行——各 worker 在本地分区上计算梯度/残差后做全局归约——它是理解分布式优化器的良好示例。
十四、基类与 Mixin:分布式估算器的骨架
文档页的最后一部分列出了cuml.dask.common.base中的基础设施类,它们支撑了上述所有算法,值得单独解读(源码见 python/cuml/cuml/dask/common/base.py)。
14.1BaseEstimator:一切分布式估算器的根
构造函数__init__(*, client=None, verbose=False, **kwargs):
client:显式传入dask.distributed.Client;不传时通过get_client()从全局获取或自动创建;verbose:既保存在本地,也会注入kwargs["verbose"]一并传给 worker 端模型;- 其余超参数收进
self.kwargs,在提交分布式任务时透传给*_mg模型。
核心机制是internal_model:它是一个dask.distributed.Future[cuml.Base]、本地cuml.Base实例或None三态对象,统一通过_set_internal_model管理。_check_internal_model会做类型校验(Future 的类型必须是cuml.Base子类),get_combined_model()则把集群上训练的模型收敛为一个可序列化的单 GPU 模型——这对pickle保存/加载分布式模型至关重要(__getstate__/__setstate__已实现该逻辑)。
此外,BaseEstimator.__getattr__(base.py)实现了属性代理:用户访问model.cluster_centers_时,如果本地没有该属性,会优先查找_cluster_centers_,再递归代理到internal_model——若模型在远端,则通过dask.delayed惰性取回。这就是为什么分布式模型用起来和单 GPU 模型一样直观。
14.2DelayedParallelFunc与四个 Mixin
DelayedParallelFunc._run_parallel_func(base.py)把predict/transform/score这类纯并行操作抽象成统一通道:
- 把
X切成X.to_delayed()分区列表; - 把内部模型与目标函数分别包装为
dask.delayed(模型带pure=True以便跨任务复用); - 对每个分区提交
func(model_delayed, part, **kwargs); - 按
output_collection_type('cupy'或'cudf',缺省沿用 fit 时的输入类型)把结果拼接为dask.array或dask_cudf.DataFrame,并支持delayed(惰性)与output_futures(直接返回 Future)两种返回形态。
基于它派生的四个 Mixin 分别对应:DelayedPredictionMixin(_predict,封装model.predict)、DelayedPredictionProbaMixin(_predict_proba,封装model.predict_proba)、DelayedTransformMixin(_transform,封装model.transform)、DelayedInverseTransformMixin(_inverse_transform,封装model.inverse_transform)。几乎所有文档列出的算法类,都是"BaseEstimator+ 若干 Mixin + 算法特有 fit 逻辑"的组合。
十五、从入门到实战:一个完整的 MNMG 工作流
综合上文,一个典型的多 GPU KMeans 训练流程如下(代码模式取自 make_blobs 与 KMeans 的 docstring 示例):
from dask_cuda import LocalCUDACluster from dask.distributed import Client from cuml.dask.cluster import KMeans from cuml.dask.datasets import make_blobs # 1. 启动本地 GPU 集群(每个 worker 一个线程) cluster = LocalCUDACluster(threads_per_worker=1) client = Client(cluster) # 2. 生成跨 worker 分布的示例数据(1000 行、10 特征、42 个簇) workers = list(client.scheduler_info()["workers"].keys()) X, y = make_blobs( n_samples=1000, n_features=10, centers=42, cluster_std=0.1, workers=workers, ) # 3. 分布式训练 model = KMeans(n_clusters=42, random_state=42) model.fit(X) # 4. 读取结果(cluster_centers_ 经属性代理取回;labels_ 保持分布式) print(model.cluster_centers_) print(model.inertia_) labels = model.labels_ # dask_cudf / dask.array labels.compute() # 需要时再物化 # 5. 清理 client.close() cluster.close()工作流要点回顾:
- 数据形态:
fit接受dask_cudf.DataFrame或 CuPy 后端的dask.array,内部由DistributedDataHandler(python/cuml/cuml/dask/common/input_utils.py)统一处理,因此生成器产出的分布式集合可直接喂给所有算法; - client 生命周期:估算器构造时
get_client()会解析client参数,建议显式传入并统一管理; - 惰性与物化:
predict/transform默认返回惰性 Dask 集合(delayed=True),调用.compute()/.persist()才真正执行,便于构建流水线; - 可复现性:各算法均支持
random_state,KMeans 内部还会用check_random_seed做归一化,确保集群上每次调用种子一致(见 kmeans.py)。
十六、测试与验证:如何确认分布式行为
仓库为cuml.dask提供了完备的测试套件,位于 python/cuml/tests/dask/(如test_dask_kmeans.py、test_dask_linear_regression.py、test_dask_datasets.py、test_dask_dbscan.py、test_dask_kneighbors_classifier.py、test_dask_base.py等)。这些测试覆盖了:分布式输入处理、基类属性代理、各算法与单 GPU 结果的一致性对比、数据集生成器的形状与可复现性等。在 CI 中它们由 ci/run_cuml_dask_pytests.sh 驱动;本地验证分布式功能可参考该脚本的 pytest 入口方式执行对应测试文件。
对开发者而言,这些测试既是回归保障,也是"分布式估算器应该如何被使用"的权威示例——例如test_dask_base.py演示了BaseEstimator的内部模型管理,test_dask_datasets.py验证了n_parts与 worker 的映射关系,阅读它们能比文档更快理解每个类的契约。
结语
cuml.dask把 cuML 的单 GPU 算法库扩展到了多节点多 GPU 维度:数据分布在 Dask 集合上,通信由 RAFTComms承接,算法主体是各 worker 上的*_mg实现,预测/变换则统一走DelayedParallelFunc的纯并行通道。本文依据 API 文档页 docs/source/api/cuml.dask.rst 的分类框架,逐一解析了聚类、分解、集成、线性模型、流形、朴素贝叶斯、近邻、预处理、特征提取、数据集生成、求解器以及基类 Mixin 的接口与实现要点。上手时只需三步:装齐dask-cudf/raft-dask依赖、用LocalCUDACluster起集群、把单 GPU 的 fit/predict 换成cuml.dask的对应类——其余交给 Dask 与 NCCL 编排。
【免费下载链接】cumlNVIDIA cuML: GPU-Accelerated Machine Learning项目地址: https://gitcode.com/GitHub_Trending/cu/cuml
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考