GenAI System Design
2. Inference, GPUs and Serving

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.

Lesson 6 of 7 11 min

Two different questions

Scaling inference answers two separate questions:

  1. 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.
  2. 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 parallelL1L2L3Each layer split 4 ways.All-reduce after every layer.Cuts latency; stays in one node.Pipeline parallelLayers 1–20Layers 21–40Layers 41–60Layers 61–80Layers split into stages.Activations hop between stages.Fits huge models across nodes.Expert parallel (MoE)RouterE1E2E3E4Experts live on different GPUs.Each token uses only top-k experts.All-to-all traffic between GPUs.
Colours are GPUs. Tensor parallelism talks every layer, so it needs NVLink; pipeline and expert parallelism tolerate slower links.

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 70BMoE, 109B total / 17B active (Llama 4 Scout)
Memory for weights (BF16)140 GB218 GB (all experts must be loaded)
Compute per token70B params≈ 17B params
ImplicationMemory and compute scale togetherMemory-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:

Router
cache-aware
Prefill pool
compute-heavy, builds KV cache
KV transfer
NVLink / RDMA
Decode pool
bandwidth-heavy, big batches
Stream to client

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.

Go deeper