Scaling Out: Tensor, Pipeline and Expert Parallelism
How a model too big for one GPU is split across several, why tensor parallelism stays inside a server, and how replicas, MoE and disaggregated serving fit together.
Two different questions
Scaling inference answers two separate questions:
- How do I serve more traffic? Add replicas: independent copies of the model behind a load balancer. This is plain horizontal scaling, and it is the first answer.
- How do I run a model that doesn't fit on one GPU, or run it faster? Split the model across GPUs. Each replica is then a group of GPUs.
A deployment is described as "R replicas × G GPUs each". For example, 40 replicas of Llama 3.1 70B at TP = 2 is 80 GPUs.
The three ways to split a model
Tensor parallelism (TP)
Each layer's weight matrices are cut into slices, one per GPU. All GPUs work on every token at once, then exchange partial results with an all-reduce after each layer.
- Upside: every GPU reads only 1/TP of the weights per step. Latency per token drops almost linearly, and memory is pooled.
- Downside: two all-reduces per layer, on every step. That needs very fast interconnect. NVLink gives ≈ 900 GB/s per GPU on H100, versus ≈ 50 GB/s for a 400 Gb/s InfiniBand link.
- Rule: keep TP within one server, typically TP = 2, 4 or 8.
Pipeline parallelism (PP)
Layers are divided into consecutive stages on different GPUs or nodes. A token's activations pass from stage to stage, and only one small activation tensor crosses the link per stage.
- Upside: low communication, so it works across nodes. This is how the largest models span several 8-GPU servers.
- Downside: each token still passes through every stage in sequence, so PP doesn't cut per-token latency. Keeping all stages busy needs many requests in flight (micro-batches), otherwise stages sit idle.
Expert parallelism (EP)
In a mixture-of-experts (MoE) model, each layer has many expert sub-networks and a router that sends each token to a few of them, for example the top 2 of 16. Expert parallelism places different experts on different GPUs and shuffles tokens to wherever their experts live, an all-to-all exchange.
MoE changes the capacity maths:
| Dense 70B | MoE, 109B total / 17B active (Llama 4 Scout) | |
|---|---|---|
| Memory for weights (BF16) | 140 GB | 218 GB (all experts must be loaded) |
| Compute per token | 70B params | ≈ 17B params |
| Implication | Memory and compute scale together | Memory-heavy, but fast and cheap per token at high batch |
Combining them
Real deployments mix strategies:
- 70B dense on H100s: TP = 2 (FP8) or TP = 4 (BF16) per replica, then replicas for traffic.
- 400B+ dense: TP = 8 within a node, plus PP = 2 across two nodes.
- Large MoE (hundreds of billions of parameters): EP across 8–64 GPUs, often combined with data parallelism for attention layers.
Disaggregated prefill and decode
Prefill is compute-bound and decode is memory-bound. Running both on the same GPUs means a big prefill arriving mid-batch stalls every user's next token. Disaggregated serving runs them on separate pools:
Each pool can be sized and even use different hardware for its own bottleneck, and TTFT and TPOT stop interfering. The cost is moving the KV cache between pools, which needs fast networking. Frameworks like NVIDIA Dynamo, llm-d and SGLang support this pattern.
Autoscaling and placement
- Scale on queue depth, KV cache usage or TTFT, not CPU. Scaling up takes minutes (pull image, load weights), so scale early and keep a warm buffer.
- Place TP groups on one node, and PP stages on nodes with fast interconnect.
- Multi-region: GPUs are often capacity-constrained per region. Design the gateway to fail over to another region, or to a hosted API, rather than assuming you can always add GPUs where you are.
Key takeaways
- Scale out first with replicas, each a full copy of the model. Split a model across GPUs only when it does not fit, or to cut latency.
- Tensor parallelism splits every layer and needs NVLink-class bandwidth, so it stays within one 8-GPU server.
- Pipeline parallelism splits layers into stages. It tolerates slower links and spans nodes for the largest models.
- Mixture-of-experts models need memory for all experts but compute for only a few. Expert parallelism spreads the experts across GPUs.