大模型训练框架,训练梯度需要用到 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;

}

更多推荐