Skip to main content

录屏回看

本章概要

分布式 AI 系统的核心瓶颈经常不是 GPU 算得慢,而是 GPU 在等数据。训练时,梯度、参数分片和激活值需要在多张 GPU 之间同步;推理时,KV cache、请求状态和模型分片需要在不同 worker 之间移动。规模越大,通信路径、同步时机和网络拓扑越影响整体吞吐。 本章以一个梯度 tensor 在分布式训练中的完整旅程为主线:从 loss.backward() 产生梯度开始,经过 DDP bucket、NCCL collective、NVLink/InfiniBand 硬件路径,到最终所有 GPU 拿到全局平均梯度。Part 1 用直觉建立全貌,并深入 PyTorch DDP 的框架层实现——为什么 GPU 之间需要通信、DDP 如何用 bucket 机制组织梯度同步、还有哪些通信类型。Part 2 深入 NCCL 与 NVIDIA 通信栈的每一层:Collective 算法(Ring/Tree)、节点内互联(NVLink/NVSwitch)、跨节点互联(RDMA/GPUDirect RDMA)、网络内计算(SHARP)。Part 3 是动手实验,把原理落到可观察的操作上。

1. 为什么 GPU 之间需要通信

1.1 一个直觉:8 个人合算全班平均分

假设一次考试有 800 个学生,分给 8 位老师批改,每人批 100 份。每位老师算出自己那 100 份的平均分后,怎么得到全班 800 人的平均分? 最笨的办法:每位老师把自己的平均分喊给其他 7 个人,每个人都收到 7 份数据,自己算一遍全班平均。通信量是 N×(N-1) 次——8 个人就是 56 次传话,80 个人就是 6320 次,规模一大就爆了。 聪明的办法:8 个人排成一圈,每人只跟左右邻居交换和累加,转一圈就够了。这就是 ring all-reduce 的直觉——不需要每个人跟所有人通信,只需要跟邻居交换,经过足够多的步骤,所有人都能拿到全局结果。
Eight participants arranged in a ring with neighbor-to-neighbor arrows, illustrating ring all-reduce intuition

Ring all-reduce 的直觉:8 个参与者只和相邻节点交换数据,经过多轮传递后每个节点都得到同一个全局结果

分布式训练里 GPU 之间同步梯度就是这件事:

1.2 一个梯度 tensor 的宏观旅程

理解了”为什么”之后,来看一个具体的训练步骤里发生了什么。假设 8 张 GPU 做数据并行训练,每张 GPU 持有完整模型副本,处理不同的 mini-batch: 第一步,loss.backward() 让每张 GPU 都算出自己的本地梯度。第二步,这些梯度不会立刻单独发送,而是先进入 DDP 预先划分好的 bucket;bucket 里的梯度都准备好后,才触发一次同步。第三步,通信库接手这次同步,把这个 bucket 里的数据沿着 GPU 之间的高速链路发出去。第四步,同步完成后,每张 GPU 都拿到相同的全局平均梯度。最后,optimizer.step() 用这份平均梯度更新模型副本,所以所有 GPU 上的模型继续保持一致。 这就是一个梯度 tensor 的完整旅程:从本地梯度出发,先在 bucket 里等待同伴,再通过高速链路和其他 GPU 合算,最终变成每张 GPU 都相同的全局平均梯度。Part 2 会深入通信栈本身:NCCL 怎么组织 collective,NVLink、InfiniBand/RoCE 和 GPUDirect RDMA 又如何承载真正的数据传输。
Simple journey of one gradient tensor from local backward gradients through DDP bucket synchronization over high-speed GPU links to identical averaged gradients on all GPUs

一个梯度 tensor 的直觉旅程:先产生本地梯度,进入 DDP bucket,bucket 就绪后发起同步,数据经高速链路在 GPU 间交换,最后所有 GPU 得到相同的全局平均梯度

关键问题是第 ② 步:梯度同步要花时间,GPU 在等同步的时候就没在算东西。聪明的做法是边算边同步——不等所有层的梯度都算完,只要某个 bucket 里的梯度都 ready,就立刻对这个 bucket 发起 all-reduce,同时继续计算后续层的梯度。
DDP timeline showing backward computation overlapping with all-reduce launched after buckets become ready

DDP 通信计算重叠:backward 继续推进时,已就绪的 bucket 可以先进入 AllReduce;触发条件是 bucket ready,而不是单个梯度 ready

这就是 PyTorch DDP 的 bucket 机制——把梯度分成若干桶,每个桶一满就发起同步,不等所有层算完。接下来看 DDP 在代码层面具体怎么实现这件事。

1.3 DDP:PyTorch 如何组织梯度同步

DDP(DistributedDataParallel)是 PyTorch 里最常用的数据并行训练封装。它的基本假设是:每个进程控制一张 GPU,每张 GPU 上都有一份完整模型副本,但每个进程处理不同的 mini-batch。forward 和 backward 都在本地 GPU 上独立执行;真正需要跨 GPU 协作的是 backward 产生的梯度同步。 DDP 要保证两件事:第一,每张 GPU 最终看到相同的全局平均梯度;第二,梯度同步尽量不要阻塞后续 backward 计算。为此,DDP 会在初始化时把参数预先分进 bucket,并给参数梯度注册 autograd hook。backward 过程中,某个 bucket 里的梯度都 ready 后,DDP 就通过通信后端发起 all-reduce,而不是等整个 backward 全部结束。

1.3.1 代码:DataParallel vs DistributedDataParallel

理解了 DDP 的核心思想后,先看它和旧式 DataParallel 在使用上的区别: Before — DataParallel(主卡瓶颈):
After — DistributedDataParallel(NCCL all-reduce,无主卡瓶颈):

1.3.2 DDP 代码架构

DDP 的核心逻辑在三层之间分工:distributed.py 是 Python 入口,负责初始化和 forward;reducer.h 是 C++ 核心,负责 bucket 分配、autograd hook 注册和梯度同步调度;ProcessGroup.hpp 是通信抽象层,把 broadcast 和 allreduce 等操作分发给具体的后端实现(NCCL 用于 NVIDIA GPU,RCCL 用于 AMD GPU,Gloo 用于 CPU)。
DDP architecture: distributed.py → reducer.h + comm.h → ProcessGroup → NCCL/Gloo/MPI/RR/XCCL

DDP 分层架构:Python 入口层(distributed.py)→ C++ 梯度同步与广播(reducer.h / comm.h)→ 通信后端抽象(ProcessGroup.hpp)→ 具体实现(NCCL / Gloo / MPI / RCCL / XCCL)

DDP bucket allreduce: Process 0 and Process 1 each have params grouped into buckets, with allreduce between corresponding buckets

DDP 梯度同步:参数按 reverse 顺序分到 bucket 里(bucket1 = grad0+grad1,bucket0 = grad2+grad3),每个 bucket 就绪后独立发起 allreduce,两个 Process 的 bucket 结构完全对称

Gradient blocks are grouped into buckets first, then completed buckets launch all-reduce before optimizer step

DDP bucket 机制细节:梯度先进入预先划分的 bucket,bucket 内所有梯度 ready 后,才对这个 bucket 发起 AllReduce

1.3.3 DDP backward 全流程:CPU 调度 + GPU 执行

Bucket 按 reverse parameter 顺序排列——backward 先算完的梯度尽量排在更早 ready 的 bucket,使通信尽早开始、与计算 overlap。为了保证所有 rank 发起 collective 的顺序一致,DDP 会按 bucket index 顺序调度 all-reduce;bucket 是预先划分好的分组,all-reduce 作用在已经 ready 的 bucket 上。
DDP backward flow with CPU reducer hooks, GPU compute stream, and GPU communication stream for bucket all-reduce

DDP backward 全流程:autograd hook 只标记梯度 ready,reducer 检查预先划分的 bucket;bucket 完成后按 bucket index 顺序提交 AllReduce 到通信流

CPU 只做轻量调度(hook 触发、bucket 检查、提交 NCCL 调用),重活全在 GPU 上。通信完成后通知 compute stream 用的是 CUDA Event——GPU 到 GPU 的 stream 同步,不需要 CPU 介入。

1.4 不只是梯度:分布式 AI 系统需要通信什么

到这里,我们已经完整追踪了一个梯度 tensor 的旅程——从 loss.backward() 到 DDP bucket 到 NCCL all-reduce。但分布式系统中需要通信的远不止梯度。 几个观察:
  1. 所有性能敏感的通信都是 tensor/buffer,走 GPU 间的高速链路。控制面走 TCP,数据量小,不是瓶颈。
  2. 训练通信以 collective 为主(all-reduce, reduce-scatter, all-gather, all-to-all)——所有参与者必须同步推进,最慢的 rank 决定整体速度。
  3. 推理通信以 point-to-point 为主(KV cache 迁移、pipeline 激活传递)——不需要全员参与,但对延迟极度敏感,每一毫秒直接加到用户感知的响应时间上。
  4. 通信路径相同,优化目标不同:训练追求吞吐(把通信藏在计算背后),推理追求延迟(让单次传输尽可能快)。
这张表就是后续所有内容的作用对象。Part 2 深入通信栈的每一层:NCCL 怎么在 NVLink、NVSwitch、InfiniBand/RoCE 等硬件之上完成 collective 通信,RDMA 和 GPUDirect RDMA 怎么绕过慢路径,SHARP 怎么把 reduction 下沉到交换机。Part 3 是动手实验,把原理落到可观察的操作上。

2. NCCL 与 NVIDIA 通信栈

NCCL 是训练侧最重要的 GPU collective 通信库。PyTorch DDP、FSDP、Megatron、Horovod 等框架中的 all-reduce、reduce-scatter、all-gather、broadcast,通常都由 NCCL 在底层完成。与 MPI 等传统库不同,NCCL 把每个 collective 实现在单个 GPU kernel 中,通信和计算操作在同一个 kernel 内完成,从而实现快速同步并最小化达到峰值带宽所需的资源开销。 NCCL 不单独存在,它位于 NVIDIA Magnum IO 通信栈里。Magnum IO 把存储 I/O、网络 I/O、网络内计算和 I/O 管理放在同一套性能优化框架下:NCCL 负责 collective,GPUDirect RDMA 负责跨节点 GPU 数据路径,SHARP/NVLS 负责把部分 reduction 下沉到网络或 NVSwitch fabric。
Four components of NVIDIA Magnum IO acceleration platform

Magnum IO 把 storage、network、in-network compute 和 I/O management 组合成 NVIDIA 的 I/O 加速栈;NCCL 属于其中的 Network I/O / collective 通信路径

NCCL architecture showing algorithms, channels, and transport layers

NCCL 架构:上层 collective 算法(Ring/Tree/PAT)通过 channel 抽象调用下层传输(NVLink P2P、SHM、NET/IB),每层可独立扩展

NCCL 的性能来自两层能力:上层是 collective 算法,下层是 NVIDIA 通信栈。

2.1 Collective 算法:Ring、Tree 与协议

最典型的 collective 是梯度 all-reduce。每张 GPU 先算出本地梯度,然后 collective 让所有 GPU 得到全局平均梯度。

2.1.1 Ring All-Reduce

2 节点 × 4 GPU/节点 为例,512 MB 梯度(BF16),NCCL 切成 8 块(每块 64 MB),在逻辑环上完成同步:
Ring all-reduce process for two nodes with four GPUs each, including reduce-scatter and all-gather phases

2 节点 × 4 GPU 的 Ring AllReduce:NCCL 根据拓扑构建逻辑环,节点内走 GPU 互联,跨节点边走 InfiniBand/RoCE

阶段一:Reduce-Scatter(7 步) — 让每张 GPU 最终持有一块全局归约后的 chunk。 初始状态下,每张 GPU 都持有完整 512 MB 本地梯度,也就是 8 个 chunk。每一步中,GPU 把一个 chunk 发给环上的下一跳,接收方把收到的数据累加到本地对应 chunk 上。7 步之后,GPU k 持有 chunk k 的全局 sum,已经包含 8 张 GPU 的贡献。 阶段二:All-Gather(7 步) — 把每张 GPU 持有的已归约 chunk 广播给所有 GPU。 All-gather 继续沿同一条 ring 路径传播已归约的 chunk,但不再做累加。7 步之后,所有 8 张 GPU 都持有完整的 512 MB 全局平均梯度。 理论通信量(每张 GPU):
Ring 的优势是带宽最优——每个 GPU 发送和接收的数据量都逼近理论下界。但它的弱点在延迟:步数随 GPU 数量线性增长,24,576 张卡就需要数万步,规模越大越不可接受。

2.1.2 Tree All-Reduce: Double Binary Tree

朴素 binary tree 能把步数压到 O(log N)——8 张卡只需 3 步。但朴素 tree 有带宽问题:根节点需要处理所有数据,而叶子节点大部分时间闲着。 NCCL 2.4 引入的 double binary tree(双互补二叉树)同时解决了这两个问题。核心思想来自 2009 年 MPI 领域的论文(Sanders, Speck & Träff):
Binary tree diagram using power-of-two pattern

单棵二叉树使用 power-of-two 模式构建,最大化节点局部性

在单棵二叉树中,半数或更多的 rank 是叶子节点,半数或更少是内部节点。关键洞察:构建第二棵互补树,让原来的叶子变成内部节点,原来的内部节点变成叶子。
Double complementary binary tree where leaves and nodes are flipped between trees

两棵互补二叉树——每个 rank 在一棵树中最多是内部节点,在另一棵树中是叶子

两棵树各处理一半数据。叠加后,每个 rank 最多接收一半数据两次、发送一半数据两次——和 ring 一样是带宽最优的,但延迟从 O(N) 降到了 O(log N)。 在 Summit 超算上 24,576 张 GPU 的实测数据验证了理论预测:
NCCL latency comparison on Summit supercomputer up to 24,576 GPUs

NCCL 在 Summit 上最多 24,576 张 GPU 的延迟——tree 在大规模下比 ring 延迟低最多 180 倍

NCCL bus bandwidth comparison on Summit up to 24,576 GPUs

NCCL 在 Summit 上的总线带宽——double binary tree 在大规模下仍能维持接近满带宽

实测中 tree 的带宽在跨 L3 交换机时略有下降(InfiniBand 路由算法与 tree 通信模式的匹配问题),但整体优势明显。与分层 2D ring(节点内 reduce-scatter → 跨节点 all-reduce → 节点内 all-gather)相比,2D ring 的优势随规模保持恒定,而 tree 的优势随 GPU 数量持续增长

2.1.3 不同算法适合不同场景

NCCL 会自动在两种算法间切换:大消息且 ring 带宽更高时用 ring,小/中消息或大规模集群时用 tree。大多数时候自动选择已经足够好。

2.1.4 Simple / LL / LL128 协议

算法(NCCL_ALGO)决定数据走什么拓扑路径,协议(NCCL_PROTO)决定数据在这条路径上怎么传输和同步。两者是独立的选择维度: 选择逻辑和 Ring/Tree 一致——大消息优化吞吐(Simple),小消息优化延迟(LL/LL128)。NCCL 会根据消息大小自动在协议之间切换,通常不需要手动干预。 算法 × 协议的组合构成 NCCL 的完整决策空间: 第一层决策是 NCCL_ALGO,决定走 Ring、Tree、CollNet 或 PAT 等拓扑路径。第二层决策是 NCCL_PROTO,决定用 Simple、LL 或 LL128 等传输协议。 例如,一个 512 MB 的梯度 all-reduce 可能走 Ring + Simple(大消息,集装箱海运走环形航线),而一个 8 KB 的 bias 梯度走 Tree + LL(小消息,电动车闪送走树形快递网络)。

2.1.5 NCCL Device API(2.28):去掉 CPU 中间人

传统 NCCL 流程中,每个 bucket ready 后需要 CPU 介入提交 NCCL kernel。NCCL 2.28 Device API 让 GPU kernel 直接发起通信,去掉 CPU round-trip。 这和 NVSHMEM 的思路一致——把通信的发起权从 CPU 下放到 GPU。传统 MPI 模型中 CPU 是通信的调度者,GPU 必须等 CPU 提交;GPU-initiated 模型中 GPU 直接操作网卡或对端显存,消除 CPU round-trip。
Comparison of MPI collective communication pattern vs NVSHMEM GPU-initiated communication

MPI 与 NVSHMEM 通信模式对比:MPI 依赖 CPU 调度 collective,NVSHMEM 让 GPU 直接发起细粒度通信

三种模式:

2.2 节点内互联:NVLink、NVSwitch 与 NVLS

节点内 GPU 之间通过 NVLink 和 NVSwitch 互联,提供远高于 PCIe 的带宽。
Node-internal GPU0 to GPU1 data path through NVLink and NVSwitch with no CPU in the data path

节点内一跳:GPU0 的 NCCL kernel 从 HBM 读出数据,经 NVLink / NVSwitch 到达 GPU1,GPU1 侧写入 recvbuff 并完成 reduce

CPU 完全不参与,整个过程在 GPU 硬件内完成。 NVLS 在 NVSwitch fabric 中做 multicast + reduction。利用 NVSwitch 的硬件 multicast 能力,一次写入广播到多个 GPU,同时在 switch 内完成部分归约。节点内生效。

PXN(PCIe × NVLink)

没有直连网卡的 GPU,通过 NVLink 借邻居 GPU 的网卡发数据(NCCL 2.12)。对 MoE all-to-all(消息聚合降 8 倍)和 Model Parallel 子通信器闭环收益最大。

Copy Engine Collectives

用 GPU 上专用的 Copy Engine 硬件代替 SM 做 NVLink 数据搬运。传统 NCCL 通信占用 SM/CTA 资源,和用户计算抢算力;CE collectives 实现 zero-SM 通信,让通信和计算完全不争资源。适用于只需数据搬运的 collective(AllGather、AlltoAll),不适用于需要计算的 AllReduce(reduce 阶段仍需 SM)。

2.3 跨节点互联:InfiniBand、RoCE 与 GPUDirect RDMA

RDMA 解决的是跨节点数据传输路径过长的问题。传统 socket 通信中,数据从应用层出发,要经过 Socket API → TCP/UDP → IP → NIC Driver 整个内核协议栈,到达对端后再逆序走一遍。每一层都意味着一次数据拷贝和上下文切换。 RDMA 的核心是 kernel bypass——应用直接和网卡通信,整个内核协议栈被绕过。这消除了三样东西:上下文切换、中间数据拷贝、协议处理开销。
RDMA vs Traditional Messaging: Socket-based path goes through full kernel network stack, RDMA bypasses kernel entirely

RDMA vs 传统 Socket 通信:传统方式数据要穿过完整的内核协议栈(Socket → TCP/UDP → IP → NIC Driver),RDMA 直接绕过内核,应用与网卡之间零拷贝通信

GPUDirect RDMA 在此基础上更进一步:不仅绕过内核,还绕过 CPU 内存——网卡直接读写 GPU 显存(HBM),数据路径变成 GPU HBM → NIC → 网络 → NIC → GPU HBM。
GPU-to-GPU direct data transfer with RoCE

GPUDirect RDMA:RDMA NIC 直接读写远端 GPU memory,避免数据经过 host CPU 和 system memory 中转

三条路径可以这样理解:
Comparison of TCP/IP, RDMA, and GPUDirect RDMA data paths across two nodes

TCP/IP、RDMA 与 GPUDirect RDMA 的数据路径对比:越往下,CPU 和内核网络栈参与越少

对 AI 训练来说,RDMA 的价值不只是”网络更快”,而是减少 CPU 参与,让 GPU 间大块数据传输更接近硬件极限。多节点 DDP、FSDP、张量并行、专家并行都会受益,尤其是梯度或参数分片很大时。

2.3.1 跨节点一跳的硬件数据路径

Cross-node GPUDirect RDMA path from GPU3 HBM through NICs and IB or RoCE network to GPU4 HBM, bypassing host memory

跨节点一跳:GPUDirect RDMA 让 NIC 通过 DMA 直接读写 GPU HBM,数据路径绕过 host memory 和内核网络栈

CPU 不在数据路径上——仅在初始化时注册内存区域(ibv_reg_mr / nvidia_peermem)。
⚠️ 注意: 如果 nvidia_peermem 模块没加载,或容器没暴露 /dev/infiniband/*,NIC 无法直接访问 GPU HBM。NCCL 会静默退回 CPU staging 路径:GPU3 → cudaMemcpy D2H → CPU 内存 → NIC → 网络 → NIC → CPU 内存 → cudaMemcpy H2D → GPU4。多了两次 PCIe 拷贝和 CPU 内存中转,延迟增加 2–3×,且业务代码不报错。

2.3.2 先分清层级:InfiniBand、RoCE 与以太网

这里最容易混淆的是把 InfiniBand、RoCE 和普通以太网当成同一层的三个名字。更准确的看法是:先看物理网络和网卡能力,再看上面跑的是 RDMA 还是 TCP/IP,最后看 NCCL 选择了哪个通信后端。 因此,InfiniBand 和 RoCE 都可以提供 RDMA 语义,只是承载网络不同:InfiniBand 是专用 IB fabric;RoCE 是在以太网上承载 RDMA。普通以太网本身通常走 TCP/IP,不提供 RDMA 快路径。对 NCCL 来说,真正重要的是它有没有成功走到 RDMA 后端,以及 GPUDirect RDMA 是否可用。 注意 RoCE 网卡和普通以太网卡在 Linux 层面可能有相同的设备名(eth0bond0),不能只靠接口名判断。RoCE 是否生效取决于网卡能力、网络配置、驱动/插件和权限。NCCL 日志里如果看到 NET/IB,通常说明它走到了 RDMA 路径;如果只看到 NET/Socket,即使底层是高速以太网,也是在走 TCP/socket 慢路径。 Bond 本身不是一种网络技术,而是 Linux 把多张物理网卡捆绑成一个逻辑接口的机制(提高带宽或做冗余),InfiniBand 和以太网卡都可以做 bond。 容器和 Kubernetes 环境中最容易出现”看起来能跑,但其实走慢路径”的问题。常见原因包括 /dev/infiniband 没有暴露给容器、nvidia_peermem 没加载、GID/权限不匹配、NCCL 选错网卡,或者 RoCE 网络没有正确配置。结果是 NCCL 悄悄退回 socket/TCP,吞吐明显下降,但业务代码不一定报错。
Bypassing CPU bottlenecks with direct connectivity between GPUs and NICs

直接连接 NIC 可以绕过 CPU 到 PCIe switch 的瓶颈,让 GPU 与 NIC 之间获得更完整的链路吞吐

2.4 网络内计算:SHARP

传统 all-reduce 中,每张 GPU 要收发 2(N-1)/N × data_size 的数据,reduction 运算(求和/平均)在 GPU 端点完成。SHARP 的思路是:把 reduction 下沉到 InfiniBand Quantum 交换机里做,端点只收最终结果。 GPU 把本地梯度发到交换机,交换机直接求和后把结果返回。需要交换机 firmware 支持 + Aggregation Manager 进程。跨节点生效。节点内的对应方案是 NVLS(见上文节点内互联)。 效果:大规模训练(数百~千卡)下 all-reduce 加速 2×–5×,且不需要改用户代码——NCCL 检测到硬件支持后自动启用。

2.5 管理与诊断

2.5.1 排查入口

常用排查入口: 调优 NCCL 时不要把环境变量当作固定模板。先回答三个问题:通信是否真的成为瓶颈,NCCL 是否走了最快链路,collective 是否已经和计算重叠。只有这些问题有明确证据后,算法和环境变量调优才有意义。 验证跨节点路径时看四类信号:
  • lsmod | grep nvidia_peermem:确认 GPUDirect RDMA 相关模块已加载。
  • NCCL_DEBUG=INFO:确认日志中出现 NET/IB,而不是只走 socket。
  • RDMA perftest:使用 GPU buffer,例如带 --use_cuda 的测试。
  • Nsight Systems / 网卡计数器:确认网络传输和 GPU 计算能够并行发生。

2.5.2 生命周期原则

不要在每个 step 里创建和销毁 NCCL communicator;初始化和拓扑发现应放在训练或服务启动阶段,并在生命周期内复用。不要忽略 NCCL warning——很多 warning 指向的不是”日志噪音”,而是版本不匹配、端口耗尽、网络 fallback、异步错误或 rank 卡死。

2.6 通信优化的基本判断

通信优化先看一个问题:GPU 是在计算,还是在等通信。如果 GPU 大部分时间在跑 kernel,通信优化收益有限;如果 timeline 中有明显的 all-reduce、copy、socket 或 network wait,通信路径和同步策略就会成为主要瓶颈。 先按现象建立判断框架。看到 GPU 等待后,不要直接套环境变量模板,而是先判断等待来自框架通信模式、同步频率、传输数据量、跨节点路径,还是 collective 拓扑。不同现象对应的优化手段不同。 框架通信模式问题。 单机多卡训练中,如果仍在使用 DataParallel,GPU0 往往会更忙。DataParallel 是单进程控制多张 GPU,主 GPU 负责 scatter 输入、gather 输出和聚合梯度,容易形成主卡瓶颈。DistributedDataParallel 通常是一进程一 GPU,用 NCCL 做梯度 all-reduce,并在 backward 过程中按 bucket 发起异步通信,因此更适合多 GPU 训练。 通信启动太晚。 朴素做法是所有 backward 结束后再统一 all-reduce,这会把通信时间直接加到每轮迭代上。更好的做法是把梯度切成 bucket,某个 bucket 一准备好就发起 NCCL all-reduce,同时继续计算后续层的梯度。
Overlapping host-to-device and device-to-host communication with computation on multiple CUDA streams

通信与计算重叠:同步模式会让 GPU 等待 all-reduce;DDP bucket 可以在 backward 过程中提前启动通信

需要避免的常见同步点包括 torch.cuda.synchronize()tensor.item()。前者会等待当前设备上的工作全部完成,后者会把 GPU 标量搬回 CPU,也可能触发同步。计时或日志统计应尽量放在迭代末尾,避免打断重叠流水线。 同步频率太高。 如果每个 step 都有固定通信开销,可以用梯度累积减少同步次数。多个 micro-batch 先在本地累积梯度,再做一次 all-reduce,相当于用更多本地计算换更少的跨 GPU 同步。代价是有效 batch size 变大,可能需要重新调学习率和内存预算。 Before — 每步同步:
After — 梯度累积,4 步同步一次:
传输数据量太大。 如果网络链路持续打满,压缩、量化或稀疏化可以减少每次通信的数据量。这类方法直接作用在”传多少”上,但会改变梯度信息,需要关注收敛稳定性和最终精度。 跨节点路径太慢。 如果跨节点吞吐远低于预期,先确认通信后端是否正确。GPU 训练应优先使用 NCCL;如果误用 Gloo,或者 NCCL 因配置问题退回 socket/TCP,通信就会变成 CPU 和内核网络栈主导。 collective 拓扑不匹配。 如果 GPU 数量增加后 all-reduce 扩展性变差,问题可能不在单条链路,而在 collective 算法和硬件拓扑没有匹配。还要检查是否存在 straggler:一个 rank 的 CPU 线程、网卡或远端内存访问变慢,就会拖住整个 collective。 实践中的第一步仍然是看 profiler。Nsight Systems 或 PyTorch profiler 里如果看到计算和通信串行排列,就先检查 DDP bucket、同步点和通信后端;如果通信已经被计算覆盖,再继续调 NCCL 算法或网络参数,收益通常会变小。

2.7 推理通信:从 NCCL 到 NIXL

训练通信的核心是 NCCL collective(all-reduce 同步梯度),推理通信的核心是点对点 KV cache 传输——没有 backward,没有梯度,不需要 all-reduce。

2.7.1 为什么要 PD 分离

传统推理系统把 prefill(计算 KV cache)和 decode(逐 token 生成)放在同一组 GPU 上,两个阶段互相干扰: 放在一起时,prefill 抢算力导致 decode 变慢,为保 decode 又让 prefill 排队。DistServe(OSDI 2024)和 Splitwise(Microsoft)的解法:拆到不同 GPU 上,各自独立优化资源和并行策略。 DistServe 在同样 GPU 预算下服务 7.4x 更多请求。
Disaggregated inference with prefill workers sending KV cache blocks point-to-point to decode workers across a high-speed network

PD 分离推理的数据路径:prefill worker 产出 KV cache,通过点对点异步传输交给 decode worker,必要时在 GPU HBM、CPU DRAM 和 NVMe 之间分层移动

2.7.2 NIXL:推理的数据搬运层

PD 分离后,KV cache 必须从 prefill GPU 搬到 decode GPU。搬运目标可能是 GPU HBM、CPU DRAM、甚至 NVMe SSD(显存不够时把冷请求的 KV cache 换出,用户回来再换入)。 NIXL(NVIDIA Inference Xfer Library)是推理场景的数据搬运抽象层,屏蔽内存类型和传输协议的复杂性,自动选择最快路径:NVLink > IB/RoCE > PCIe > GDS/NVMe。 三者关系:DistServe/Splitwise 定义了为什么要拆(架构设计),NIXL 解决拆了之后数据怎么搬(通信实现),NVIDIA Dynamo 把两者整合成完整的推理框架。
注意区分: NVIDIA Dynamo(github.com/ai-dynamo/dynamo)是分布式推理 serving 框架,负责 prefill/decode worker 调度和 NIXL 数据搬运编排。TorchDynamo(torch.compile 的前端)是 PyTorch 编译器组件,负责捕获 Python 计算图并交给 TorchInductor 优化。两者完全无关,只是撞名。

2.8 其他通信组件

GDS 面向 GPU 与 NVMe 存储之间的数据路径,更多属于下一章存储 I/O 优化。

3. 实验

3.1 bucket_cap_mb 调参(MI300X × 8,XGMI 896 GB/s,ROCm 6.2)

下面是实验用的 ddp_bench.py 脚本,构建一个可配置层数和 hidden size 的 Transformer 模型,用 DDP 运行若干 step 并记录平均耗时:
用下面的脚本扫描不同的 bucket_cap_mb 配置:

3.1.1 小模型(88 MB,6 层 2048-hidden)

小模型最优 bucket = 50 MB,远大于默认值。模型只有 88 MB,bucket 太小导致通信启动次数多、每次数据量小,启动开销占比大。

3.1.2 大模型(1628 MB,24 层 4096-hidden)

大模型最优 bucket = 5 MB,远小于默认值。模型大、backward 时间长,小 bucket 能更早发起通信,overlap 更充分。

3.1.3 Profiler 通信/计算分解(大模型 1628 MB)

通信和计算接近 1:1,梯度同步已是显著瓶颈,这解释了为什么小 bucket(更早 overlap)在大模型上表现最好。

核心思想:先判断等待在哪里

  1. 如果 GPU 在等 backward 后的梯度同步,优先看 DDP、bucket 和通信计算重叠。
  2. 如果跨节点通信慢,优先确认 RDMA / GPUDirect RDMA 是否走通。
  3. 如果 collective 放大到多节点,重点看 NCCL 是否正确利用 NVLink、NVSwitch、InfiniBand/RoCE、SHARP 等通信栈。
  4. 如果只有少数 rank 慢,按 straggler 排查 CPU affinity、NUMA、网卡和远端内存访问。

FAQ

PD 分离架构下,prefill 和 decode 运行在不同 GPU 上。请求交给 decode worker 继续生成时,不仅要搬 KV cache,还要搬请求的元信息——当前生成到第几个 token、sampling 参数、stop conditions、已生成的 token ids。没有这些状态,decode worker 不知道从哪接着生成。
Bucket 是 DDP 对参数梯度的预分组(默认每桶上限 25 MB)。单个 bucket 的 all-reduce 确实是 collective——所有 rank 必须共同完成。但 DDP 不是对整个模型做一次 all-reduce,而是按 bucket 拆成多次独立的 all-reduce。Backward 从最后一层往前算,最后几层的梯度凑满一个 bucket 后,这个 bucket 的 all-reduce 立即启动,同时 backward 继续计算前面的层。通信和计算流水线式重叠。这就是为什么 bucket 大小很重要——太大则通信启动太晚,太小则启动次数多、开销大。Bucket 和层不是一一对应的。一个层通常有多个参数(如 FFN 的 W1b1W2b2),这些参数可能分属不同 bucket;反过来,多个小层的参数也可能合并进同一个 bucket。DDP 按参数的 reverse 顺序依次填充 bucket,填满一个再开下一个,不关心层边界。但单个参数的梯度不会被拆到多个 bucket——如果某个参数超过 bucket_cap_mb,它独占一个 bucket。Autograd 对每个参数一次性算完整个 grad tensor,算完才标记 ready;bucket 内所有参数都 ready 后才触发 all-reduce。所以不存在”梯度算了一半就被 all-reduce 走”的情况。
构造方法来自 [Sanders, Speck & Träff (2009)]https://www.sciencedirect.com/science/article/abs/pii/S0167819109000957),对 2^k 个 rank 有确定性构造——把第一棵树的叶子和内部节点角色翻转即可。两棵树各走一半数据,叠加后每个 rank 收发量与 ring 一致,接近带宽最优。唯一的例外是 root 在单棵树中负载略低,但因为两棵树互补,整体接近满带宽。 具体也可以看NV的这个博客 [Massively Scale Your Deep Learning Training with NCCL 2.4](https://developer.nvidia.com/blog/massively-scale-deep-learning-training-nccl-2-4/#ref3)
没有固定卡数阈值,取决于消息大小 × GPU 数量。大消息 + 少量 GPU 倾向 Ring(带宽打满,步数少时延迟不是问题);小消息或大规模集群倾向 Tree(O(log N) 步,启动快)。单机八卡或小规模多机(2-4 台),Ring 通常够用。64+ 卡时 Tree 的 O(log N) 优势开始显现。实际中 NCCL 会对不同 collective 自动选择——大 bucket 走 Ring 吃带宽,小 bucket 走 Tree 低延迟,通常不需要手动指定。
核心数据路径(梯度/参数/KV cache 搬运)已不需要 CPU。CPU 仍负责:初始化与拓扑发现(communicator 创建、内存注册)、控制面操作(bootstrap、barrier、health check,走 TCP,数据量极小)、传统模式下每次 collective 的 CPU 端提交(NVSHMEM 的 device-side API 或 CUDA Graphs 可消除这一步)、数据加载与预处理。DDP 的 autograd hook 也在 CPU 上,但它做的是轻量调度(标记梯度 ready、检查 bucket、提交 NCCL 调用),微秒级开销。在静态图场景下(torch.compile / CUDA Graphs),整个 backward + NCCL 调用可以被捕获为 graph 一次性 replay,消除每轮迭代的 CPU 提交开销。但动态图(eager mode)下 autograd 图遍历本身依赖 CPU 调度,无法消除。实际中这部分开销是微秒级,不是性能瓶颈。数据面是 NVLink/InfiniBand 上实际搬运 tensor 的路径,追求高带宽;控制面是协调”谁跟谁通信、何时开始”的信令,走 TCP,KB 级,不是瓶颈。
不一样。Registered(pinned)memory 通过 cudaHostRegister() 注册后页锁定——OS 不会 swap 到磁盘,GPU DMA 和 NIC 可以直接访问,物理地址已知,硬件可以直接 DMA。普通 CPU 内存(pageable)随时可能被换出,需要先拷贝到 pinned buffer 再传输,多一次拷贝。
不需要,默认行为。NCCL 根据消息大小、GPU 数量和拓扑自动选 Ring/Tree 和 Simple/LL/LL128。手动覆盖用 NCCL_ALGONCCL_PROTO,但一般不建议,除非 profiler 显示自动选择不是最优。DDP 层面的 bucket_cap_mb 需要手动调。
不建议。NCCL bootstrap(进程发现和初始化握手)必须走 TCP。应该做的是确保数据面不走 TCP:正确配置 RDMA 环境,用 NCCL_DEBUG=INFO 确认日志中数据传输走 NET/IB 而非 NET/Socket。如果看到后者,排查 RDMA 配置而不是禁 TCP。
Collective 让所有 rank 等最慢的那个,表面上”全卡了”。排查方法:Per-rank profiling(Nsight Systems 对每个 rank 分别采集 timeline,比较同一个 all-reduce 中各 rank 的到达时间);网卡计数器(ethtool -S / InfiniBand perfquery,看错误数、丢包、重传);CPU affinity 检查(numactl --show 确认进程绑在正确的 NUMA node);排除法(每次同一个 rank 慢 → 查硬件,随机漂移 → 网络拥塞或 OS 抖动)。线上不方便挂 profiler 时,可以用轻量方式定位:在 all-reduce 前后各插一个 torch.cuda.synchronize() + time.perf_counter(),每个 rank 把耗时写到各自的日志文件,跑几十个 step 后离线比较哪个 rank 的 skew(最大值 - 最小值)最大。 也可以不改代码,直接看 nvidia-smi dmon -s u 的 GPU 利用率波动——straggler 通常是计算慢的那个(GPU util 持续高),而其他节点因为提前算完、在 collective 里等 straggler,会出现周期性的 util 低谷。如果某个节点每轮迭代都比别人多一段 idle gap,说明它在等别人(它不是 straggler);如果某个节点几乎没有 idle gap 但别人都在等,它就是 straggler。网卡侧用 ethtool -S 或 InfiniBand perfquery 看重传和错误计数,不需要 profiler 也能发现链路级问题。
Channel 是 NCCL 的并行通信流水线。更多 channel = 更高带宽利用,但占更多 SM 资源。默认值通常够用。调之前先看 Nsight Systems 确认瓶颈在通信带宽——带宽没打满且 SM 有余量可以调大;通信 kernel 抢 SM 影响计算,应调小或改用 Copy Engine Collectives。注意:NCCL 2.21+ 将 NCCL_MAX_NCHANNELS 改名为 NCCL_MAX_CTAS(对应 NCCL_MIN_CTAS)。旧名称在部分版本仍兼容,但新部署建议用新名称。
不需要,profiling 本身有 5-15% 开销。发现异常后下一次运行中加 profiler,采集几个 step 即可(nsys profile --duration 30 或 PyTorch Profiler 指定步数)。如果问题是卡死(hang)而非”慢”,需要不同工具:NCCL_DEBUG=WARN 看异步错误、py-spy 看调用栈、nvidia-smi 看 GPU 利用率是否归零。
不是。即使 NCCL group 只有 2 个节点,做的仍然是 collective 语义——双方必须同步到达同步点。NIXL 是异步点对点语义:发送方 put 数据到接收方指定地址,不需要接收方同时参与;支持多种内存类型(GPU HBM / CPU DRAM / NVMe)任意组合;worker 动态加入离开;自动选最快路径。NCCL 为训练设计(固定 rank、collective),NIXL 为推理设计(动态 worker、point-to-point)。
NCCL 假设固定数量的 rank、固定的 communicator、重复执行相同的 collective,不匹配推理场景(动态 worker 池、点对点 KV cache 传输、连接随请求建立/销毁)。这正是 NIXL 存在的原因——专为推理场景的异步点对点传输设计。
先澄清:训练数据的切分只发生在数据并行维度。其他并行策略(TP、PP、EP)切的是模型,不是数据——同一份输入数据会被所有参与计算的 GPU 共同处理,只是每张卡负责模型的不同部分。具体来说,训练数据在各策略下的去向:
  • 数据并行(DDP/FSDP/ZeRO):每张 GPU 拿到不同的 mini-batch 子集,独立做 forward/backward,最后同步梯度。这是唯一真正”切数据”的维度。
  • 张量并行(TP):同一份输入数据广播到 TP group 内所有 GPU,每张卡只算权重矩阵的一部分,结果再 all-reduce/all-gather 拼回来。数据没切,切的是矩阵乘法。
  • 流水线并行(PP):同一份 mini-batch 被切成若干 micro-batch,依次流过各 stage。每个 stage 看到完整的 micro-batch,只是在时间上错开。数据没在空间上切分,而是在时间上流水。
  • 专家并行(EP):每个 token 被路由到少数 expert,all-to-all 把 token 发给对应 GPU。这里切的是”哪些 token 给哪个 expert”,由 router 动态决定,不是预先均分。
所以”能否不按单卡均分”——在数据并行维度,PyTorch DistributedSampler 默认均分,但可以自定义不均匀采样(比如按 GPU 显存大小分配不同 batch size)。代价是梯度聚合时需要按实际 sample 数加权平均,否则大 batch 的 GPU 对梯度贡献被低估。

参考资料