Local LLM Hosting Series — Presentation 05

Multi-GPU Parallelism for Serving

Tensor, pipeline, data, and expert parallel — what they do, what they cost, and when to use each. NVLink, NVSwitch, NCCL, InfiniBand, and the realities of scaling across boxes.

Tensor Parallel Pipeline Parallel Data Parallel Expert Parallel NCCL NVLink InfiniBand
Why → TP → PP → DP → EP → Fabric → Multi-node
00

Topics We'll Cover

01

Why Parallelise at All?

Two distinct reasons — they look the same on the outside but need different strategies:

The model doesn't fit

70B FP16 = 140 GB. No single GPU on the market except an H200/B100 can hold that. You need to shard the weights. Tensor-parallel or pipeline-parallel.

One GPU can't keep up

Weights fit, but throughput is the bottleneck. You want more independent instances of the model. Data-parallel with a load balancer.

The four axes

AxisSplits…Adds…Wins when…
TP Tensoreach weight matrix across GPUsall-reduce per layersingle model bigger than one GPU; fast intra-node fabric
PP Pipelinelayers across GPUspoint-to-point send/recv between stagesvery deep models; cheaper fabric; slower
DP Datarequests across full replicasnothing between GPUsmodel fits on one GPU; throughput-bound
EP ExpertMoE experts across GPUsall-to-all token dispatchMoE models (Mixtral, DeepSeek-V3, Qwen3-MoE)

Real deployments combine these: world_size = TP × PP × DP × EP. For 70B on 2 GPUs: TP=2, PP=1, DP=1. For DeepSeek-V3 on 16 GPUs: TP=2, PP=4, EP=2 is one valid choice.

02

Tensor Parallel (TP)

Split each weight matrix across GPUs by columns (or rows) and reconstruct the output with an all-reduce. Megatron-LM's attention and MLP sharding is now the industry standard and is what vLLM uses when you pass --tensor-parallel-size N.

TP=2 on one attention + MLP block X (hidden) W_qkv [shard 0] · GPU0 W_qkv [shard 1] · GPU1 attn head 0..H/2 attn head H/2..H W_o [shard 0] partial W_o [shard 1] partial all-reduce SUM Y (post-attn) W_up [col 0..N/2] W_up [col N/2..N] SiLU · GPU0 SiLU · GPU1 W_down[row 0..N/2] W_down[row N/2..N] all-reduce SUM Two all-reduces per transformer layer (one after attn, one after MLP). For Llama-70B with 80 layers → 160 all-reduces per token — fabric matters.
Key constraint

TP world size must evenly divide the number of attention heads (GQA: key/value heads). Llama-3.1-70B has 64 heads / 8 KV heads → TP can be 1, 2, 4, 8. Not 3 or 5. vLLM checks this at load and errors out with a clear message if you pick a bad factor.

03

Pipeline Parallel (PP)

Split the model by layer: GPU 0 runs layers 0–19, GPU 1 runs 20–39, etc. Each mini-batch flows through the pipeline. Inter-GPU traffic is only the hidden state at the stage boundary — tiny compared to TP's all-reduces.

When PP beats TP

  • No NVLink between GPUs (PCIe-only box, or cross-node)
  • Very deep models (ultra-large reasoning models)
  • Heterogeneous GPUs (different memory per stage)
  • Multi-node where TP over the network would be fatal

When PP hurts

  • Serving a single stream — the "pipeline bubble" means only one stage is busy at a time
  • Low-concurrency workloads — bubbles dominate
  • Strict low-latency SLOs — PP adds stage-count hops to TTFT

The Bubble Problem

With PP=4 and one request at a time, only 25% of GPUs are active. vLLM fills bubbles with micro-batches: multiple sequences in-flight simultaneously, each at a different pipeline stage. This is fine for throughput-oriented serving; it hurts latency on a single slow request.

vLLM with PP on 4 GPUs
vllm serve meta-llama/Meta-Llama-3.1-70B-Instruct \
    --pipeline-parallel-size 4 \
    --tensor-parallel-size 1 \
    --gpu-memory-utilization 0.92

In practice, PP is usually combined with TP: TP=2 within a node, PP=N across nodes. See slide 09.

04

Data Parallel (DP) & Replica Routing

If the model fits on one GPU, the simplest scale-out is: launch N identical vLLM replicas and put a load balancer in front. No cross-GPU communication at all.

DP=4 — independent replicas behind a router with KV-aware stickiness clients N.. router / LB (session-stick) GPU 0 · vllm replica 0 · KV cache 0 GPU 1 · vllm replica 1 · KV cache 1 GPU 2 · vllm replica 2 · KV cache 2 GPU 3 · vllm replica 3 · KV cache 3 Linear scale in throughput, near-zero change in per-stream latency, no fabric cost.
Session stickiness matters

KV cache is per-replica. If your router load-balances round-robin, every multi-turn follow-up will land on a different replica and pay full prefill again. Stick sessions to replicas by API key, conversation ID, or hashed prefix. vLLM's --override-neuron-config is not this — use a real LB (Envoy, NGINX with hash, or a purpose-built router like LiteLLM's).

vLLM's own DP mode

Since vLLM 0.6, the engine has a built-in data-parallel mode: --data-parallel-size N. It spawns N copies in one process group and does prefix-cache–aware routing between them. Useful on a DGX Spark pair (two independent nodes with GPU-aware router) or on a single 8-GPU box where you want 4 instances at TP=2.

05

Expert Parallel (EP) for MoE

Mixture-of-Experts models (DeepSeek-V3, Qwen3-MoE, Mixtral) only activate a few experts per token. That means:

EP=4 — 64 experts sharded 16-per-GPU, all-to-all for dispatch and combine. tokens (batch) router (top-k) GPU0 · E0..15 GPU1 · E16..31 GPU2 · E32..47 GPU3 · E48..63 combine (all-to-all) All-to-all volume scales with batch × top-k. On 200 Gb/s per-GPU link it's < 1 ms for typical batches; on PCIe-only it can dominate, which is why DeepSeek-V3-scale MoE wants NVLink or IB.
vLLM with EP on 8 GPUs (Mixtral-8x7B)
vllm serve mistralai/Mixtral-8x7B-Instruct-v0.1 \
    --tensor-parallel-size 2 \
    --data-parallel-size 4 \
    --enable-expert-parallel \
    --dtype bfloat16 \
    --gpu-memory-utilization 0.92

Rule of thumb: use EP only when the active-params / total-params ratio is low (< 0.4). Below that ratio, EP wins throughput hugely because you're not paying for idle experts. Above it, TP is simpler.

06

The Fabric — NVLink, NVSwitch, PCIe, InfiniBand

Every parallelism axis turns into a collective operation. The fabric decides how fast those collectives run, and therefore which axis is viable.

LinkPer-GPU bandwidthTopologyGood for
PCIe 4.0 x16~32 GB/shost busDP, PP, tiny TP (2×)
PCIe 5.0 x16~64 GB/shost busTP≤4 workstation
NVLink 3 (A100)600 GB/spoint-to-point meshTP up to 8
NVLink 4 (H100)900 GB/sthrough NVSwitchTP up to 8 / node
NVLink C2C (Grace-Blackwell)900 GB/s CPU↔GPUintegratedunified-memory loading & Spark clusters
NVSwitch / NVL721.8 TB/s per GPUall-to-all switch72-GPU single domain (GB200 NVL72)
InfiniBand NDR 40050 GB/sClos fabriccross-node TP/PP/EP
Ethernet RoCE 200G25 GB/sswitchedDGX Spark pair / commodity
Ethernet 10G1.25 GB/sswitchedDP-only across nodes
The one rule

TP stays inside one NVLink domain. The moment you cross PCIe or (worse) the network with TP, tok/s falls off a cliff — sometimes 5–10×. PP and DP tolerate slower fabric; TP does not. If you must scale a single model across nodes, put TP inside each node and PP between them.

Checking your topology

See what NVIDIA thinks you have
nvidia-smi topo -m       # matrix: NV#, SYS, PIX, NODE tags
nvidia-smi nvlink -s     # link state per lane
ibstat                   # InfiniBand HCAs
nvidia-smi -q -d TOPOLOGY | head -50

If nvidia-smi topo -m shows SYS (UPI / CPU-to-CPU) between two GPUs you want to TP across, do not TP — use two independent replicas (DP) instead.

07

NCCL — What Actually Moves Bytes

vLLM, TGI, SGLang, and TensorRT-LLM all sit on top of NCCL (NVIDIA Collective Communications Library). NCCL is the thing that actually does all-reduce, all-gather, all-to-all, broadcast, and point-to-point on GPU memory. It picks the best available transport automatically: NVLink → PCIe P2P → shared memory → IB/RDMA → sockets.

Environment variables you'll reach for

VariableEffect
NCCL_DEBUG=INFOPrint topology and ring selection at init — essential first-run check
NCCL_P2P_DISABLE=1Disable PCIe P2P; forces shm. Useful when IOMMU breaks P2P.
NCCL_IB_DISABLE=1Don't use InfiniBand (force TCP) — debugging only
NCCL_SOCKET_IFNAME=eth0Pick a specific NIC on multi-homed hosts
NCCL_TOPO_FILE=/path/to.xmlOverride autodetected topology (rare)
NCCL_ASYNC_ERROR_HANDLING=1Fail fast on lost peer — recommended in production
NCCL_CUMEM_ENABLE=1Use CUDA VMM API for buffers — faster re-allocation

Collective operations at a glance

TP → all-reduce

Every GPU contributes a partial sum; result is summed and broadcast. Dominated by ring-algorithm latency for small payloads and by bandwidth for large ones. This is why two all-reduces per layer make TP fabric-sensitive.

EP → all-to-all

Every GPU sends a different chunk to every other. Worst-case communication pattern; needs real switched fabric (NVSwitch, NVL72, IB) to not serialise.

PP → send/recv

Point-to-point. Minimal traffic; works fine over PCIe or even 10 GbE. The cost is pipeline bubbles, not bandwidth.

DP → nothing

Replicas are independent processes. Only the load balancer sees them. Network does not matter (beyond routing).

08

Interactive: Plan a Parallelism Strategy

Pick the model and your hardware budget. The planner applies the rules from slides 02–07.

16
09

Multi-Node Deployment Patterns

Pattern A — Ray-launched multi-node vLLM

vLLM supports multi-node TP+PP via Ray. A head node starts the Ray cluster; workers join; vLLM spawns TP workers on each node.

head node
ray start --head --port=6379 --num-gpus=8
vllm serve meta-llama/Meta-Llama-3.1-405B-Instruct-FP8 \
    --tensor-parallel-size 8 \
    --pipeline-parallel-size 2 \
    --distributed-executor-backend ray \
    --gpu-memory-utilization 0.92
worker node
ray start --address=<head-ip>:6379 --num-gpus=8
# just sits there; vLLM schedules TP/PP ranks onto it

Pattern B — DP across nodes, single vLLM each

If the model fits in a node, don't bother with cross-node TP. Run one vLLM per node, front with a router (LiteLLM, Envoy). Nodes fail independently; scaling is linear.

Pattern C — Disaggregated prefill

Cutting-edge vLLM / SGLang / DistServe pattern: put prefill on one pool of GPUs and decode on another. Prefill is compute-bound and benefits from TP; decode is memory-bandwidth-bound and benefits from DP. A KV-cache transfer moves the computed prefix from the prefill pool to the decode pool.

Prefill pool — 2× H100, TP=2, optimised for prompt compute
↓
KV transfer (NIXL / Mooncake, over RDMA)
↓
Decode pool — 8× H100, DP=8, optimised for tok/s per stream

Disaggregation is worth the complexity at scale (≥ 16 GPUs, mixed long-prompt / short-output workload). Below that, it's overengineering.

10

Summary & Rules of Thumb

Next

Deck 06 compares the frameworks head-to-head. Deck 07 goes into quantisation, which is how you make "the model fits on one GPU" true in the first place.