/ AI Infra
[Infra-10] All-to-All Optimization: From MoE Token Routing to Communication-Computation Overlap
A systems guide to All-to-All optimization in MoE, covering dispatch and combine, uneven routing, small messages, topology, permutation overhead, hierarchical communication, kernel fusion, overlap, expert placement, and profiling.
在 MoE、Expert Parallel、Sequence Parallel 和分布式推理系统里,All-to-All(也常写作 All2All 或 A2A)不是一个孤立的通信 API,而是一条完整的数据重分布路径。
它从 Router 的输出开始,经过 token 计数、重排和打包,将不同 token 送到不同 GPU 上的 Expert;Expert 计算完成后,再沿反方向把结果送回原位置。网络传输只是其中一段,前后的 permutation、内存拷贝、同步和负载不均衡都可能比 collective 本身更贵。
因此,All-to-All 优化真正要回答的不是“怎样把 NCCL 调用变快”,而是:
怎样用更少的数据、更少的拷贝和同步,
沿着更合适的拓扑完成 token 重分布,
同时把无法消除的通信隐藏在 Expert 计算后面。本文从通信语义、MoE 数据流、性能模型、优化方法和 profiling 五个层次,建立一套分析 All-to-All 的系统框架。
1. All-to-All 到底做了什么#
假设有 4 个 rank,每个 rank 都准备了 4 份数据,第 份要发给 rank :
发送前:
Rank 0: [A00, A01, A02, A03]
Rank 1: [A10, A11, A12, A13]
Rank 2: [A20, A21, A22, A23]
Rank 3: [A30, A31, A32, A33]
接收后:
Rank 0: [A00, A10, A20, A30]
Rank 1: [A01, A11, A21, A31]
Rank 2: [A02, A12, A22, A32]
Rank 3: [A03, A13, A23, A33]从全局看,这很像一次分布式矩阵转置:发送方按目标 rank 切分,接收方按来源 rank 拼接。
它与几个常见 collective 的区别是:
| Collective | 每个 rank 发送什么 | 每个 rank 得到什么 | 常见用途 |
|---|---|---|---|
| AllGather | 自己的一整份 shard | 所有 rank 的 shard | TP activation、参数收集 |
| AllReduce | 自己的一整份 tensor | 所有 tensor 的规约结果 | 梯度、TP partial result |
| ReduceScatter | 一整份 tensor | 规约结果的一份 shard | ZeRO、并行规约 |
| All-to-All | 给不同目标的不同分片 | 来自不同来源的目标分片 | MoE dispatch、序列转置 |
这里还有一个重要区别:连接数增加,不代表每张卡的有效 payload 必然按 倍增加。若每个 rank 的输入总量固定,只是均匀切给 个目标,那么远端发送量约为总量的 ;真正随规模恶化的往往是 peer 数、消息粒度、同步轮次、链路竞争与跨节点比例。
1.1 固定大小与可变大小#
规则的 All-to-All 要求发往每个目标的数据长度固定。但 MoE Router 很少生成完全均匀的 token 分布:
Rank 0 -> Rank 1: 1000 tokens
Rank 0 -> Rank 2: 10 tokens
Rank 0 -> Rank 3: 300 tokens这更接近 all_to_all_v:每个 source-destination pair 的长度都可能不同。系统不仅要搬数据,还要交换 split size 或 offset,处理空消息、容量限制和动态 buffer。
这也是通用 collective benchmark 与真实 MoE 性能经常不一致的原因。前者通常测等长、连续的大 buffer;后者面对的是由 Router 动态产生的、带 token 元数据的不规则流量。
2. 为什么 MoE 每层通常需要两次 All-to-All#
假设 8 个 Expert 分布在 4 张 GPU 上:
GPU 0: Expert 0, 1
GPU 1: Expert 2, 3
GPU 2: Expert 4, 5
GPU 3: Expert 6, 7每张 GPU 起初持有自己的 token shard。Router 为每个 token 选择 Top-K Expert 后,token 所在 GPU 与 Expert 所在 GPU 往往不同。例如 GPU 0 上的 token 被路由为:
token 0 -> Expert 1 -> GPU 0
token 1 -> Expert 3 -> GPU 1
token 2 -> Expert 4 -> GPU 2
token 3 -> Expert 7 -> GPU 3于是一个 MoE 层会出现两次方向相反的数据交换:
hidden states on original ranks
|
v
Router / Top-K
|
v
Permute + Pack by destination
|
v
All-to-All Dispatch: token rank -> expert rank
|
v
Expert grouped GEMM
|
v
All-to-All Combine: expert rank -> token rank
|
v
Unpack + weighted combine
|
v
hidden states restored to original order第一次 dispatch 把 token 送到 Expert,第二次 combine 把 Expert 输出送回 token 原来的并行域。Top-2 路由通常还会让一个 token 产生两份 Expert 输入和两份返回结果,因此通信量、Expert 计算量和 combine 工作量都会高于 Top-1。
2.1 通信的不只有 hidden states#
真实实现还可能携带:
- token 对应的 Expert ID;
- 原始 token index 或反向映射;
- routing weight;
- 每个 Expert 或 rank 的 token count;
- quantization scale 等附加信息。
如果只用 tokens * hidden_size * dtype_bytes 估算流量,会漏掉元数据、padding、对齐以及 Top-K 复制。一个更实用的近似是:
其中 是本 rank 的 token 数, 是 Top-K, 是 hidden size, 是每个元素的字节数。若只统计跨 rank 流量,还要乘以非本地路由比例,并额外计入 metadata 和 padding。
3. All-to-All 为什么容易变成瓶颈#
3.1 Peer 多,通信调度复杂#
在 个 rank 中,每个 rank 最多需要与其余 个 rank 通信。若按有向 source-destination pair 计数,全局最多存在:
条远端通信关系。实际实现不会简单地让所有 pair 无序争抢,但 peer 数越多,连接管理、通信轮次、NIC queue、credit 和网络竞争越难控制。
AllReduce 的数据流通常更规则,可以稳定使用 ring、tree 等算法;MoE All-to-All 的 split size 每步变化,热点目标也会变化,通信库更难提前形成同样稳定的调度。
3.2 Router 不均匀会同时制造计算和通信长尾#
MoE 性能不是由平均 token 数决定,而常由最忙的 rank 或 Expert 决定:
Rank 0: 10 ms
Rank 1: 12 ms
Rank 2: 35 ms <- critical path
Rank 3: 11 ms
layer latency: at least 35 ms热门 Expert 所在 rank 不仅要接收更多 bytes,还要执行更多 GEMM;其他 rank 即使早已完成,也必须在后续同步点等待。因而 profile 中看到的“通信等待”不一定意味着网络本身慢,也可能是上游路由倾斜或 Expert 计算不均衡。
3.3 Decode 经常受小消息延迟支配#
小 batch decode 的每步 token 数很少,切到多个目标后,每个 peer 可能只收到几个 token。此时近似模型:
其中 表示一次通信阶段的固定启动与调度成本, 是有效轮次, 是 payload。小消息时第二项很小,kernel launch、CPU 调度、协议握手、队列管理和同步反而成为主体。
所以 decode MoE 往往不是“带宽没跑满”的问题。消息本来就不足以进入带宽饱和区间,真正应优化的是固定开销和关键路径延迟。
3.4 节点内和节点间不是同一种通信#
节点内通常经过 NVLink、NVSwitch 或 PCIe;跨节点还要经过 GPU/CPU memory path、PCIe、NIC、交换网络和远端 NIC。不同路径的带宽、延迟与并发能力可能相差很大:
intra-node: GPU -> NVLink / NVSwitch / PCIe -> GPU
inter-node: GPU -> PCIe -> NIC -> fabric
-> NIC -> PCIe -> GPU若对拓扑一视同仁,同样数量的 token 可能产生完全不同的代价。更糟的是,大量跨节点小包可能同时争抢有限 NIC,而节点内高速链路没有充分利用。
3.5 Pack 和 Unpack 也在关键路径上#
Router 输出通常是离散的 Expert ID。通信前至少要完成:
count / histogram
-> prefix sum / offset
-> sort or scatter by destination
-> pack into contiguous send buffers通信后还要:
receive
-> unpack by local expert
-> grouped GEMM
-> reverse dispatch
-> restore original token order
-> weighted combine因此端到端 MoE dispatch 时间更接近:
只盯着通信 kernel 会漏掉完整关键路径。一次额外的 HBM round trip,或 count 结果从 GPU 回到 CPU 后再发起通信,都可能让“网络优化”失去意义。
4. 优化第一层:减少必须跨域传输的数据#
所有通信优化里,最可靠的收益仍然来自不发送。
4.1 控制 Top-K 与通信精度#
Top-2 会让每个 token 进入两个 Expert,通常比 Top-1 产生更多 dispatch、combine 和 Expert 计算。降低 Top-K 可以直接减少数据,但它改变模型结构与质量,不能只作为系统开关处理。
另一条路线是在通信路径采用 FP8 或 INT8:先量化 hidden states,传输后再反量化。它可能把 payload 压缩到 BF16/FP16 的一半甚至更低,但收益必须扣除量化 kernel、scale metadata、精度损失和额外同步。
4.2 提高 Expert locality#
Router 可以对本地或同节点 Expert 增加偏好,或限制 token 只在少数节点内选择 Expert:
global routing: token may select any expert in the cluster
node-limited routing: token selects experts from a small node subset
locality-aware routing: prefer local experts when scores are close这能把跨节点流量转化为节点内流量,减少 NIC 压力。但它也缩小了 Router 的选择空间,可能损害负载均衡或模型质量。正确评估必须同时看 loss/quality、Expert 分布、跨节点 bytes 和尾延迟。
4.3 避免无效 padding 与重复拷贝#
capacity-based 实现有时为每个 Expert 分配固定槽位,再把未使用位置一并发送。这让 buffer 规则,但会浪费带宽。可变长度 dispatch 可以避免 padding,却引入 count exchange、动态 offset 和更复杂的 kernel。
同样,应该尽量让 permutation 直接写入通信 buffer:
less efficient:
input -> temporary permutation buffer -> communication buffer
preferred:
input -> destination-contiguous communication buffer优化目标不是追求零临时 buffer,而是减少关键路径上的全量数据搬运,并保证写入模式足够连续、合并访存友好。
5. 优化第二层:让 Pack、通信与 Combine 成为一条流水线#
5.1 融合 token permutation#
朴素实现可能依次启动 histogram、prefix sum、scatter、metadata pack 等多个 kernel。每一步都读写中间 tensor,还要建立 kernel 间依赖。
可行的融合方向包括:
- Router 输出 Expert ID 时同步生成部分计数信息;
- 将 count、offset 计算与 scatter 合并或减少中间落盘;
- 一次写出 hidden states、token index、routing weight 等 payload;
- combine 时融合反向索引、routing weight 乘法和归并。
融合并不意味着强行做成一个巨型 kernel。目标是减少全量 HBM 往返、launch 和全局同步;若融合导致寄存器压力过高、occupancy 下降或原子冲突严重,拆成少数几个阶段可能更好。
5.2 将消息按目标 rank 连续布局#
DMA 和网络传输更适合连续区域。典型 send buffer 会组织为:
[to rank 0][to rank 1][to rank 2] ... [to rank N-1]每段内部还可按 local Expert 分桶:
[rank 2 / expert 4][rank 2 / expert 5]这样既减少零散传输,也让接收侧更容易直接形成 grouped GEMM 所需布局。理想情况下,网络输出 buffer 就是 Expert kernel 的输入,避免接收后再做一次全量转置。
5.3 合并小消息,但不要无限等待#
合并小包可提高带宽利用率、减少 launch 与 queue 开销;等待大包形成却会增加排队延迟。拆成小 chunk 更容易流水,但固定开销更多。
因此 chunk size 是一个系统参数,而不是越大越好:
| Chunk 较大 | Chunk 较小 |
|---|---|
| 带宽利用率更高 | 更早启动通信和计算 |
| launch 与消息数更少 | 更容易形成流水线 |
| 等待数据就绪更久 | 固定开销占比更高 |
| 需要更大 buffer | event、queue 压力更大 |
合适的值取决于 hidden size、dtype、token 数、Expert GEMM 大小、节点内外链路以及目标是吞吐还是延迟。
6. 优化第三层:选择适合拓扑和消息规模的通信算法#
6.1 Direct P2P 与 pairwise exchange#
Direct P2P 让 rank 直接向目标 rank 发 send/recv,适合稀疏或高度不均匀的流量。实现要避免无序 peer 并发造成死锁、拥塞或大量细碎 launch。
Pairwise exchange 将 peer 交互排成若干阶段:每轮只与指定 rank 交换。它的优点是调度规则、瞬时连接数可控;缺点是轮数随 world size 增长,且某轮中的热点或空消息会降低利用率。
6.2 Bruck 类算法#
Bruck 类算法通过多轮转发减少通信轮次,常被用于小消息 All-to-All。代价是部分数据需要中转,总传输字节数和本地重排可能增加。
所以它不是普遍更快:小消息、启动延迟主导时,减少轮次可能获益;大消息、带宽主导时,额外转发反而可能更贵。
6.3 Hierarchical All-to-All#
多机多卡系统更常用分层思路:
1. intra-node regroup / aggregate
2. inter-node exchange through NICs
3. intra-node scatter to destination GPUs例如两台 8 卡机器之间,不必让每个 GPU 都独立产生许多跨节点小消息。可以先在节点内按远端节点聚合,通过少量大块跨 NIC 传输,再在对端节点内分发。
分层通信的价值是让不同硬件各做擅长的事:NVLink/NVSwitch 负责节点内重排,NIC 负责更连续、更大的节点间数据。代价是增加节点内 hop、buffer 和同步,因此仍要比较“聚合成本”与“减少跨节点碎片”的收益。
6.4 拓扑感知必须进入 rank mapping#
算法本身之外,rank 到物理 GPU、NUMA、NIC 和节点的映射同样重要。常见检查项包括:
- GPU 是否连接到预期 NIC 与 CPU NUMA node;
- 多个通信 rank 是否争抢同一 NIC;
- 节点内 P2P 是否实际走 NVLink/NVSwitch,而非绕行 host;
- GPUDirect RDMA 或对应设备直连路径是否生效;
- 网络是否存在 oversubscription 或热点 spine/link。
没有物理拓扑信息,仅凭逻辑 rank 号很难解释 All-to-All 的带宽和长尾。
7. 优化第四层:通信与 Expert 计算重叠#
完全串行的执行是:
pack all tokens
-> wait for full All-to-All
-> run all Expert GEMMs
-> wait for full combine All-to-All
-> restore all tokens更理想的方式是把 token 分 chunk 或按 Expert 建立 ready 条件:部分输入一到达,就立即执行对应 Expert,同时继续接收其他 token。
time -------------------------------------------------------->
dispatch chunk 0: [ network ]
expert chunk 0: [ GEMM ]
dispatch chunk 1: [ network ]
expert chunk 1: [ GEMM ]
combine chunk 0: [ network ]理想情况下,串行时间:
可以接近:
7.1 重叠成立需要哪些条件#
常见实现会用到通信 stream、计算 stream、CUDA event、异步 P2P、persistent kernel 或 signal buffer。但使用两个 stream 不等于自动重叠,至少还要满足:
- Expert 计算不依赖尚未到达的 token;
- 通信和 GEMM 不完全争抢同一 HBM 带宽或执行资源;
- chunk 足够大,能摊薄 launch 与 event 成本;
- 依赖通过设备侧 event/signal 传递,CPU 不成为串行控制点;
- 接收布局能直接被 Expert kernel 消费;
- combine 不因全局 barrier 再次把流水线压平。
重叠也不等于通信消失。如果 GEMM 很小,通信没有足够计算可隐藏;如果通信 kernel 占用大量 SM 或内存带宽,反而可能拖慢 GEMM。最终应比较端到端 MoE layer latency,而不是只看 timeline 中是否出现彩色 kernel 重叠。
8. 优化第五层:Expert 放置和负载均衡#
8.1 Auxiliary loss、capacity 与 token dropping#
训练中常用 auxiliary load balancing loss 鼓励 Router 均匀使用 Expert。Capacity factor 给每个 Expert 设置可接收 token 上限,超过容量的 token 可能被丢弃、回退或重新路由。
它们改善系统长尾,但都带模型语义:过强的均衡约束可能影响 Router 学习;过小 capacity 可能丢 token;过大 capacity 又会增加 padding、显存和最坏情况 buffer。
8.2 Expert replication#
若少数 Expert 长期热门,可以将其复制到多个 GPU:
GPU 0: E0, E1
GPU 1: E0, E2
GPU 2: E3, E4
GPU 3: E5, E6Router 或调度器再把访问 E0 的 token 分散到不同 replica。这可以降低热点和远端流量,却会占用额外权重显存,并引入 replica 一致性、路由决策和负载统计问题。推理中多占的显存还可能挤压 KV Cache 与 batch capacity。
8.3 动态与拓扑感知放置#
若 Expert 热度随 workload 变化,可以基于近期路由统计迁移或复制 Expert;若某些 Expert 经常被同一节点的 token 访问,也可以把它们放到近端。
动态方案必须考虑权重搬运本身的成本和稳定性。若 placement 调整周期短于热度变化周期,系统可能频繁迁移却来不及回收收益。实践中通常需要平滑统计、迁移阈值和冷却时间。
9. DeepEP 一类 MoE 通信库优化的是什么#
通用接口通常把问题抽象为“给定 send buffer 与 split size,完成一次 collective”。MoE 专用通信库则能利用更强的先验:数据是 token,目标由 Expert 决定,dispatch 后紧跟 grouped GEMM,combine 还需要恢复 token 顺序。
因此,DeepEP 一类方案优化的是完整链路,而不只是替换一个 all_to_all_single:
routing result
-> token pack
-> intra-node exchange
-> inter-node exchange
-> receive layout / unpack
-> expert compute handoff
-> reverse dispatch
-> weighted combine这类方案常见的设计方向包括:
- 为不均匀 token 数量设计 dispatch/combine kernel;
- 节点内走高速 GPU 互联,节点间走 RDMA 等路径;
- 减少 CPU 参与和 host synchronization;
- 使用 persistent kernel 或设备侧 signal 降低启动延迟;
- 分别提供高吞吐与低延迟模式;
- 暴露 hook 或 event,让 Expert GEMM 与通信重叠;
- 融合 token layout 转换与通信 buffer 管理。
“专用”不代表在所有 workload 上都更快。高吞吐模式适合训练或大 prefill batch;低延迟模式更适合 decode 小消息。硬件拓扑、rank 数、hidden size、Top-K、token 分布和 Expert kernel 都会改变最佳配置。
10. 训练与推理的优化目标不同#
| 维度 | 训练 | 推理 Prefill | 推理 Decode |
|---|---|---|---|
| Token 规模 | 通常大 | 可很大 | 单步通常小 |
| 消息特征 | 较大、吞吐导向 | 中大消息 | 大量小消息、动态变化 |
| 主要目标 | samples/tokens per second | TTFT 与吞吐 | TPOT、ITL 与尾延迟 |
| 常见瓶颈 | 带宽、重排、反向通信 | 带宽与调度干扰 | 启动、同步、负载长尾 |
| Chunk 策略 | 倾向较大 | 按 prompt/batch 调整 | 倾向低延迟小流水 |
| 额外约束 | 反向与梯度、路由训练 | 长 prompt、batching | 动态 batch、每步变化 |
训练中,一个 MoE 层还要在 backward 传播 activation gradient,并处理 Router/Expert 参数的梯度路径。通信次数和 buffer 生命周期比纯 forward 更复杂,但较大的 token 规模也更容易进入带宽饱和区并隐藏 launch 开销。
推理 decode 则常是延迟问题:每一步 payload 很小,却要重复 dispatch 和 combine。此时 persistent kernel、设备侧控制、Expert replication、node-limited routing 和减少 barrier 往往比单纯追求峰值 GB/s 更重要。
11. 如何确认 All-to-All 真的是瓶颈#
首先不要把所有 idle 都归因于网络。一个可操作的排查顺序是:
1. 确认 MoE layer 的端到端占比
2. 拆出 routing / pack / dispatch / expert / combine / unpack
3. 比较各 rank 的 token count、bytes 与阶段耗时
4. 区分节点内、跨节点和不同消息大小
5. 检查通信是否与计算真实重叠
6. 再决定优化网络、kernel、路由还是 placement11.1 应同时记录的指标#
| 类别 | 指标 |
|---|---|
| 路由 | 每 Expert token 数、最大/平均值、Top-K、本地路由比例 |
| 通信量 | 每 peer bytes、跨节点 bytes、padding ratio、metadata bytes |
| 时间 | pack、dispatch、Expert GEMM、combine、unpack、barrier |
| 网络 | effective bandwidth、NIC 利用率、链路热点、重传/拥塞 |
| GPU | SM 利用率、HBM 带宽、kernel gap、stream overlap |
| 服务 | TTFT、TPOT/ITL、TPS、P50/P95/P99 latency |
最大值和分布比平均值更重要。平均每个 Expert 100 个 token,可能实际是一个 Expert 800 个、多个 Expert 接近空闲;平均通信带宽也可能掩盖某个 rank 的拥塞长尾。
11.2 Profile 中容易看到什么#
常见热点包括:
send / recv / all_to_all kernels
token_permute / index_select / scatter / gather
histogram / prefix sum
device-to-device memcpy
stream or event synchronization
grouped GEMM gaps观察 timeline 时重点问四个问题:
- 通信前为什么有空洞,是否在等 CPU 或 count 回读?
- pack 是否产生了多次全量 HBM copy?
- 第一批 token 到达后,Expert GEMM 是否立刻开始?
- 最快 rank 在 barrier 上等待最慢 rank 多久?
11.3 Benchmark 必须保留 workload 维度#
单个大 buffer 的 All-to-All benchmark 只能说明链路上限。评估 MoE 需要至少扫描:
- rank 与 node 数;
- token 数与 hidden size;
- Top-1 / Top-2;
- 均匀、Zipf 热点和真实路由分布;
- BF16/FP16 与低精度通信;
- prefill 与 decode batch;
- 不同 chunk size;
- overlap 开启与关闭;
- 节点内、跨节点及混合拓扑。
同时保留 correctness 检查:token 顺序、Top-K weight、丢弃策略、量化误差和空 Expert 都可能让性能优化产生静默错误。
12. 一组更实用的优化优先级#
All-to-All 优化可以按下面的顺序推进:
第一步:测完整 MoE layer,而不是只测 collective
第二步:消除不必要 payload、padding 和中间 copy
第三步:检查 Router/Expert 的最大负载与跨节点比例
第四步:让 buffer 布局匹配目标 rank 和 Expert GEMM
第五步:按拓扑选择 P2P、pairwise 或 hierarchical 路径
第六步:用 chunk、stream、event 建立 dispatch-compute-combine 流水
第七步:再调 Expert replication、动态 placement 与路由约束这个顺序的核心是先解决确定性的浪费,再处理更复杂的调度和模型取舍。若一次 dispatch 仍有两次无意义全量拷贝,过早调整通信算法通常不会得到稳定收益;若某个 Expert 长期过载,再高的平均网络带宽也消除不了 barrier 长尾。
13. 总结#
All-to-All 可以理解为“每个 rank 都把不同数据发送给不同目标”的分布式转置。在 MoE 中,它承担 token 到 Expert 的 dispatch,以及 Expert 输出回到原 token 位置的 combine。
但真正的优化对象不是某一个 collective kernel,而是整条路径:
Router
-> count / permute / pack
-> topology-aware dispatch
-> Expert compute
-> topology-aware combine
-> unpack / weighted restore最关键的抓手可以归纳为七类:
- 减少 Top-K、精度、padding 或非本地路由带来的通信量;
- 融合 count、pack、unpack 和 combine,减少 HBM 往返;
- 将小消息聚合为合适的 chunk,而不是盲目追求大包;
- 按节点内外拓扑选择 direct、pairwise 或 hierarchical 通信;
- 让部分 token 到达后立即计算,形成通算流水;
- 用负载均衡、Expert replication 和 placement 控制长尾;
- 用端到端 MoE latency、P95/P99 和最慢 rank 评价收益。
最终要降低的不是一个孤立的网络数字,而是:
dispatch latency
combine latency
MoE layer latency
rank-to-rank tail waiting time只有这些指标真正下降,All-to-All 优化才转化成训练吞吐或推理延迟上的端到端收益。