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.
Two distinct reasons — they look the same on the outside but need different strategies:
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.
Weights fit, but throughput is the bottleneck. You want more independent instances of the model. Data-parallel with a load balancer.
| Axis | Splits… | Adds… | Wins when… |
|---|---|---|---|
| TP Tensor | each weight matrix across GPUs | all-reduce per layer | single model bigger than one GPU; fast intra-node fabric |
| PP Pipeline | layers across GPUs | point-to-point send/recv between stages | very deep models; cheaper fabric; slower |
| DP Data | requests across full replicas | nothing between GPUs | model fits on one GPU; throughput-bound |
| EP Expert | MoE experts across GPUs | all-to-all token dispatch | MoE 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.
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 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.
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.
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 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.
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.
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).
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.
Mixture-of-Experts models (DeepSeek-V3, Qwen3-MoE, Mixtral) only activate a few experts per token. That means:
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.
Every parallelism axis turns into a collective operation. The fabric decides how fast those collectives run, and therefore which axis is viable.
| Link | Per-GPU bandwidth | Topology | Good for |
|---|---|---|---|
| PCIe 4.0 x16 | ~32 GB/s | host bus | DP, PP, tiny TP (2×) |
| PCIe 5.0 x16 | ~64 GB/s | host bus | TP≤4 workstation |
| NVLink 3 (A100) | 600 GB/s | point-to-point mesh | TP up to 8 |
| NVLink 4 (H100) | 900 GB/s | through NVSwitch | TP up to 8 / node |
| NVLink C2C (Grace-Blackwell) | 900 GB/s CPU↔GPU | integrated | unified-memory loading & Spark clusters |
| NVSwitch / NVL72 | 1.8 TB/s per GPU | all-to-all switch | 72-GPU single domain (GB200 NVL72) |
| InfiniBand NDR 400 | 50 GB/s | Clos fabric | cross-node TP/PP/EP |
| Ethernet RoCE 200G | 25 GB/s | switched | DGX Spark pair / commodity |
| Ethernet 10G | 1.25 GB/s | switched | DP-only across nodes |
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.
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.
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.
| Variable | Effect |
|---|---|
NCCL_DEBUG=INFO | Print topology and ring selection at init — essential first-run check |
NCCL_P2P_DISABLE=1 | Disable PCIe P2P; forces shm. Useful when IOMMU breaks P2P. |
NCCL_IB_DISABLE=1 | Don't use InfiniBand (force TCP) — debugging only |
NCCL_SOCKET_IFNAME=eth0 | Pick a specific NIC on multi-homed hosts |
NCCL_TOPO_FILE=/path/to.xml | Override autodetected topology (rare) |
NCCL_ASYNC_ERROR_HANDLING=1 | Fail fast on lost peer — recommended in production |
NCCL_CUMEM_ENABLE=1 | Use CUDA VMM API for buffers — faster re-allocation |
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.
Every GPU sends a different chunk to every other. Worst-case communication pattern; needs real switched fabric (NVSwitch, NVL72, IB) to not serialise.
Point-to-point. Minimal traffic; works fine over PCIe or even 10 GbE. The cost is pipeline bubbles, not bandwidth.
Replicas are independent processes. Only the load balancer sees them. Network does not matter (beyond routing).
Pick the model and your hardware budget. The planner applies the rules from slides 02–07.
vLLM supports multi-node TP+PP via Ray. A head node starts the Ray cluster; workers join; vLLM spawns TP workers on each 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
ray start --address=<head-ip>:6379 --num-gpus=8
# just sits there; vLLM schedules TP/PP ranks onto it
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.
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.
Disaggregation is worth the complexity at scale (≥ 16 GPUs, mixed long-prompt / short-output workload). Below that, it's overengineering.
nvidia-smi topo -m) before picking flags. Engineers who skip this step debug NCCL hangs for days.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.