Skip to content

分布式计算 — 概念

并行策略概览

Tensor Parallelism(TP)

将模型的权重矩阵按列或行切分到多个 GPU:

ColumnParallelLinear

Y = XW, W 按列切分为 [W₁, W₂]
GPU 0: Y₁ = XW₁
GPU 1: Y₂ = XW₂
Y = [Y₁, Y₂]  (无需通信)

RowParallelLinear

Y = XW, W 按行切分为 [W₁; W₂]
GPU 0: Y₁ = X₁W₁  (X 按列切分)
GPU 1: Y₂ = X₂W₂
Y = Y₁ + Y₂  (需要 All-Reduce)

TP 的通信开销

每个 Transformer 层需要 2 次 All-Reduce:

  • Attention 后的 output projection
  • MLP 后的 output

通信量与 hidden_size 成正比,与 sequence length 无关。

序列并行 MoE(reduce-scatter)

对 GLM-5.2 / DeepSeek V2 系 MoE,新增 use_sequence_parallel_moe(PR #46635):把 MoE 输出的聚合从旧模式(每 TP rank 切一块 hidden → 本地算 MoE → all-gather 拼回)改为(全 hidden → 算 MoE → tensor_model_parallel_reduce_scatter)。每个 rank 只发送 hidden_dim/tp_size 而非整段 hidden_dim,把 MoE 输出通信量减半,端到端吞吐提升 3.1~3.2%。实现上 o_proj 与注意力设 reduce_results=False(延迟 all-reduce),注意力输出后把 hidden pad 到 tp 整除再 reduce-scatter,MoE 以 already_sequence_parallel=True 跳过自身切分/all-gather。

Pipeline Parallelism(PP)

将模型的不同层分配到不同 GPU:

PP 的挑战

  • 气泡(bubble):GPU 之间存在等待时间
  • 微批次(micro-batching):将请求分成小批次以减少气泡
  • KV 缓存:需要跨 GPU 传递 KV 缓存

Data Parallelism(DP)

复制整个模型到多个 GPU,每个 GPU 处理不同的请求:

  • 不需要模型权重通信
  • 需要请求在 GPU 间均衡分配
  • 适合小模型高并发场景

Expert Parallelism(EP)

MoE 模型中,将专家分布到不同 GPU:

  • 需要 All-to-All 通信:token 发送到对应专家的 GPU
  • 负载均衡是关键挑战
  • 支持弹性专家并行(动态调整专家分配)

Elastic EP(弹性专家并行)

Elastic EP 允许在线扩展/缩减 EP 大小而不中断服务:

关键机制:

  • Staged 量化:在 scale-up 前预创建替换 MoE 内核
  • NIXL EP All2All:新 rank 先以 masked 状态加入集体通信,commit 时 unmask
  • 即时参与:新 rank 可以立即加入 EP 集体通信

EPLB(专家级负载均衡)

EPLB 在线迁移专家权重以平衡各 GPU 负载。异步 EPLB 在 v0.23 已成为默认use_async=True,非阻塞迁移),并扩展支持 DeepSeek V4 的 Mega MoE。通信后端有多种选择:torch_nccl / torch_gloo / nixl / pynccl。其中 NixlEplbCommunicator 基于 NIXL READ 实现零拷贝传输,性能最佳;注意异步 EPLB 与 NCCL/pynccl 多流冲突,必须用 nixltorch_gloo。弹性 EP 场景下 NIXL EPLB 通信器会延迟远程 agent 元数据交换(defer_remote_setup)到首次 set_transfer_context

负载记录的 padding 掩码(PR #38128):EPLB 依据 router 的 expert load 统计决定迁移,但被 padding 的空 token 也会被路由并计入 load,导致负载被高估、触发不必要的迁移。修复为每 ubatch 预分配一次 num_unpadded_tokens int32 张量(CUDA graph 安全),每次前向前由 EplbState.prepare_forward() 填入真实 token 数,Triton router kernel 经 HAS_NUM_UNPADDED 编译期标志在 atomic_add 前掩掉越界 token。EPLB 也扩展到推测解码的草稿模型(DraftModelSpeculator.set_eplb_state(),AutoRegressive/DFlash speculator 初始化时调用),并新增 compute_hash_cached() 避免每步重算昂贵的 config hash。

NIXL EP All2All 重构

NixlEPAll2AllManager 引入了分阶段生命周期,_NixlEPBufferState 现维护三态:

  1. _stage_ep_size():仅连接 standby rank,保持 maskedconnected_ep_size 已连接含 masked,但不活跃)
  2. _commit_staged_state():通过 update_mask_buffer(rank, mask=False) 逐个 unmask,提升到 active_ep_size
  3. scale-down:延迟断开旧 rank 直到 commit
  4. 支持跨异构 MoE 层的 workspace 增长

WideEP + DeepEP V2

WideEP 后端集成了 DeepEP V2(DeepEPV2PrepareAndFinalize),两种 dispatch 模式按阶段选择:

  • Decodedo_expand=False(cudagraph 安全,recv buffer 预分配为最坏情况 R × num_max_tokens_per_rank),do_cpu_sync=False
  • Prefilldo_expand=True(按实际接收量分配省内存),do_cpu_sync=True

DeepEP V2 会为每个不同的 num_max_tokens_per_rank 值 JIT 编译独立 dispatch kernel,因此将 token 数向上取整为 2 的幂以限制重编译次数。此外,NIXL EP 现已支持与 DBO(Dynamic Batch Overlap)同时启用,之前二者互斥的限制已解除。

容错框架(Fault Tolerance,DP+EP 外部 LB)

v0.26 为 MoE DP+EP 外部负载均衡部署新增容错框架(#44428/#43637),使单个 engine 出错时整个部署不崩溃、可被外部恢复。开启 enable_fault_tolerance(需 --data-parallel-external-lb 外部 LB 模式)后,run_busy_loop@fault_tolerant_wrapper 装饰,恢复窗口由 FaultToleranceConfig.engine_recovery_timeout_sec(默认 120s)控制:

  • EngineCoreSentinelvllm/v1/fault_tolerance/engine_core_sentinel.py):on_fault() 把所有请求置 FINISHED_ABORTED、清空 batch_queue、状态转 UNHEALTHY/DEAD_push_status(),随后阻塞在 resumed.wait(timeout) 上等恢复指令;指令经输入 socket 循环里的 FT_UTILITY_METHODhandle_fault_tolerance)分支到达。
  • retry:用 stateless 进程组 API(stateless_destroy/init_torch_distributed_process_group,新端口经 dp_store 交换、递增 _dp_reinit_epoch)重建 DP 进程组,重置 step_counter,再经 executor.collective_rpc("handle_ft_command") 把恢复动作下发到各 worker,最后翻回 HEALTHY
  • WorkerSentinelvllm/v1/worker/sentinel/gpu_worker_sentinel.py,经 collective_rpc 在各 worker 执行):retry() 同步加速器、_clean_worker_state()(清 execute_model_state、驱逐所有 req、condense batch)、get_ep_all2all_manager().clean_buffers()、重建 DP CPU group。构造器强制要求支持容错的 all2all 后端(仅 deepep_low_latency / nixl_ep)。
  • 对端故障检测(#43637):DeepEPLLAll2AllManager / NIXL-EP manager 新增 query_active_mask() / query_fault()(基于 DeepEP low_latency_query_mask_buffer),GPU runner 据此在输出前丢弃因对端 rank 掉线而被污染的 batch,避免发出损坏的 token。
  • 控制端点entrypoints/serve/fault_tolerance/api_router.pyPOST /fault_tolerance/apply(仅 retry 指令白名单)。device communicator 也新增 process-checkpoint 生命周期钩子(#46877,首见于 FlashInfer),为状态保存/恢复提供基础。

DCP(Decode Context Parallelism)

DCP(解码上下文并行,--decode-context-parallel-size N / -dcp N)在不增加 world size 的前提下把一个 TP 组切成 tp_size // dcp_size 个 DCP 子组,每个子组只持有 KV cache 的一份切片,从而在 decode 阶段跨 GPU 分摊长上下文的 KV,支持更长的上下文。动机在 MLA:KV head 数 H 很小(通常 1),TP 部署下每张卡的 KV cache 复制 tp_size/H 份;DCP 改在 token 维分片——本 rank 的一个 block 覆盖全局序列的 block_size × dcp_size 个 token,decode 时 query 经 all-to-all 路由到 KV owner rank,各 rank 算局部 output + LSE 再用 log-sum-exp 加权合并。--dcp-comm-backendag_rs(all-gather + reduce-scatter)或更新的 a2a(all-to-all)。集合通信实现收敛到 vllm/v1/attention/ops/dcp.pyMLADCPManager(#52839 合并 dcp_alltoall 与 common 的 CP 部分):combine 优先用对称内存直连 A2A(DirectDCPA2AWorkspace),否则退到 dcp_a2a_lse_reduce / PCP 的 cp_lse_ag_out_rs/cp_lse_ag_out_ar;相关 Triton kernel 已迁入 JIT warmup(#53564)。MLA decode 路径现已支持 FP8 KV cache(model_executor/layers/attention/mla_attention.py)—— 移除了原先"DCP 不支持 FP8 / scaled KV cache"的硬断言,改为按 chunk gather 并逐序列反量化。

DCP × 前缀缓存也已打通(Kimi K3 混合模型):resolve_dcp_kv_block_size 把 full-attention/MLA spec 的 block_size 乘 dcp_world_size(Mamba/SW 保持复制),各 cache manager 与 find_longest_cache_hit 带上自己的 dcp_world_size/pcp_world_size,partial hash hit 不再要求 dcp_world_size==1(#50493/#53324/#53598);MooncakeStore 亦支持混合 DCP 远端命中(见 topics/kv-cache/)。

PCP(Prefill Context Parallelism)

与 DCP 对偶的 prefill 侧序列并行(--prefill-context-parallel-size Nconfig/parallel.py 定义为"切分 prefill 序列计算的 rank 数,扩大 world size 但不增加 KV-cache 分片数")。运行时由 MRV2 的 PCPManagerv1/worker/gpu/pcp_manager.py)管理:持有全局 batch,把每步 InputBatch 重写成 rank-local 的 DualChunkSwap 行,采样/后处理前再还原全局形状(RankSegment 记录各 rank 分段);无 PCP 时 DCP 复用 TP rank。CUDA graph 侧配套 PIECEWISE 模式的持久输入 buffer / slot mapping(#53515/#53869)。PCP×DP 组合目前不支持——校验从全局配置移到了设备管理层(platforms/cuda.py / rocm.py,#54523),兼容性检查委托 PCPManager(#53853)。DCP 的 dcp_kv_cache_interleave_size 已被 cp_kv_cache_interleave_size 取代,将在 PCP 完全就绪后废弃。

通信后端演进(all-reduce)

  • FlashInfer all-reduce 默认开启(#52998,distributed/device_communicators/flashinfer_all_reduce.py),SM103 阈值调优(#53318/#53606)。
  • 跨节点 MNNVL custom all-reduce 按组能力 gate(#53253,Lamport mailbox 修复 #53000)。
  • 新增 opt-in 的 FlashInfer PCIe IPC all-reduce(#53576,flashinfer_pcie_ipc_all_reduce.py)。
  • ROCm 侧 fused allreduce + GemmaRMSNorm fast path(#54787,MiniMax M3:MLC allreduce 与 GemmaRMSNorm 熔合,fused_allreduce_gemma_rms_norm.py);NVIDIA 侧对应物是 Kimi K3 的 CuTeDSL latent-MoE tail 熔合(reduce + RMSNorm + reduce-scatter + early-exit + Lamport 无锁同步,models/kimi_k3/nvidia/ops/cute_dsl/latent_moe_tail/)。SM90+ 还有 PDL(Programmatic Dependent Launch,内核提前放行依赖内核并重叠执行)逐步铺开:fusedQKNormRopeKernel(#55755,小 batch launch-bound 时与 QKV GEMM 重叠)、fused QK-RMSNorm 权重预取(#55020)。

并行策略组合

组合适用场景通信量
TP=2单节点 2 GPU中等
TP=4单节点 4/8 GPU较高
TP=8 + PP=2双节点 16 GPU中等
TP=8 + EP=8MoE 模型 64 GPU中高
DP=4高并发小模型

分布式 KV 缓存传输

vLLM 支持 KV 缓存在 GPU 间传输(PD 分离部署):

  • Disaggregated Prefill:prefill 和 decode 在不同 GPU 上执行
  • NIXL 连接器:拆分为 Pull(D 端 READ 拉取)与 Push(P 端 WRITE 直接写入 D 端显存)两种模式(详见 topics/kv-cache/)。连接器按 region 区分 KV 布局 —— MLA region 用 REPLICATE(整块、key-only),full-attn region 用 SPLIT(按 TP head 切片),使混合 attention 架构传输正确。kv_role='kv_both' 已进入弃用周期,应使用 kv_producer / kv_consumer
  • Mooncake 连接器MooncakeStoreConnector 支持多组 KV cache,并新增流水线并行(PP)下的 PD 分离(PP-aware handshake 聚合中间 stage 输出)与异步 lookup;开启 DCP(>1)时 lookup key prefix 必须显式枚举所有 (tp_rank, pcp_rank, dcp_rank, pp_rank) 命名空间(DCP 复用 TP worker、KV-head 去重不再适用),否则会漏掉 DCP rank 1+ 的块(详见 topics/kv-cache/
  • KV Events:通过事件通知机制协调 KV 缓存的生产和消费,OffloadingConnector 支持自描述事件

相关概念