> ## Documentation Index
> Fetch the complete documentation index at: https://se7en.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

# 第四章：分布式网络通信调优

## 录屏回看

<iframe src="https://player.bilibili.com/player.html?bvid=BV1HTEK6qEak&autoplay=0&high_quality=1" width="100%" height="480" scrolling="no" frameBorder="0" allowFullScreen />

## 本章概要

分布式 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** 的直觉——不需要每个人跟所有人通信，只需要跟邻居交换，经过足够多的步骤，所有人都能拿到全局结果。

<Frame caption="Ring all-reduce 的直觉：8 个参与者只和相邻节点交换数据，经过多轮传递后每个节点都得到同一个全局结果">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-ring-allreduce-intuition.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=cd9798b87f2e70760aa07855dd4ea7cf" alt="Eight participants arranged in a ring with neighbor-to-neighbor arrows, illustrating ring all-reduce intuition" width="1672" height="941" data-path="books/ai-systems-performance-engineering/images/ch04-generated-ring-allreduce-intuition.png" />
</Frame>

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

| 考试类比           | 分布式训练                  |
| -------------- | ---------------------- |
| 每位老师批改 100 份试卷 | 每张 GPU 处理一份 mini-batch |
| 每位老师算出的平均分     | 每张 GPU 算出的本地梯度         |
| 合算全班平均分        | all-reduce 得到全局平均梯度    |
| 8 个人排成一圈传递     | 8 张 GPU 组成逻辑环通信        |
| 每人最终拿到相同的全班平均分 | 每张 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 又如何承载真正的数据传输。

<Frame caption="一个梯度 tensor 的直觉旅程：先产生本地梯度，进入 DDP bucket，bucket 就绪后发起同步，数据经高速链路在 GPU 间交换，最后所有 GPU 得到相同的全局平均梯度">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-gradient-intuitive-journey.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=7a869c9cc31f88ea0fa2f37b11ea769a" alt="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" width="1774" height="887" data-path="books/ai-systems-performance-engineering/images/ch04-generated-gradient-intuitive-journey.png" />
</Frame>

关键问题是第 ② 步：梯度同步要花时间，GPU 在等同步的时候就没在算东西。聪明的做法是**边算边同步**——不等所有层的梯度都算完，只要某个 bucket 里的梯度都 ready，就立刻对这个 bucket 发起 all-reduce，同时继续计算后续层的梯度。

<Frame caption="DDP 通信计算重叠：backward 继续推进时，已就绪的 bucket 可以先进入 AllReduce；触发条件是 bucket ready，而不是单个梯度 ready">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-ddp-overlap-timeline.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=ba053e5f2f31a4fb749423c9284cbbe0" alt="DDP timeline showing backward computation overlapping with all-reduce launched after buckets become ready" width="1754" height="897" data-path="books/ai-systems-performance-engineering/images/ch04-generated-ddp-overlap-timeline.png" />
</Frame>

这就是 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（主卡瓶颈）：**

```python theme={null}
import torch
import torch.nn as nn

model = MyModel().cuda()
# 单进程控制多卡，GPU0 负责 scatter/gather/梯度聚合
# → GPU0 比其他卡更忙，其他卡等 GPU0 完成聚合
model = nn.DataParallel(model, device_ids=[0, 1, 2, 3])

for data, target in dataloader:
    data, target = data.cuda(), target.cuda()
    output = model(data)        # scatter 到 4 卡，gather 回 GPU0
    loss = criterion(output, target)
    loss.backward()             # 梯度全部回 GPU0 聚合
    optimizer.step()
    optimizer.zero_grad()
```

**After — DistributedDataParallel（NCCL all-reduce，无主卡瓶颈）：**

```python theme={null}
import torch
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP

dist.init_process_group("nccl")
local_rank = int(os.environ["LOCAL_RANK"])
torch.cuda.set_device(local_rank)

model = MyModel().cuda(local_rank)
# 每个进程控制一张卡，梯度同步由 NCCL all-reduce 完成
# bucket 机制让通信和 backward 计算重叠
model = DDP(model, device_ids=[local_rank])

for data, target in dataloader:
    data = data.cuda(local_rank)
    target = target.cuda(local_rank)
    output = model(data)
    loss = criterion(output, target)
    loss.backward()             # 梯度按 bucket 异步 all-reduce
    optimizer.step()
    optimizer.zero_grad()
```

| 指标      | DataParallel            | DDP                      |
| ------- | ----------------------- | ------------------------ |
| 梯度同步    | GPU0 gather + broadcast | NCCL all-reduce（无主卡）     |
| 通信/计算重叠 | 无                       | bucket 机制，backward 中异步发起 |
| 进程模型    | 单进程多线程（GIL 限制）          | 多进程（无 GIL，每进程一卡）         |
| 扩展性     | 单机                      | 单机 + 多机                  |

#### 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）。

<Frame caption="DDP 分层架构：Python 入口层（distributed.py）→ C++ 梯度同步与广播（reducer.h / comm.h）→ 通信后端抽象（ProcessGroup.hpp）→ 具体实现（NCCL / Gloo / MPI / RCCL / XCCL）">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-figure-ddp-architecture.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=89a7a1b8631f89ea9378be0390fdd290" alt="DDP architecture: distributed.py → reducer.h + comm.h → ProcessGroup → NCCL/Gloo/MPI/RR/XCCL" width="560" height="368" data-path="books/ai-systems-performance-engineering/images/ch04-figure-ddp-architecture.png" />
</Frame>

<Frame caption="DDP 梯度同步：参数按 reverse 顺序分到 bucket 里（bucket1 = grad0+grad1，bucket0 = grad2+grad3），每个 bucket 就绪后独立发起 allreduce，两个 Process 的 bucket 结构完全对称">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-figure-ddp-bucket-allreduce.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=5eaf096978dabdbfcf0f73fdae2dd1b6" alt="DDP bucket allreduce: Process 0 and Process 1 each have params grouped into buckets, with allreduce between corresponding buckets" width="1458" height="688" data-path="books/ai-systems-performance-engineering/images/ch04-figure-ddp-bucket-allreduce.png" />
</Frame>

<Frame caption="DDP bucket 机制细节：梯度先进入预先划分的 bucket，bucket 内所有梯度 ready 后，才对这个 bucket 发起 AllReduce">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-gradient-bucket-before-allreduce.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=23fa8517db5cd541d176c36a4ba971cb" alt="Gradient blocks are grouped into buckets first, then completed buckets launch all-reduce before optimizer step" width="1672" height="941" data-path="books/ai-systems-performance-engineering/images/ch04-generated-gradient-bucket-before-allreduce.png" />
</Frame>

#### 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 上。

<Frame caption="DDP backward 全流程：autograd hook 只标记梯度 ready，reducer 检查预先划分的 bucket；bucket 完成后按 bucket index 顺序提交 AllReduce 到通信流">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-ddp-backward-cpu-gpu-flow.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=5de889a6ae07c4c306a4cf0d2c477ccc" alt="DDP backward flow with CPU reducer hooks, GPU compute stream, and GPU communication stream for bucket all-reduce" width="1672" height="941" data-path="books/ai-systems-performance-engineering/images/ch04-generated-ddp-backward-cpu-gpu-flow.png" />
</Frame>

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

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

到这里，我们已经完整追踪了一个梯度 tensor 的旅程——从 `loss.backward()` 到 DDP bucket 到 NCCL all-reduce。但分布式系统中需要通信的远不止梯度。

| 传输内容              | 出现场景                                       | 通信模式                        | 大小量级                             |
| ----------------- | ------------------------------------------ | --------------------------- | -------------------------------- |
| **梯度** ← 刚才追踪的    | DDP / ZeRO 训练                              | all-reduce / reduce-scatter | 整个模型参数量（GB 级）                    |
| **参数分片**          | FSDP — 每层 forward 前 all-gather 完整权重        | all-gather                  | 单层参数（数百 MB）                      |
| **激活值**           | Pipeline parallelism — stage 间传递           | point-to-point send/recv    | per-microbatch per-layer（MB\~GB） |
| **KV cache**      | Disaggregated inference — prefill → decode | point-to-point 异步传输         | per-request（MB\~百 MB，长上下文可达 GB）  |
| **Expert tokens** | MoE all-to-all — token dispatch 和 combine  | all-to-all                  | 动态，取决于路由分布                       |
| **控制信号**          | NCCL bootstrap、barrier、health check        | TCP / gRPC                  | 极小（KB 级）                         |

几个观察：

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。

<Frame caption="Magnum IO 把 storage、network、in-network compute 和 I/O management 组合成 NVIDIA 的 I/O 加速栈；NCCL 属于其中的 Network I/O / collective 通信路径">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-figure-4-2-magnum-io.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=11f4636e7504a3bc5ea2be751a4ff6eb" alt="Four components of NVIDIA Magnum IO acceleration platform" width="1075" height="491" data-path="books/ai-systems-performance-engineering/images/ch04-figure-4-2-magnum-io.png" />
</Frame>

<Frame caption="NCCL 架构：上层 collective 算法（Ring/Tree/PAT）通过 channel 抽象调用下层传输（NVLink P2P、SHM、NET/IB），每层可独立扩展">
  <img src="https://developer-blogs.nvidia.com/wp-content/uploads/2020/10/nccl-architecture-1.png" alt="NCCL architecture showing algorithms, channels, and transport layers" />
</Frame>

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

| 层次            | 组件                                   | 作用                        |
| ------------- | ------------------------------------ | ------------------------- |
| Collective 算法 | Ring、Tree、CollNet、PAT                | 根据消息大小和 GPU 数量组织通信        |
| 节点内互联         | NVLink、NVSwitch、NVLS                 | 提供 GPU 间高带宽低延迟链路          |
| 跨节点互联         | InfiniBand、RoCE、GPUDirect RDMA       | 让 GPU 跨服务器高速交换数据          |
| 网络内计算         | SHARP                                | 在交换机中做部分 reduction，降低端点压力 |
| 管理与诊断         | NCCL 日志、拓扑文件、UFM/NetQ、Nsight Systems | 验证链路、算法和瓶颈                |

### 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），在逻辑环上完成同步：

<Frame caption="2 节点 × 4 GPU 的 Ring AllReduce：NCCL 根据拓扑构建逻辑环，节点内走 GPU 互联，跨节点边走 InfiniBand/RoCE">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-ring-allreduce-process.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=6e0f9d394ef680254c81b4bfe976461c" alt="Ring all-reduce process for two nodes with four GPUs each, including reduce-scatter and all-gather phases" width="1672" height="941" data-path="books/ai-systems-performance-engineering/images/ch04-generated-ring-allreduce-process.png" />
</Frame>

**阶段一：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）：**

```
单向发送 = 2 × (N-1)/N × data_size
         = 2 × (7/8) × 512 MB
         = 896 MB

瓶颈链路耗时（假设跨节点 IB 是瓶颈）：
  Ring 中有 2 条跨节点链路（GPU3→GPU4 和 GPU7→GPU0）
  有效跨节点带宽 ≈ 50 GB/s per direction（单 NIC）
  每步跨节点传输 64 MB：64 MB / 50 GB/s ≈ 1.3 ms
  14 步（7 reduce-scatter + 7 all-gather）中有 14 次跨节点跳
  但 pipeline 化后，耗时 ≈ 14 × 1.3 ms ≈ 18 ms（理论下界）
```

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）：

<Frame caption="单棵二叉树使用 power-of-two 模式构建，最大化节点局部性">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-nccl-btree.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=4117de74535dce418cf049fbab8b0de3" alt="Binary tree diagram using power-of-two pattern" width="600" height="497" data-path="books/ai-systems-performance-engineering/images/ch04-nccl-btree.png" />
</Frame>

在单棵二叉树中，半数或更多的 rank 是叶子节点，半数或更少是内部节点。关键洞察：**构建第二棵互补树，让原来的叶子变成内部节点，原来的内部节点变成叶子。**

<Frame caption="两棵互补二叉树——每个 rank 在一棵树中最多是内部节点，在另一棵树中是叶子">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-nccl-double-btree.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=e7d65120d633e919d8540496f79a88d8" alt="Double complementary binary tree where leaves and nodes are flipped between trees" width="800" height="443" data-path="books/ai-systems-performance-engineering/images/ch04-nccl-double-btree.png" />
</Frame>

两棵树各处理一半数据。叠加后，每个 rank 最多接收一半数据两次、发送一半数据两次——**和 ring 一样是带宽最优的**，但延迟从 O(N) 降到了 O(log N)。

| 性质             | Ring        | Double Binary Tree |
| -------------- | ----------- | ------------------ |
| **延迟**         | O(N) — 线性增长 | O(log N) — 对数增长    |
| **带宽**         | 最优          | 最优（与 ring 相同）      |
| **8 卡延迟**      | 14 步        | \~6 步              |
| **24,576 卡延迟** | \~49,150 步  | \~30 步             |

在 Summit 超算上 24,576 张 GPU 的实测数据验证了理论预测：

<Frame caption="NCCL 在 Summit 上最多 24,576 张 GPU 的延迟——tree 在大规模下比 ring 延迟低最多 180 倍">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-nccl-summit-latency.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=ef0e3241c7519201da98f6b3b5189997" alt="NCCL latency comparison on Summit supercomputer up to 24,576 GPUs" width="597" height="512" data-path="books/ai-systems-performance-engineering/images/ch04-nccl-summit-latency.png" />
</Frame>

<Frame caption="NCCL 在 Summit 上的总线带宽——double binary tree 在大规模下仍能维持接近满带宽">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-nccl-summit-bw.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=17c29e524975579694a92464e33fe008" alt="NCCL bus bandwidth comparison on Summit up to 24,576 GPUs" width="646" height="501" data-path="books/ai-systems-performance-engineering/images/ch04-nccl-summit-bw.png" />
</Frame>

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

#### 2.1.3 不同算法适合不同场景

| 算法                     | 适合场景           | 直觉                       |
| ---------------------- | -------------- | ------------------------ |
| **Ring**               | 大消息、高带宽场景      | 把大张量切块，持续喂满链路            |
| **Tree**               | 小消息、低延迟场景      | 通信步数少，启动快                |
| **CollNet / CollTree** | 多节点训练          | 节点内和节点间分层通信              |
| **PAT**                | 大规模 collective | 兼顾 tree 的低延迟和 ring 的带宽利用 |

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

#### 2.1.4 Simple / LL / LL128 协议

算法（`NCCL_ALGO`）决定数据走什么拓扑路径，协议（`NCCL_PROTO`）决定数据在这条路径上怎么传输和同步。两者是独立的选择维度：

| 协议                   | 核心思路                                              | 比喻                       |
| -------------------- | ------------------------------------------------- | ------------------------ |
| **Simple**           | 大块传输，GPU 线程做 memory copy，高吞吐                      | 海运集装箱——装满再发，单次运量大        |
| **LL** (Low Latency) | 每 4 字节数据附带 4 字节 flag 做内联同步，无需额外同步原语，延迟极低但带宽浪费 50% | 电动车闪送——每个小件立刻发，零等待但运力浪费大 |
| **LL128**            | 每 128 字节附带少量 flag，在吞吐和延迟之间折中                      | 小箱快发——介于集装箱和闪送之间         |

选择逻辑和 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。

<Frame caption="MPI 与 NVSHMEM 通信模式对比：MPI 依赖 CPU 调度 collective，NVSHMEM 让 GPU 直接发起细粒度通信">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-mpi-nvshmem-comparison.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=25ef59e32301f8d898b1ff5b53007475" alt="Comparison of MPI collective communication pattern vs NVSHMEM GPU-initiated communication" width="706" height="314" data-path="books/ai-systems-performance-engineering/images/ch04-mpi-nvshmem-comparison.png" />
</Frame>

三种模式：

| 模式                                | 路径                        | 适用            |
| --------------------------------- | ------------------------- | ------------- |
| **LSA**（Load/Store Accessible）    | CUDA P2P 直接读写对端显存         | 节点内 NVLink 设备 |
| **Multimem**                      | NVLink SHARP 硬件 multicast | 节点内 reduction |
| **GIN**（GPU-Initiated Networking） | GPU 直接操作网卡，不经 CPU         | 跨节点通信         |

### 2.2 节点内互联：NVLink、NVSwitch 与 NVLS

节点内 GPU 之间通过 NVLink 和 NVSwitch 互联，提供远高于 PCIe 的带宽。

<Frame caption="节点内一跳：GPU0 的 NCCL kernel 从 HBM 读出数据，经 NVLink / NVSwitch 到达 GPU1，GPU1 侧写入 recvbuff 并完成 reduce">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-nvswitch-hop.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=5c7479643c6c74f35a6168d9fa10bf23" alt="Node-internal GPU0 to GPU1 data path through NVLink and NVSwitch with no CPU in the data path" width="1672" height="941" data-path="books/ai-systems-performance-engineering/images/ch04-generated-nvswitch-hop.png" />
</Frame>

| 步骤             | 发生了什么                                                  | 耗时                             |
| -------------- | ------------------------------------------------------ | ------------------------------ |
| NCCL kernel 启动 | GPU0 的 SM 执行 NCCL 内部 copy+reduce kernel                | \~μs                           |
| NVLink 传输      | 数据从 GPU0 NVLink Engine → NVSwitch → GPU1 NVLink Engine | 64 MB / 900 GB/s ≈ **0.07 ms** |
| 远端 reduce      | GPU1 SM 执行 `recvbuff[i] += incoming[i]`                | \~μs                           |

CPU 完全不参与，整个过程在 GPU 硬件内完成。

#### NVLS（NVLink SHARP）

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**——应用直接和网卡通信，整个内核协议栈被绕过。这消除了三样东西：上下文切换、中间数据拷贝、协议处理开销。

<Frame caption="RDMA vs 传统 Socket 通信：传统方式数据要穿过完整的内核协议栈（Socket → TCP/UDP → IP → NIC Driver），RDMA 直接绕过内核，应用与网卡之间零拷贝通信">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-rdma-yuque-diagram.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=40aab96a3f354429f36f6a0f6f85e6a7" alt="RDMA vs Traditional Messaging: Socket-based path goes through full kernel network stack, RDMA bypasses kernel entirely" width="1948" height="1110" data-path="books/ai-systems-performance-engineering/images/ch04-rdma-yuque-diagram.png" />
</Frame>

GPUDirect RDMA 在此基础上更进一步：不仅绕过内核，还绕过 CPU 内存——网卡直接读写 GPU 显存（HBM），数据路径变成 GPU HBM → NIC → 网络 → NIC → GPU HBM。

<Frame caption="GPUDirect RDMA：RDMA NIC 直接读写远端 GPU memory，避免数据经过 host CPU 和 system memory 中转">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-figure-4-3-gpudirect-rdma.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=f97ce571477f64dcc2e65b49e51a6f4f" alt="GPU-to-GPU direct data transfer with RoCE" width="1201" height="485" data-path="books/ai-systems-performance-engineering/images/ch04-figure-4-3-gpudirect-rdma.png" />
</Frame>

三条路径可以这样理解：

| 路径                 | 数据怎么走                                                               | 适用直觉                  |
| ------------------ | ------------------------------------------------------------------- | --------------------- |
| **普通 TCP/IP**      | GPU -> CPU 内存 -> kernel TCP/IP -> NIC -> 网络 -> NIC -> CPU 内存 -> GPU | 兼容性最好，但多拷贝、CPU 参与多    |
| **RDMA**           | 应用内存 -> NIC -> 网络 -> NIC -> 远端应用内存                                  | 绕过内核数据路径，降低延迟和 CPU 开销 |
| **GPUDirect RDMA** | GPU HBM -> NIC -> InfiniBand/RoCE -> NIC -> GPU HBM                 | 多节点 GPU 训练和推理的关键快路径   |

<Frame caption="TCP/IP、RDMA 与 GPUDirect RDMA 的数据路径对比：越往下，CPU 和内核网络栈参与越少">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-tcp-rdma-gdr-comparison.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=00f9505257624853e60d5a13ddd89f3f" alt="Comparison of TCP/IP, RDMA, and GPUDirect RDMA data paths across two nodes" width="1672" height="941" data-path="books/ai-systems-performance-engineering/images/ch04-tcp-rdma-gdr-comparison.png" />
</Frame>

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

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

<Frame caption="跨节点一跳：GPUDirect RDMA 让 NIC 通过 DMA 直接读写 GPU HBM，数据路径绕过 host memory 和内核网络栈">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-gpudirect-rdma-hop.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=7a8cc36360bdfa249fcb7da150b134f9" alt="Cross-node GPUDirect RDMA path from GPU3 HBM through NICs and IB or RoCE network to GPU4 HBM, bypassing host memory" width="1693" height="929" data-path="books/ai-systems-performance-engineering/images/ch04-generated-gpudirect-rdma-hop.png" />
</Frame>

| 步骤                       | 发生了什么                                  | 耗时                       |
| ------------------------ | -------------------------------------- | ------------------------ |
| 1. NCCL kernel 触发传输      | GPU3 SM 通知 NIC "数据在 HBM 这个地址，发出去"      | \~μs                     |
| 2. NIC DMA 读 GPU3 HBM    | ConnectX-7 NIC 通过 PCIe Gen5 直接读 GPU 显存 | 64 MB / 64 GB/s ≈ 1.0 ms |
| 3. InfiniBand 传输         | NDR 400 Gb/s = 50 GB/s                 | 64 MB / 50 GB/s ≈ 1.3 ms |
| 4. 远端 NIC DMA 写 GPU4 HBM | 对端 NIC 通过 PCIe 直接写入 GPU4 显存            | \~1.0 ms                 |
| 5. GPU4 SM 累加            | reduce kernel 执行                       | \~μs                     |

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×**，且业务代码不报错。

| 路径             | 步骤数                                | 单跳延迟（64 MB） | CPU 参与 |
| -------------- | ---------------------------------- | ----------- | ------ |
| GPUDirect RDMA | GPU→NIC→网络→NIC→GPU                 | \~3.3 ms    | 仅初始化   |
| CPU staging    | GPU→CPU→NIC→网络→NIC→CPU→GPU         | \~6–8 ms    | 每次传输   |
| 纯 TCP/IP       | GPU→CPU→内核栈→NIC→网络→NIC→内核栈→CPU→GPU | \~10+ ms    | 全程     |

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

这里最容易混淆的是把 InfiniBand、RoCE 和普通以太网当成同一层的三个名字。更准确的看法是：先看物理网络和网卡能力，再看上面跑的是 RDMA 还是 TCP/IP，最后看 NCCL 选择了哪个通信后端。

| 层级            | 需要问的问题                             | 常见情况                                                             |
| ------------- | ---------------------------------- | ---------------------------------------------------------------- |
| **物理网络 / 网卡** | 机器接的是 IB 交换机，还是以太网交换机？网卡是否支持 RDMA？ | InfiniBand HCA；支持 RoCE 的以太网 NIC；普通以太网 NIC                        |
| **传输语义**      | 数据能否绕过内核 TCP/IP 路径，直接 RDMA 读写远端内存？ | InfiniBand 原生 RDMA；RoCE = RDMA over Converged Ethernet；普通 TCP/IP |
| **GPU 数据路径**  | NIC 能否直接访问 GPU HBM？                | GPUDirect RDMA 快路径；否则退回 CPU staging                              |
| **NCCL 后端**   | NCCL 最后走的是 RDMA 后端还是 socket 后端？    | `NET/IB` 表示 RDMA 路径；`NET/Socket` 表示 TCP/socket 路径                |

因此，InfiniBand 和 RoCE 都可以提供 RDMA 语义，只是承载网络不同：InfiniBand 是专用 IB fabric；RoCE 是在以太网上承载 RDMA。普通以太网本身通常走 TCP/IP，不提供 RDMA 快路径。对 NCCL 来说，真正重要的是它有没有成功走到 RDMA 后端，以及 GPUDirect RDMA 是否可用。

注意 RoCE 网卡和普通以太网卡在 Linux 层面可能有相同的设备名（`eth0` 或 `bond0`），不能只靠接口名判断。RoCE 是否生效取决于网卡能力、网络配置、驱动/插件和权限。NCCL 日志里如果看到 `NET/IB`，通常说明它走到了 RDMA 路径；如果只看到 `NET/Socket`，即使底层是高速以太网，也是在走 TCP/socket 慢路径。

Bond 本身不是一种网络技术，而是 Linux 把多张物理网卡捆绑成一个逻辑接口的机制（提高带宽或做冗余），InfiniBand 和以太网卡都可以做 bond。

容器和 Kubernetes 环境中最容易出现"看起来能跑，但其实走慢路径"的问题。常见原因包括 `/dev/infiniband` 没有暴露给容器、`nvidia_peermem` 没加载、GID/权限不匹配、NCCL 选错网卡，或者 RoCE 网络没有正确配置。结果是 NCCL 悄悄退回 socket/TCP，吞吐明显下降，但业务代码不一定报错。

<Frame caption="直接连接 NIC 可以绕过 CPU 到 PCIe switch 的瓶颈，让 GPU 与 NIC 之间获得更完整的链路吞吐">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-figure-4-4-direct-nic.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=9a95827a1b978ce4a93ecc488f715fca" alt="Bypassing CPU bottlenecks with direct connectivity between GPUs and NICs" width="1314" height="639" data-path="books/ai-systems-performance-engineering/images/ch04-figure-4-4-direct-nic.png" />
</Frame>

### 2.4 网络内计算：SHARP

传统 all-reduce 中，每张 GPU 要收发 `2(N-1)/N × data_size` 的数据，reduction 运算（求和/平均）在 GPU 端点完成。SHARP 的思路是：**把 reduction 下沉到 InfiniBand Quantum 交换机里做，端点只收最终结果。**

|               | 传统 All-Reduce | SHARP    |
| ------------- | ------------- | -------- |
| reduction 在哪做 | GPU SM        | 交换机      |
| 每个 GPU 接收     | B×(N-1)/N     | B/N      |
| 流量缩减          | —             | \~(N-1)× |

GPU 把本地梯度发到交换机，交换机直接求和后把结果返回。需要交换机 firmware 支持 + Aggregation Manager 进程。**跨节点**生效。节点内的对应方案是 NVLS（见上文节点内互联）。

效果：大规模训练（数百\~千卡）下 all-reduce 加速 2×–5×，且不需要改用户代码——NCCL 检测到硬件支持后自动启用。

### 2.5 管理与诊断

#### 2.5.1 排查入口

常用排查入口：

| 目标               | 工具 / 环境变量                         |
| ---------------- | --------------------------------- |
| 查看初始化、拓扑和网络路径    | `NCCL_DEBUG=INFO`                 |
| 聚焦网络和 collective | `NCCL_DEBUG_SUBSYS=INIT,NET,COLL` |
| 指定网卡，避免走错接口      | `NCCL_SOCKET_IFNAME=<iface>`      |
| 验证 RDMA 路径       | NCCL 日志、RDMA perftest、网卡计数器       |
| 看通信是否覆盖计算        | Nsight Systems timeline           |

调优 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 拓扑。不同现象对应的优化手段不同。

| 现象                           | 可能原因                         | 优先考虑的优化                                | 直觉                           |
| ---------------------------- | ---------------------------- | -------------------------------------- | ---------------------------- |
| 单机多卡时 GPU0 更忙，其他 GPU 等待      | 使用 `DataParallel`，主 GPU 负责聚合 | 换成 `DistributedDataParallel`           | 避免 GPU0 聚合瓶颈，让 NCCL 负责同步     |
| backward 结束后出现一大段 all-reduce | 梯度同步启动太晚                     | DDP bucket、CUDA streams、通信计算重叠         | 边算边同步，减少 GPU 空转              |
| 每个 step 都有固定同步开销             | 同步频率太高                       | 梯度累积                                   | 多个 micro-batch 后再同步一次        |
| 网络链路持续打满，GPU 等通信             | 每次传输数据量太大                    | 梯度压缩、量化、稀疏化                            | 少传数据，但要关注收敛影响                |
| all-reduce 时 CPU 很忙，GPU 通信慢  | 后端误用 Gloo 或 socket 路径        | 使用 NCCL，并确认走 GPU 通信路径                  | GPU collective 不应主要靠 CPU 搬数据 |
| 跨节点吞吐远低于预期                   | NCCL 走 TCP 或 CPU staging 慢路径 | RDMA / GPUDirect RDMA                  | 让 NIC 直接访问注册内存或 GPU HBM      |
| GPU 数量增加后 all-reduce 扩展性变差   | collective 没有匹配硬件拓扑          | NCCL 拓扑感知、SHARP、NVLS                   | 让节点内、节点间和网络内计算各走合适路径         |
| 某个 rank 经常拖慢整轮迭代             | straggler、NUMA 绑定不当或网络局部拥塞   | 检查 rank 级 timeline、CPU affinity 和网卡计数器 | collective 的总耗时由最慢参与者决定      |
| 分离式推理中 decode 等 KV cache     | KV cache 点对点迁移成为瓶颈           | NIXL                                   | 大块点对点数据异步移动                  |

**框架通信模式问题。** 单机多卡训练中，如果仍在使用 `DataParallel`，GPU0 往往会更忙。`DataParallel` 是单进程控制多张 GPU，主 GPU 负责 scatter 输入、gather 输出和聚合梯度，容易形成主卡瓶颈。`DistributedDataParallel` 通常是一进程一 GPU，用 NCCL 做梯度 all-reduce，并在 backward 过程中按 bucket 发起异步通信，因此更适合多 GPU 训练。

**通信启动太晚。** 朴素做法是所有 backward 结束后再统一 all-reduce，这会把通信时间直接加到每轮迭代上。更好的做法是把梯度切成 bucket，某个 bucket 一准备好就发起 NCCL all-reduce，同时继续计算后续层的梯度。

<Frame caption="通信与计算重叠：同步模式会让 GPU 等待 all-reduce；DDP bucket 可以在 backward 过程中提前启动通信">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-figure-4-1-overlap.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=ba08a51a83c9d03c4153c085fe7b30a0" alt="Overlapping host-to-device and device-to-host communication with computation on multiple CUDA streams" width="1387" height="853" data-path="books/ai-systems-performance-engineering/images/ch04-figure-4-1-overlap.png" />
</Frame>

需要避免的常见同步点包括 `torch.cuda.synchronize()` 和 `tensor.item()`。前者会等待当前设备上的工作全部完成，后者会把 GPU 标量搬回 CPU，也可能触发同步。计时或日志统计应尽量放在迭代末尾，避免打断重叠流水线。

**同步频率太高。** 如果每个 step 都有固定通信开销，可以用梯度累积减少同步次数。多个 micro-batch 先在本地累积梯度，再做一次 all-reduce，相当于用更多本地计算换更少的跨 GPU 同步。代价是有效 batch size 变大，可能需要重新调学习率和内存预算。

**Before — 每步同步：**

```python theme={null}
for data, target in dataloader:
    output = model(data)
    loss = criterion(output, target)
    loss.backward()         # 每个 micro-batch 都触发 all-reduce
    optimizer.step()
    optimizer.zero_grad()
# 如果一步 all-reduce 耗时 5ms，100 步 = 500ms 纯通信
```

**After — 梯度累积，4 步同步一次：**

```python theme={null}
accumulation_steps = 4
for i, (data, target) in enumerate(dataloader):
    # no_sync() 关闭 DDP 的自动 all-reduce，梯度只在本地累积
    context = model.no_sync() if (i + 1) % accumulation_steps != 0 else nullcontext()
    with context:
        output = model(data)
        loss = criterion(output, target) / accumulation_steps
        loss.backward()

    if (i + 1) % accumulation_steps == 0:
        optimizer.step()    # 只在第 4 步才 all-reduce + 更新
        optimizer.zero_grad()
# 同样 100 步，all-reduce 只发生 25 次 = 125ms，通信量减少 75%
```

**传输数据量太大。** 如果网络链路持续打满，压缩、量化或稀疏化可以减少每次通信的数据量。这类方法直接作用在"传多少"上，但会改变梯度信息，需要关注收敛稳定性和最终精度。

**跨节点路径太慢。** 如果跨节点吞吐远低于预期，先确认通信后端是否正确。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                          |
| ---- | ------------------ | ------------------------------- |
| 特征   | compute-bound，大矩阵乘 | memory-bound，逐 token 读 KV cache |
| 关键指标 | TTFT（首 token 延迟）   | TPOT（每 token 延迟）                |

放在一起时，prefill 抢算力导致 decode 变慢，为保 decode 又让 prefill 排队。DistServe（OSDI 2024）和 Splitwise（Microsoft）的解法：**拆到不同 GPU 上，各自独立优化资源和并行策略。** DistServe 在同样 GPU 预算下服务 7.4x 更多请求。

<Frame caption="PD 分离推理的数据路径：prefill worker 产出 KV cache，通过点对点异步传输交给 decode worker，必要时在 GPU HBM、CPU DRAM 和 NVMe 之间分层移动">
  <img src="https://mintcdn.com/se7en/vQcSEdezjjIiLrs-/books/ai-systems-performance-engineering/images/ch04-generated-pd-kv-cache-transfer.png?fit=max&auto=format&n=vQcSEdezjjIiLrs-&q=85&s=436569d5dad829d7436c93c7661655a8" alt="Disaggregated inference with prefill workers sending KV cache blocks point-to-point to decode workers across a high-speed network" width="1672" height="941" data-path="books/ai-systems-performance-engineering/images/ch04-generated-pd-kv-cache-transfer.png" />
</Frame>

#### 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。

|      | NCCL            | NIXL                |
| ---- | --------------- | ------------------- |
| 场景   | 训练              | 推理                  |
| 操作   | collective（多对多） | point-to-point（一对一） |
| 参与者  | 固定 N 个 rank     | 动态 worker 池         |
| 数据   | 梯度              | KV cache / tensor   |
| 抽象层级 | 同层，都在框架和硬件之间    | 同层                  |

三者关系：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 并记录平均耗时：

```python theme={null}
# ddp_bench.py — DDP bucket_cap_mb benchmark
import argparse, os, time, torch, torch.distributed as dist, torch.nn as nn
from torch.nn.parallel import DistributedDataParallel as DDP

def parse_args():
    p = argparse.ArgumentParser()
    p.add_argument("--bucket_cap_mb", type=int, default=25)
    p.add_argument("--num_layers", type=int, default=6)
    p.add_argument("--hidden_size", type=int, default=2048)
    p.add_argument("--batch_size", type=int, default=32)
    p.add_argument("--seq_len", type=int, default=128)
    p.add_argument("--warmup_steps", type=int, default=5)
    p.add_argument("--bench_steps", type=int, default=20)
    p.add_argument("--profile", action="store_true")
    return p.parse_args()

class SimpleTransformer(nn.Module):
    def __init__(self, num_layers, hidden_size):
        super().__init__()
        self.embed = nn.Linear(hidden_size, hidden_size)
        encoder_layer = nn.TransformerEncoderLayer(
            d_model=hidden_size, nhead=8, dim_feedforward=hidden_size * 4, batch_first=True,
        )
        self.encoder = nn.TransformerEncoder(encoder_layer, num_layers=num_layers)
        self.head = nn.Linear(hidden_size, hidden_size)

    def forward(self, x):
        return self.head(self.encoder(self.embed(x)))

def main():
    args = parse_args()
    dist.init_process_group("nccl")
    local_rank = int(os.environ["LOCAL_RANK"])
    torch.cuda.set_device(local_rank)

    model = SimpleTransformer(args.num_layers, args.hidden_size).cuda()
    model = DDP(model, device_ids=[local_rank], bucket_cap_mb=args.bucket_cap_mb)
    optimizer = torch.optim.AdamW(model.parameters(), lr=1e-4)

    if local_rank == 0:
        param_mb = sum(p.numel() * p.element_size() for p in model.parameters()) / 1e6
        print(f"Model: {param_mb:.0f} MB | bucket_cap_mb={args.bucket_cap_mb}")

    data = torch.randn(args.batch_size, args.seq_len, args.hidden_size, device="cuda")
    target = torch.randn_like(data)

    def run_step():
        optimizer.zero_grad()
        loss = nn.functional.mse_loss(model(data), target)
        loss.backward()
        optimizer.step()

    # warmup
    for _ in range(args.warmup_steps):
        run_step()
    torch.cuda.synchronize()

    if args.profile:
        with torch.profiler.profile(
            activities=[torch.profiler.ProfilerActivity.CPU, torch.profiler.ProfilerActivity.CUDA],
            record_shapes=True,
        ) as prof:
            for _ in range(args.bench_steps):
                run_step()
            torch.cuda.synchronize()
        if local_rank == 0:
            print(prof.key_averages().table(sort_by="cuda_time_total", row_limit=30))
    else:
        torch.cuda.synchronize()
        t0 = time.perf_counter()
        for _ in range(args.bench_steps):
            run_step()
        torch.cuda.synchronize()
        elapsed = (time.perf_counter() - t0) / args.bench_steps * 1000
        if local_rank == 0:
            print(f"bucket_cap_mb={args.bucket_cap_mb}  ms/step={elapsed:.2f}")

    dist.destroy_process_group()

if __name__ == "__main__":
    main()
```

用下面的脚本扫描不同的 bucket\_cap\_mb 配置：

```bash theme={null}
# 扫描不同 bucket_cap_mb，记录每个配置的 ms/step
for bucket in 1 5 10 25 50 100 200; do
  RCCL_DEBUG=INFO torchrun --nproc_per_node=8 ddp_bench.py \
    --bucket_cap_mb $bucket \
    --num_layers 6 --hidden_size 2048 \
    2>&1 | tee log_bucket_${bucket}.txt
done
```

#### 3.1.1 小模型（88 MB，6 层 2048-hidden）

| bucket\_cap\_mb | ms/step  | 相对最优   |
| --------------- | -------- | ------ |
| 1               | 12.87    | +108%  |
| 5               | 10.55    | +71%   |
| 10              | 7.90     | +28%   |
| 25（默认）          | 9.82     | +59%   |
| **50**          | **6.18** | **基准** |
| 100             | 7.78     | +26%   |
| 200             | 7.32     | +18%   |

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

#### 3.1.2 大模型（1628 MB，24 层 4096-hidden）

| bucket\_cap\_mb | ms/step   | 相对最优   |
| --------------- | --------- | ------ |
| 1               | 28.57     | +6%    |
| **5**           | **26.95** | **基准** |
| 10              | 27.36     | +2%    |
| 25（默认）          | 28.35     | +5%    |
| 50              | 36.83     | +37%   |
| 100             | 35.73     | +33%   |
| 200             | 35.90     | +33%   |

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

#### 3.1.3 Profiler 通信/计算分解（大模型 1628 MB）

```bash theme={null}
# 用 torch.profiler 采集 kernel 级耗时，按名称分类为通信和计算
torchrun --nproc_per_node=8 ddp_bench.py \
  --bucket_cap_mb 5 --num_layers 24 --hidden_size 4096 \
  --profile
```

| bucket\_cap\_mb | 通信 (ms) | 计算 (ms) | 通信占比  |
| --------------- | ------- | ------- | ----- |
| 5               | 128.97  | 119.36  | 43.8% |
| 25              | 113.72  | 107.20  | 39.0% |
| 50              | 142.41  | 126.25  | 40.9% |
| 100             | 119.59  | 114.86  | 38.6% |

通信和计算接近 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

<Accordion title="推理时，请求状态为什么也需要移动？">
  PD 分离架构下，prefill 和 decode 运行在不同 GPU 上。请求交给 decode worker 继续生成时，不仅要搬 KV
  cache，还要搬请求的元信息——当前生成到第几个 token、sampling 参数、stop conditions、已生成的 token ids。没有这些状态，decode
  worker 不知道从哪接着生成。
</Accordion>

<Accordion title="DDP 的 bucket 是什么？梯度同步必须所有 GPU 一起完成吗？">
  Bucket 是 DDP 对参数梯度的预分组（默认每桶上限 25 MB）。单个 bucket 的 all-reduce 确实是 collective——所有 rank 必须共同完成。但
  DDP 不是对整个模型做一次 all-reduce，而是按 bucket 拆成多次独立的 all-reduce。

  Backward 从最后一层往前算，最后几层的梯度凑满一个 bucket 后，这个 bucket 的 all-reduce 立即启动，同时 backward
  继续计算前面的层。通信和计算流水线式重叠。这就是为什么 bucket 大小很重要——太大则通信启动太晚，太小则启动次数多、开销大。

  Bucket 和层不是一一对应的。一个层通常有多个参数（如 FFN 的 `W1`、`b1`、`W2`、`b2`），这些参数可能分属不同
  bucket；反过来，多个小层的参数也可能合并进同一个 bucket。DDP 按参数的 reverse 顺序依次填充 bucket，填满一个再开下一个，不关心层边界。但**单个参数的梯度不会被拆到多个 bucket**——如果某个参数超过 `bucket_cap_mb`，它独占一个
  bucket。Autograd 对每个参数一次性算完整个 grad tensor，算完才标记 ready；bucket 内所有参数都 ready 后才触发
  all-reduce。所以不存在"梯度算了一半就被 all-reduce 走"的情况。
</Accordion>

<Accordion title="互补树好构造吗？能把带宽用满吗？">
  构造方法来自 \[Sanders, Speck & Träff (2009)][https://www.sciencedirect.com/science/article/abs/pii/S0167819109000957），对](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）](https://developer.nvidia.com/blog/massively-scale-deep-learning-training-nccl-2-4/#ref3）)
</Accordion>

<Accordion title="Ring 和 Tree 有没有卡数判断基线？单机八卡多机还有必要走 Tree 吗？">
  没有固定卡数阈值，取决于消息大小 × GPU 数量。大消息 + 少量 GPU 倾向
  Ring（带宽打满，步数少时延迟不是问题）；小消息或大规模集群倾向 Tree（O(log N) 步，启动快）。

  单机八卡或小规模多机（2-4 台），Ring 通常够用。64+ 卡时 Tree 的 O(log N) 优势开始显现。实际中 NCCL 会对不同 collective
  自动选择——大 bucket 走 Ring 吃带宽，小 bucket 走 Tree 低延迟，通常不需要手动指定。
</Accordion>

<Accordion title="有了 GPUDirect 之后，还有哪些操作在 CPU？Autograd hook 能移到 GPU 吗？">
  核心数据路径（梯度/参数/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
  级，不是瓶颈。
</Accordion>

<Accordion title="Registered host memory 和普通 CPU memory 一样吗？">
  不一样。Registered（pinned）memory 通过 `cudaHostRegister()` 注册后页锁定——OS 不会 swap 到磁盘，GPU DMA 和 NIC
  可以直接访问，物理地址已知，硬件可以直接 DMA。普通 CPU 内存（pageable）随时可能被换出，需要先拷贝到 pinned buffer
  再传输，多一次拷贝。
</Accordion>

<Accordion title="NCCL 的算法/协议自动选择需要手动开启吗？">
  不需要，默认行为。NCCL 根据消息大小、GPU 数量和拓扑自动选 Ring/Tree 和 Simple/LL/LL128。手动覆盖用 `NCCL_ALGO` 和
  `NCCL_PROTO`，但一般不建议，除非 profiler 显示自动选择不是最优。DDP 层面的 `bucket_cap_mb` 需要手动调。
</Accordion>

<Accordion title="可以把 TCP 路径禁掉吗？">
  不建议。NCCL bootstrap（进程发现和初始化握手）必须走 TCP。应该做的是确保数据面不走 TCP：正确配置 RDMA 环境，用 `NCCL_DEBUG=INFO`
  确认日志中数据传输走 `NET/IB` 而非 `NET/Socket`。如果看到后者，排查 RDMA 配置而不是禁 TCP。
</Accordion>

<Accordion title="Straggler 怎么找？所有 rank 都卡住了，很难定位">
  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 也能发现链路级问题。
</Accordion>

<Accordion title="NCCL_MAX_NCHANNELS 怎么调？">
  Channel 是 NCCL 的并行通信流水线。更多 channel = 更高带宽利用，但占更多 SM 资源。默认值通常够用。调之前先看 Nsight Systems
  确认瓶颈在通信带宽——带宽没打满且 SM 有余量可以调大；通信 kernel 抢 SM 影响计算，应调小或改用 Copy Engine Collectives。

  注意：NCCL 2.21+ 将 `NCCL_MAX_NCHANNELS` 改名为 `NCCL_MAX_CTAS`（对应 `NCCL_MIN_CTAS`）。旧名称在部分版本仍兼容，但新部署建议用新名称。
</Accordion>

<Accordion title="Profiling 是不是要一直开着？">
  不需要，profiling 本身有 5-15% 开销。发现异常后下一次运行中加 profiler，采集几个 step 即可（`nsys profile --duration 30` 或
  PyTorch Profiler 指定步数）。如果问题是卡死（hang）而非"慢"，需要不同工具：`NCCL_DEBUG=WARN` 看异步错误、`py-spy`
  看调用栈、`nvidia-smi` 看 GPU 利用率是否归零。
</Accordion>

<Accordion title="NIXL 和 NCCL 有什么区别？NCCL group 只有 2 个节点就是 NIXL 吗？">
  不是。即使 NCCL group 只有 2 个节点，做的仍然是 collective 语义——双方必须同步到达同步点。NIXL 是异步点对点语义：发送方 put
  数据到接收方指定地址，不需要接收方同时参与；支持多种内存类型（GPU HBM / CPU DRAM / NVMe）任意组合；worker
  动态加入离开；自动选最快路径。NCCL 为训练设计（固定 rank、collective），NIXL 为推理设计（动态 worker、point-to-point）。
</Accordion>

<Accordion title="Inference 数据是动态的，NCCL 能很好处理吗？">
  NCCL 假设固定数量的 rank、固定的 communicator、重复执行相同的 collective，不匹配推理场景（动态 worker 池、点对点 KV cache
  传输、连接随请求建立/销毁）。这正是 NIXL 存在的原因——专为推理场景的异步点对点传输设计。
</Accordion>

<Accordion title="除了数据并行按 GPU 数量均分 mini-batch，其他并行策略如何切分数据？能否不按单卡均分？">
  先澄清：**训练数据的切分只发生在数据并行维度**。其他并行策略（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 对梯度贡献被低估。
</Accordion>

## 参考资料

* [NCCL Documentation](https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/) — NCCL 官方文档，涵盖安装、环境变量、拓扑配置和调优指南
* [NVIDIA Magnum IO Developer Guide](https://docs.nvidia.com/networking/display/magnum-io-dev-guide/) — Magnum IO 通信栈总览，包括 GPUDirect RDMA、NCCL 和 SHARP 的协同
* [GPUDirect RDMA Documentation](https://docs.nvidia.com/cuda/gpudirect-rdma/) — GPUDirect RDMA 的内核模块要求、内存注册和验证方法
* [NCCL Tests (GitHub)](https://github.com/NVIDIA/nccl-tests) — NCCL 官方 benchmark 工具，用于测量 all-reduce、reduce-scatter 等 collective 的带宽和延迟
* [PyTorch DDP Design Note](https://pytorch.org/docs/stable/notes/ddp.html) — DDP 的 bucket 机制、autograd hook 和通信计算重叠的设计文档
* [PyTorch DDP Tutorial](https://pytorch.org/tutorials/intermediate/ddp_tutorial.html) — DDP 入门教程，从单机多卡到多节点训练
* [DistServe: Disaggregating Prefill and Decoding for Goodput-optimized Large Language Model Serving (OSDI 2024)](https://www.usenix.org/conference/osdi24/presentation/zhong-yinmin) — PD 分离架构的原始论文
* [NIXL: NVIDIA Inference Xfer Library (GitHub)](https://github.com/ai-dynamo/nixl) — 推理场景的点对点数据搬运库
* [Splitwise: Efficient generative LLM inference with model splitting](https://arxiv.org/abs/2311.18677) — Microsoft 的 PD 分离方案
* [Massively Scale Your Deep Learning Training with NCCL 2.4](https://developer.nvidia.com/blog/massively-scale-deep-learning-training-nccl-2-4/) — NCCL 2.4 引入 double binary tree all-reduce 的技术博客，含 Summit 超算 24,576 GPU 实测数据
* [Two-tree algorithms for full bandwidth broadcast, reduction and scan](https://www.sciencedirect.com/science/article/abs/pii/S0167819109000957) — Sanders, Speck & Träff (2009)，double binary tree 的理论基础
* [Scaling Deep Learning Training with NVSHMEM](https://developer.nvidia.com/blog/scaling-deep-learning-training-nvshmem/) — NVSHMEM GPU-initiated 通信模型，MPI 与 NVSHMEM 通信模式对比
* [到底什么是All-Reduce、All-to-All？](https://zhuanlan.zhihu.com/p/1976642212977713852) - 图解各种通信原语
* [Ring Allreduce并行计算优化](https://www.deeplearn.me/2649.html) - 图解Ring Allreduce
* [Massively Scale Your Deep Learning Training with NCCL 2.4](https://developer.nvidia.com/blog/massively-scale-deep-learning-training-nccl-2-4/#ref3) - 图解Double Binaty Tree
