大模型分布式训练中的NCCL应用
·
大模型训练框架,训练梯度需要用到 DDP。
流水线如下图

Bucket 到底怎么对应网络层?
以 4 个bucket为例,假设网络总共有 20 层:layer0 ~ layer19;反向传播顺序:layer19 → layer18 … → layer0(从输出层往输入层算梯度)
bucket_cap_mb 设置 25MB,框架自动把连续的网络层打包归到同一个 Bucket:
- Bucket‑B3:layer19、layer18、layer17、layer16 (几层梯度加起来≈25MB)
- Bucket‑B2:layer15、layer14、layer13、layer12
- Bucket‑B1:layer11、layer10、layer9、layer8
- Bucket‑B0:layer7 … layer0
计算通信 overlap 流水线代码
#include <vector>
#include <cstdio>
#include <cuda_runtime.h>
#include <mpi.h>
#include <nccl.h>
#define CHECK_CUDA(err) \
if(err != cudaSuccess) {printf("CUDA error %d\n", err);exit(-1);}
#define CHECK_NCCL(err) \
if(err != ncclSuccess) {printf("NCCL error %s\n", ncclGetErrorString(err));exit(-1);}
#define BUCKET_NUM 4
#define ELEM_PER_BUCKET 1024
struct GradBucket {
float* d_grad;
bool is_computed;
};
// 模拟梯度计算kernel
__global__ void computeGradientKernel(float* grad, int elem_cnt, int bid)
{
int idx = blockIdx.x * blockDim.x + threadIdx.x;
if(idx < elem_cnt) {
grad[idx] = bid * 1.0f;
}
}
// 结果校验kernel
__global__ void verifyResultKernel(float* grad, int elem_cnt, int nranks, int bid)
{
int idx = blockIdx.x * blockDim.x + threadIdx.x;
if(idx < elem_cnt) {
// ...校验逻辑省略
}
}
int main(int argc, char** argv)
{
int rank = 0;
int nranks = 1;
// ============ 1、MPI初始化(多卡环境必需) ============
MPI_Init(&argc, &argv);
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
MPI_Comm_size(MPI_COMM_WORLD, &nranks);
// 绑定GPU: 每张卡对应一个进程
CHECK_CUDA(cudaSetDevice(rank));
printf("[rank %d] MPI initialized: rank=%d, world_size=%d\n", rank, rank, nranks);
// ============ 2、NCCL初始化 ============
ncclComm_t comm;
ncclUniqueId id;
// rank 0生成唯一ID
if (rank == 0) {
CHECK_NCCL(ncclGetUniqueId(&id));
}
// 广播ID给所有rank(关键!让所有GPU用同一个ID建群)
CHECK_NCCL(ncclBcast(&id, sizeof(id), ncclChar, 0, MPI_COMM_WORLD));
// 每个rank用自己的rank初始化NCCL通信器
CHECK_NCCL(ncclCommInitRank(&comm, nranks, id, rank));
printf("[rank %d] NCCL communicator initialized\n", rank);
// ============ 3、创建两套CUDA stream(Overlap的核心!) ============
cudaStream_t compute_stream; // GPU反向计算流: 跑梯度计算kernel
cudaStream_t comm_stream; // 通信流: 跑ncclAllReduce, 后台异步通信
CHECK_CUDA(cudaStreamCreate(&compute_stream));
CHECK_CUDA(cudaStreamCreate(&comm_stream));
// ============ 4、分配多个bucket梯度显存(模拟模型多层梯度) ============
std::vector<GradBucket> buckets(BUCKET_NUM);
for(int b = 0; b < BUCKET_NUM; b++){
CHECK_CUDA(cudaMalloc(&buckets[b].d_grad, sizeof(float) * ELEM_PER_BUCKET));
buckets[b].is_computed = false;
}
// 用于记录时间,验证Overlap效果
cudaEvent_t start_event, end_event;
CHECK_CUDA(cudaEventCreate(&start_event));
CHECK_CUDA(cudaEventCreate(&end_event));
CHECK_CUDA(cudaEventRecord(start_event, compute_stream));
// ============ 5、模拟反向传播时序(核心Overlap逻辑!) ============
// 遍历bucket:顺序 B3 → B2 → B1 → B0,对应神经网络 L → L‑1 → L‑2层
for(int bid = BUCKET_NUM - 1; bid >= 0; bid--)
{
GradBucket& bucket = buckets[bid];
// -------- Step1: GPU在compute_stream执行【梯度计算】 --------
// 这才是真实的场景: 反向传播kernel在compute_stream上异步执行
int threads = 256;
int blocks = (ELEM_PER_BUCKET + threads - 1) / threads;
// 启动梯度计算kernel(异步!立刻返回)
computeGradientKernel<<<blocks, threads, 0, compute_stream>>>(
bucket.d_grad, ELEM_PER_BUCKET, bid);
CHECK_CUDA(cudaGetLastError());
// 【关键优化点】这里不要sync! 继续往下走才能Overlap
// 标记这个bucket已经"计算完成"(实际kernel还在跑)
bucket.is_computed = true;
// -------- Step2: 异步提交NCCL AllReduce到comm_stream --------
// ✅精髓: AllReduce依赖compute_stream的结果,需要加stream等待
// 这样通信会等计算完成才开始,但CPU主线程不阻塞
CHECK_CUDA(cudaStreamWaitEvent(comm_stream, compute_stream, 0));
// 提交AllReduce到comm stream(异步!立刻返回)
CHECK_NCCL(ncclAllReduce(
bucket.d_grad, // 输入buffer(就地AllReduce)
bucket.d_grad, // 输出buffer
ELEM_PER_BUCKET, // 元素个数
ncclFloat32, // 数据类型
ncclSum, // 操作: 求和
comm, // NCCL通信器
comm_stream // 通信stream
));
// 【关键优化点】提交完NCCL立刻继续循环,不等待通信完成!
// CPU马上进入下一轮迭代,GPU并行干活:
// - compute_stream: 计算下一个bucket的梯度
// - comm_stream: 后台异步跑当前bucket的AllReduce
// 打印调试信息(可选)
if (rank == 0) {
printf("Bid %d: compute kernel launched, AllReduce submitted to comm_stream\n", bid);
}
}
// ============ 6、全部任务提交完毕,最后统一等待 ============
// 此时:
// - compute_stream: 可能还有最后一个bucket的计算没做完
// - comm_stream: 还有4个AllReduce在后台排队/执行
// 等待所有计算完成(compute_stream上的所有kernel)
CHECK_CUDA(cudaStreamSynchronize(compute_stream));
if (rank == 0) printf("All compute kernels done.\n");
// 等待所有通信完成(comm_stream上的所有AllReduce)
CHECK_CUDA(cudaStreamSynchronize(comm_stream));
if (rank == 0) printf("All AllReduce communications done.\n");
// 记录结束事件
CHECK_CUDA(cudaEventRecord(end_event, comm_stream));
CHECK_CUDA(cudaEventSynchronize(end_event));
float elapsed_ms = 0;
CHECK_CUDA(cudaEventElapsedTime(&elapsed_ms, start_event, end_event));
if (rank == 0) printf("Total time: %.2f ms\n", elapsed_ms);
// ============ 7、验证结果(可选) ============
// 检查AllReduce结果: 所有rank的grad应该相同(因为是求和)
threads = 256;
blocks = (ELEM_PER_BUCKET + threads - 1) / threads;
verifyResultKernel<<<blocks, threads>>>(buckets[0].d_grad, ELEM_PER_BUCKET, nranks, 0);
CHECK_CUDA(cudaDeviceSynchronize());
// ============ 8、清理资源 ============
for(auto& b : buckets) CHECK_CUDA(cudaFree(b.d_grad));
CHECK_CUDA(cudaStreamDestroy(compute_stream));
CHECK_CUDA(cudaStreamDestroy(comm_stream));
CHECK_CUDA(cudaEventDestroy(start_event));
CHECK_CUDA(cudaEventDestroy(end_event));
CHECK_NCCL(ncclCommDestroy(comm));
MPI_Finalize();
if (rank == 0) printf("Cleanup done, Success!\n");
return 0;
}
naive ringAllReduce 代码
#include <cuda_runtime.h>
#include <stdio.h>
#include <vector>
// ===================== 配置常量(对应图中 4 卡场景)=====================
const int GPU_NUM = 4; // GPU 数量:0,1,2,3 环形
const int SEGMENT_ELEM = 1024; // 每个分片元素数量
const int TOTAL_ELEM = GPU_NUM * SEGMENT_ELEM; // 全局梯度总长度
// 错误检查宏
#define CHECK_CUDA (err)
if (err != cudaSuccess) {
printf("CUDA Error: %s at line %d\n", cudaGetErrorString(err), **LINE**);
exit(-1);
}
// ===================== CUDA Kernel:GPU 端分片加法规约 =====================
**global** void addSegmentKernel(float* dst, const float* src, int elem_cnt)
{
int idx = blockIdx.x * blockDim.x + threadIdx.x;
if (idx < elem_cnt) {
dst[idx] += src[idx];
}
}
// ===================== 工具函数:获取当前 GPU 右邻居 ID(环形)=====================
inline int getRightNeighbor (int self_gpu)
{
return (self_gpu + 1) % GPU_NUM;
}
// ===================== 工具函数:获取当前 GPU 左邻居 ID(环形)=====================
inline int getLeftNeighbor (int self_gpu)
{
return (self_gpu - 1 + GPU_NUM) % GPU_NUM;
}
int main ()
{
// 1. 初始化:分配多卡显存、开启 P2P 互访
std::vector<float*> dev_buf (GPU_NUM, nullptr); // 每张 GPU 完整梯度显存
std::vector<cudaStream_t> streams (GPU_NUM); // 每张 GPU 独立流
for (int gpu = 0; gpu < GPU_NUM; gpu++)
{
CHECK_CUDA (cudaSetDevice (gpu));
// 分配全局梯度显存
CHECK_CUDA (cudaMalloc (&dev_buf [gpu], sizeof (float) * TOTAL_ELEM));
// 创建异步流
CHECK_CUDA (cudaStreamCreate (&streams [gpu]));
// 开启所有 GPU 两两 P2P 访问权限
for (int peer = 0; peer < GPU_NUM; peer++)
{
if (gpu != peer)
{
int can_access;
CHECK_CUDA (cudaDeviceCanAccessPeer (&can_access, gpu, peer));
if (can_access)
CHECK_CUDA (cudaDeviceEnablePeerAccess (peer, 0));
}
}
// 初始化本地梯度:GPU n 的初始梯度 g_n = 全值 n
std::vector<float> host_init(TOTAL_ELEM, (float)gpu);
CHECK_CUDA(cudaMemcpyAsync(dev_buf[gpu], host_init.data(), sizeof(float)*TOTAL_ELEM, cudaMemcpyHostToDevice, streams[gpu]));
CHECK_CUDA(cudaStreamSynchronize(streams[gpu]));
}</float>
//-------------------------- 阶段 1:Scatter-Reduce 分片归约 --------------------------
// 循环 GPU_NUM-1 = 3 轮(4 卡需要 3 轮收发归约)
for (int step = 0; step < GPU_NUM - 1; step++)
{
printf ("==== Scatter-Reduce Step % d ====\n", step);
for (int gpu = 0; gpu < GPU_NUM; gpu++)
{
CHECK_CUDA (cudaSetDevice (gpu));
int right_gpu = getRightNeighbor (gpu);
// 当前轮次需要发送的分片索引
int send_seg_idx = (gpu - step + GPU_NUM) % GPU_NUM;
// 分片显存偏移
size_t seg_offset = send_seg_idx * SEGMENT_ELEM * sizeof (float);
// 1. 异步把自己的分片 发给 右邻居 GPU
CHECK_CUDA (cudaMemcpyPeerAsync (
dev_buf [right_gpu] + send_seg_idx * SEGMENT_ELEM, // 目标:右邻居显存分片位置
right_gpu,
dev_buf [gpu] + send_seg_idx * SEGMENT_ELEM, // 源:当前 GPU 本地分片
gpu,
SEGMENT_ELEM * sizeof (float),
streams [gpu]
));
// 2. 右邻居收到分片后,原地累加归约
CHECK_CUDA (cudaSetDevice (right_gpu));
addSegmentKernel<<<(SEGMENT_ELEM + 256 -1)/256, 256, 0, streams [right_gpu]>>>(
dev_buf [right_gpu] + send_seg_idx * SEGMENT_ELEM,
dev_buf [gpu] + send_seg_idx * SEGMENT_ELEM,
SEGMENT_ELEM
);
CHECK_CUDA (cudaGetLastError ());
}
// 等待本轮所有 GPU 通信、计算完成
for (int gpu = 0; gpu < GPU_NUM; gpu++)
{
CHECK_CUDA (cudaSetDevice (gpu));
CHECK_CUDA (cudaStreamSynchronize (streams [gpu]));
}
}
// Scatter-Reduce 结束:每个 GPU 仅持有一段【全局总和分片】
//-------------------------- 阶段 2:All-Gather 分片全收集 --------------------------
for (int step = 0; step < GPU_NUM -1; step++)
{
printf ("==== All-Gather Step % d ====\n", step);
for (int gpu = 0; gpu < GPU_NUM; gpu++)
{
CHECK_CUDA (cudaSetDevice (gpu));
int right_gpu = getRightNeighbor (gpu);
int send_seg_idx = (gpu - step -1 + GPU_NUM) % GPU_NUM;
size_t seg_offset = send_seg_idx * SEGMENT_ELEM * sizeof (float);
// 仅传输,无需计算:把已求和的分片转发给右邻居
CHECK_CUDA (cudaMemcpyPeerAsync (
dev_buf [right_gpu] + send_seg_idx * SEGMENT_ELEM,
right_gpu,
dev_buf [gpu] + send_seg_idx * SEGMENT_ELEM,
gpu,
SEGMENT_ELEM * sizeof (float),
streams [gpu]
));
}
for (int gpu = 0; gpu < GPU_NUM; gpu++)
{
CHECK_CUDA (cudaSetDevice (gpu));
CHECK_CUDA (cudaStreamSynchronize (streams [gpu]));
}
}
// ===================== 验证结果 =====================
printf ("\n==== 结果校验:全局总和 g = g0+g1+g2+g3 = 0+1+2+3 = 6.0 ====\n");
std::vector<float> host_res (TOTAL_ELEM);
for (int gpu = 0; gpu < GPU_NUM; gpu++)
{
CHECK_CUDA (cudaSetDevice (gpu));
CHECK_CUDA (cudaMemcpy (host_res.data (), dev_buf [gpu], sizeof (float)*TOTAL_ELEM, cudaMemcpyDeviceToHost));
float val = host_res [0];
printf ("GPU % d 第一个梯度值 = %.1f\n", gpu, val);
}</float>
// ===================== 资源释放 =====================
for (int gpu = 0; gpu < GPU_NUM; gpu++)
{
CHECK_CUDA (cudaSetDevice (gpu));
CHECK_CUDA (cudaFree (dev_buf [gpu]));
CHECK_CUDA (cudaStreamDestroy (streams [gpu]));
for (int peer =0; peer < GPU_NUM; peer++)
{
if (gpu != peer)
cudaDeviceDisablePeerAccess (peer);
}
}
return 0;
}
更多推荐
所有评论(0)