
人工智能机器学习深度学习图计算【免费下载链接】dglPython package built to ease deep learning on graph, on top of existing DL frameworks.项目地址https://gitcode.com/gh_mirrors/dg/dgl点击查看免费下载本指南围绕 DGL 仓库中 examples/multigpu/graphbolt 目录下的多 GPU 训练示例展开讲解如何在多张 GPU 上使用 GraphBolt 数据加载器DataLoader与 PyTorch 的分布式数据并行Distributed Data ParallelDDP训练 GraphSAGE 节点分类模型。读者学完后将掌握DistributedItemSampler的分片采样原理、DDP 环境初始化、Join上下文管理器处理不均衡输入、以及跨 rank 加权聚合评估指标等完整的多卡训练技术方案。快速运行该示例的入口脚本位于 examples/multigpu/graphbolt/node_classification.py运行方式如下python node_classification.py --gpu0,1--gpu参数接受逗号分隔的 GPU 编号列表例如--gpu0,1表示在 GPU 0 与 GPU 1 上并行训练脚本会通过torch.multiprocessing.spawn为每个 GPU 派生一个子进程每个子进程对应一个 DDP rankworld_size等于 GPU 数量训练结束后会在 rank 0 上打印验证集准确率每个 epoch与测试集准确率。前置知识阅读本示例前官方建议先熟悉两个基础内容单卡 GraphBolt 节点分类示例即 examples/graphbolt/node_classification.py它演示了使用gb.ItemSampler、sample_neighbor、fetch_feature、gb.DataLoader构建端到端 GraphBolt 训练流水线的方法多卡版本正是将其中的ItemSampler替换为DistributedItemSampler后的分布式扩展。经典 GraphSAGE 实现即 examples/core/graphsage/node_classification.py用于理解 GraphSAGE 模型的训练范式。总体执行流程脚本源码顶部的注释给出了完整的流程示意图可概括为两个阶段main │ ├─── OnDiskDataset 预处理gb.BuiltinDataset(args.dataset).load() │ └─── run (multiprocessing) │ ├─── 初始化进程组并构建分布式 SAGE 模型DDP │ ├─── train │ │ │ ├─── 使用 DistributedItemSampler 构建 GraphBolt dataloader │ │ │ └─── 训练循环SAGE.forward → 验证集评估 → 收集各 rank 的指标 │ └─── 测试集评估主进程先加载数据集然后通过mp.spawn启动world_size个子进程每个子进程独立执行run(rank, world_size, args, devices, dataset)。多 GPU 环境初始化在run函数开头完成分布式环境设置见 node_classification.pydevice devices[rank] torch.cuda.set_device(device) dist.init_process_group( backendnccl, # 分布式 GPU 训练使用 NCCL 后端 init_methodtcp://127.0.0.1:12345, world_sizeworld_size, rankrank, )要点说明backendnccl多 GPU 训练推荐使用 NCCL 后端它针对 NVIDIA GPU 通信做了深度优化init_method示例使用tcp://127.0.0.1:12345作为默认的进程组初始化方式仅适用于单机多卡场景若在跨机集群中使用需要替换为实际可用的通信地址在main中还会设置os.environ[OMP_NUM_THREADS] str(mp.cpu_count() // 2 // world_size)限制线程数以避免资源竞争并通过mp.set_sharing_strategy(file_system)设置多进程共享策略。初始化进程组后脚本会将数据集中的图结构与特征搬运到对应设备if args.storage_device pinned: graph dataset.graph.pin_memory_() feature dataset.feature.pin_memory_() else: graph dataset.graph.to(args.storage_device) feature dataset.feature.to(args.storage_device)storage_device pinned时调用pin_memory_()将图与特征固定在 CPU 的页锁定内存pinned memory中支持异步传输与重叠取数否则通过.to()将数据放到 CPU 或 CUDA 设备上具体取决于--mode参数。DistributedItemSampler多卡数据分片的核心这是多 GPU 版本区别于单卡版本的核心组件其完整实现位于 python/dgl/graphbolt/item_sampler.py。它与单卡ItemSampler的最大区别在于原始 item 集合会先按 replica进程切分成互不重叠的子集每个 rank 只在自己的子集上执行可选的shuffle 与 batch 划分因此每个 replica 始终获得确定且互斥的一份数据。构造参数参数类型默认值说明item_setItemSet/HeteroItemSet必填待采样的数据例如train_set、validation_set、test_setbatch_sizeint必填mini-batch 的大小即一批处理的样本数量drop_lastboolFalse是否丢弃最后一个不完整的 batchshuffleboolFalse是否在采样前打乱数据drop_uneven_inputsboolFalse是否让所有 rank 的 batch 数量保持一致丢弃多出的部分seedintNone可复现的随机种子为None时自动生成构造时item_sampler.py采样器会通过dist.get_world_size()与dist.get_rank()获取当前分布式环境信息并在world_size 1时调用_align_seeds(src0)用dist.broadcast将 rank 0 的种子同步给所有 rankitem_sampler.py从而保证各进程的随机行为一致、便于复现实验。实际分片行为内部通过calculate_range计算每个 rank 与每个 worker 的起止范围实现在 python/dgl/graphbolt/internal/item_sampler_utils.py。以torch.arange(15)、batch_size2、4 个 replica 为例源码 docstring 给出了各类参数组合下的输出全部为False时Replica#0 得到[0,1],[2,3]Replica#1 得到[4,5],[6,7]Replica#2 得到[8,9],[10,11]Replica#3 得到[12,13],[14]各 rank 的 batch 数不均等drop_lastTrue, drop_uneven_inputsFalseReplica#3 只剩[12,13]其余不变drop_lastTrue, drop_uneven_inputsTrue所有 rank 都只保留 1 个 batchReplica#0[0,1]、Replica#1[4,5]、Replica#2[8,9]、Replica#3[12,13]各 rank batch 数量完全一致shuffleTrue时各 rank 在自己的子集内以seed epoch为随机种子打乱数据且每个 epoch 顺序都会变化。注意DistributedItemSampler特意没有使用torch.utils.data.functional_datapipe装饰即不支持函数式调用但可以继续追加其他可迭代 datapipe如copy_to、sample_neighbor、fetch_feature。构建分布式 Dataloader 流水线示例中的create_dataloader函数node_classification.py完整演示了多卡数据流水线的组装方式datapipe gb.DistributedItemSampler( item_setitemset, batch_sizeargs.batch_size, drop_lastis_train, shuffleis_train, drop_uneven_inputsis_train, ) if args.storage_device ! cpu: datapipe datapipe.copy_to(device) datapipe datapipe.sample_neighbor( graph, args.fanout, overlap_fetchargs.storage_device pinned, asynchronousargs.storage_device ! cpu, ) datapipe datapipe.fetch_feature(features, node_feature_keys[feat]) if args.storage_device cpu: datapipe datapipe.copy_to(device) dataloader gb.DataLoader(datapipe, args.num_workers)各步骤的作用gb.DistributedItemSampler按 rank 切分 item 子集并产出 mini-batch训练阶段is_trainTrue同时开启shuffle、drop_last与drop_uneven_inputs验证/测试阶段三者均关闭copy_to(device)非 CPU 存储时提前执行将数据先拷贝到目标设备使后续采样操作直接在 GPU 上运行sample_neighbor(graph, fanout, ...)为每个 batch 的种子节点采样邻居fanout长度必须与模型层数一致默认10,10,10对应三层 GraphSAGEoverlap_fetch在 pinned 内存模式下开启以重叠取数asynchronous在非 CPU 存储时开启异步采样fetch_feature(features, node_feature_keys[feat])为采样得到的子图拉取节点特征gb.DataLoader(datapipe, args.num_workers)用num_workers个进程并行加载数据。训练循环中从每个data中解出三部分node_classification.pyx data.node_features[feat] # 第一层计算图的源节点特征 y data.labels # 最后一层计算图的目标节点标签 blocks data.blocks # 逐层的消息传递块 y_hat model(blocks, x)其中blocks是邻居采样产生的多层计算图块Block逐层经过SAGE.forward中的SAGEConv聚合方式为mean隐藏层 256 维中间层经过ReLU与Dropout(0.5)node_classification.py。模型与训练循环DDP Join 加权聚合构建分布式模型model SAGE(in_size, hidden_size, out_size).to(device) model DDP(model)每个 rank 拥有一份完整的模型副本replicaDDP负责在反向传播时通过allreduce同步梯度。特征输入维度in_size通过feature.size(node, None, feat)[0]动态获取out_size为num_classes。Join 上下文管理器处理不均衡输入DDP 要求所有 rank 的输入数量一致否则程序可能报错或挂起。示例提供了两种解决方案PyTorch 的Join上下文管理器示例采用的方式with Join([model]): for data in (tqdm.tqdm(train_dataloader) if rank 0 else train_dataloader): ...drop_uneven_inputsTrue在DistributedItemSampler中设置通过丢弃多出的 batch 使各 rank 的 batch 数量一致。示例在训练阶段同时开启drop_uneven_inputs与Join双重保障训练不会因输入不均衡而中断。进度条tqdm只在 rank 0 上显示避免多进程输出刷屏。加权聚合损失与准确率由于各 GPU 处理的样本数量可能不同简单取平均会得到有偏的指标。示例实现了weighted_reducenode_classification.pydef weighted_reduce(tensor, weight, dst0): dist.reduce(tensortensor, dstdst) weight torch.tensor(weight, devicetensor.device) dist.reduce(tensorweight, dstdst) return tensor / weightdist.reduce将各 rank 的张量归约到dst指定的进程默认 rank 0默认使用ReduceOp.SUM求和损失乘以各自样本数求和后再除以总样本数得到精确的加权平均损失验证准确率同样以acc * num_val_items加权求和后除以总样本数node_classification.py。训练每个 epoch 后还会调用torch.cuda.synchronize()同步保证计时准确并在 rank 0 打印Epoch / Average Loss / Accuracy / Time。验证与测试评估evaluate函数node_classification.py在torch.no_grad()下遍历 dataloader每个 rank 在自己的验证/测试子集上推理收集预测y_hats与标签y使用torchmetrics.functional.accuracytaskmulticlass计算准确率返回(准确率, 样本数)供weighted_reduce加权平均。测试阶段同样只在 rank 0 打印Test Accuracy最后调用dist.destroy_process_group()清理进程组node_classification.py。命令行参数详解脚本通过argparse提供以下参数node_classification.py参数默认值可选值/说明--gpu0逗号分隔的 GPU 编号如0,1,2,3GPU 数量即world_size--epochs10训练轮数--lr0.001学习率Adam 优化器--batch-size1024mini-batch 大小--fanout10,10,10邻居采样扇出逗号分隔长度必须与模型层数一致--num-workers0数据加载进程数--gpu-cache-size0GPU 特征缓存容量字节--datasetogbn-products支持ogbn-arxiv、ogbn-products、ogbn-papers100M--modepinned-cuda数据存储位置与训练设备组合cpu-cuda图/特征在 CPU 内存、pinned-cuda图/特征在页锁定内存、cuda-cuda图/特征在 GPU 显存其中--mode会被拆分为storage_device与训练设备两部分args.storage_device, _ args.mode.split(-)并决定数据流水线中copy_to、overlap_fetch、asynchronous的行为。当--gpu-cache-size 0且存储设备不是cuda时示例会用gb.gpu_cached_feature实现见 python/dgl/graphbolt/impl/gpu_cached_feature.py为节点特征挂上 GPU 缓存将热点特征缓存在显存中以减少 PCIe 传输。数据集通过gb.BuiltinDataset(args.dataset).load()加载为OnDiskDataset训练/验证/测试子集分别取自dataset.tasks[0]的train_set、validation_set、test_set类别数来自dataset.tasks[0].metadata[num_classes]。使用建议单机多卡场景下默认的tcp://127.0.0.1:12345初始化即可工作跨机训练需修改init_method并配合 DGL 的分布式工具链使用--fanout长度必须与模型层数严格一致否则采样阶段会报错验证/测试阶段不要开启shuffle与drop_last以保证评估覆盖全部样本若发现 DDP 训练因输入不均衡而卡死优先确认训练 dataloader 是否同时启用了Join上下文管理器或drop_uneven_inputs进一步了解 GraphBolt 数据加载器各 datapipe 的详细用法可参考 notebooks/graphbolt/walkthrough.ipynb。赞分享人工智能机器学习深度学习图计算【免费下载链接】dglPython package built to ease deep learning on graph, on top of existing DL frameworks.项目地址https://gitcode.com/gh_mirrors/dg/dgl点击查看免费下载相关推荐DGL 多 GPU 分布式训练实战基于 PyTorch DDP 与 GraphBolt 的 GraphSAGE 节点分类DGL 多 GPU 分布式训练实战基于 PyTorch DDP 与 GraphBolt 的 GraphSAGE 节点分类 本篇技术指南以 DGL 官方多 GP人工智能机器学习深度学习图计算使用 DGL Sparse 与 GraphBolt 完成 GraphSAGE 小批量训练实战指南使用 DGL Sparse 与 GraphBolt 完成 GraphSAGE 小批量训练实战指南 本文基于 DGL 官方指南 docs/source/guide人工智能机器学习深度学习图计算DGL 节点分类/回归完整指南从 GraphSAGE 到 Heterogeneous RGCN 的实战训练DGL 节点分类/回归完整指南从 GraphSAGE 到 Heterogeneous RGCN 的实战训练 导读 节点分类/回归是图神经网络GNN最经典人工智能机器学习深度学习图计算创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考