首先先看上一节遗留的多线程NCCL

#ifndef __CUDACC__
#define __CUDACC__
#endif
#include <cuda_runtime.h>
#include <nccl.h>
#include <iostream>
#include<vector>
#include<thread>
void checkCuda(cudaError_t res) {
    if (res != cudaSuccess) {
        std::cerr << "CUDA Error: " << cudaGetErrorString(res) << std::endl;
        exit(EXIT_FAILURE);
    }
}

void checkNCCL(ncclResult_t res) {
    if (res != ncclSuccess) {
        std::cerr << "NCCL Error: " << ncclGetErrorString(res) << std::endl;
        exit(EXIT_FAILURE);
    }
}
const int nDev=2;
const int size=32;
void thread_func(int rank,ncclUniqueId id){
    int dev=rank;
    cudaSetDevice(dev);
    float *sendbuff;
    float *recvbuff;
    //分配内存
    checkCuda(cudaMalloc(&sendbuff, size * sizeof(float)));
    checkCuda(cudaMalloc(&recvbuff, size * sizeof(float)));
        // 初始化数据
    float h_data[size];
    for (int i = 0; i < size; i++) h_data[i] = float(rank + 1); // rank=0 ->1,rank=1->2
    checkCuda(cudaMemcpy(sendbuff, h_data, size * sizeof(float), cudaMemcpyHostToDevice));
    // 创建流
    cudaStream_t stream;
    checkCuda(cudaStreamCreate(&stream));
    ncclComm_t comm;
    checkNCCL(ncclCommInitRank(&comm,nDev,id,rank));
    // AllReduce
    checkNCCL(ncclAllReduce(sendbuff, recvbuff, size,
                            ncclFloat, ncclSum, comm, stream));
    // 等待完成
    checkCuda(cudaStreamSynchronize(stream));

    // 拷贝结果回 Host
    checkCuda(cudaMemcpy(h_data, recvbuff, size * sizeof(float), cudaMemcpyDeviceToHost));
    std::cout << "Thread/Rank " << rank << " GPU " << dev
              << " result[0] = " << h_data[0] << std::endl;

    // 清理
    ncclCommDestroy(comm);
    cudaFree(sendbuff);
    cudaFree(recvbuff);
    cudaStreamDestroy(stream);
}
int main() {
    // 获取 UniqueId (多进程时 rank0 生成,广播给其他进程)
    ncclUniqueId id;
    checkNCCL(ncclGetUniqueId(&id));
    // 启动两个线程,模拟两个进程
    std::vector<std::thread> threads;
    for (int rank = 0; rank < nDev; rank++) {
        threads.emplace_back(thread_func, rank, id);
    }
    for (auto& t : threads) t.join();
}

通信原语

ncclBroadcast

ncclResult_t ncclBroadcast(
    const void* sendbuff, void* recvbuff,
    size_t count, ncclDataType_t datatype,
    int root, ncclComm_t comm, cudaStream_t stream);

参数

  • sendbuff:只有 root GPU 使用,表示要广播的数据。

  • recvbuff:所有 GPU(包括 root)接收数据的缓冲区。

  • count:元素个数。

  • datatype:数据类型(如 ncclFloat, ncclInt)。

  • root:源 GPU 的 rank(0 ~ nDev-1)。

  • comm:NCCL 通信器。

  • stream:CUDA stream,支持异步。

功能

  • root GPU 把数据广播给所有 GPU。

典型应用

  • 分布式训练时,把 最新模型参数 广播给所有 worker。

ncclReduce

ncclResult_t ncclReduce(
    const void* sendbuff, void* recvbuff,
    size_t count, ncclDataType_t datatype,
    ncclRedOp_t op, int root,
    ncclComm_t comm, cudaStream_t stream);

参数

  • sendbuff:每个 GPU 上待参与 reduce 的数据。

  • recvbuff:只有 root GPU 接收 reduce 结果。

  • op:归约操作(ncclSum, ncclProd, ncclMax, ncclMin)。

  • root:结果接收 GPU 的 rank。

功能

  • 所有 GPU 的数据按照 op 进行归约,并放到 root GPU 上。

典型应用

  • 把多个 GPU 计算的 梯度聚合到一个 GPU,由它来更新参数。

ncclAllReduce

前文已经介绍过

功能

  • 所有 GPU 的数据先 reduce(归约),再广播给所有 GPU。

  • 每个 GPU 都得到同样的结果。

典型应用

  • 分布式训练 DDP(Data Parallel) 中的梯度同步。

ncclReduceScatter

ncclResult_t ncclReduceScatter(
    const void* sendbuff, void* recvbuff,
    size_t recvcount, ncclDataType_t datatype,
    ncclRedOp_t op,
    ncclComm_t comm, cudaStream_t stream);

参数

  • sendbuff:每个 GPU 上的数据(大小 = recvcount * nDev)。

  • recvbuff:每个 GPU 接收的数据块(大小 = recvcount)。

功能

  • 先做一次 reduce(所有 GPU 数据求和),然后 把结果均匀切分,分给各个 GPU。

典型应用

  • 张量并行(Tensor Parallel):比如矩阵乘法分块,每个 GPU 只保留自己那份结果。

ncclAllGather

ncclResult_t ncclAllGather(
    const void* sendbuff, void* recvbuff,
    size_t sendcount, ncclDataType_t datatype,
    ncclComm_t comm, cudaStream_t stream);

参数

  • sendbuff:每个 GPU 的本地数据(大小 = sendcount)。

  • recvbuff:每个 GPU 接收的数据(大小 = sendcount * nDev)。

功能

  • 每个 GPU 把自己的数据广播出去,所有 GPU 都收集到完整的数据。

典型应用

  • 模型并行:不同 GPU 保存不同的参数分片,最后拼接在一起。

  • 数据拼接:比如 NLP 里需要把不同 GPU 上的 token 拼成一个 batch。

总结表格

函数功能数据流向应用场景
ncclBroadcastroot → all1 → n模型参数分发
ncclReduceall → rootn → 1梯度聚合
ncclAllReduceall → alln ↔ nDDP 梯度同步
ncclReduceScatterall → 部分n → n(切分)Tensor Parallel
ncclAllGather部分 → alln ↔ n参数/数据拼接
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐