Pidoku

Parallelism: TP, PP, DP, EP and CP

Expert 50 min Difficulty 4/5 Lesson 01 of 05

Prerequisites Processes and Wires, Sizing the Cache

The idea in one minute#

More GPUs can be used in two fundamentally different ways. You can spread one copy of the model across them because it does not fit, or to make each step faster. Or you can run several copies, one per GPU or group of GPUs, to serve more requests. vLLM names five strategies with two-letter abbreviations, and they are easier to remember by what they split: tensor parallelism splits each layer, pipeline parallelism splits the stack of layers, data parallelism splits the requests, expert parallelism splits a mixture-of-experts model’s experts, and context parallelism splits one long sequence. They combine, each has a flag, and each maps onto the process diagram from the start of this course.

A picture#

flowchart TB
  subgraph TP["Tensor parallel (-tp 2): one model, each layer split"]
    direction LR
    T0[":nvidia: GPU 0<br/><small>half of every layer</small>"] <-->|"all-reduce, every layer"| T1[":nvidia: GPU 1<br/><small>other half</small>"]
  end
  subgraph PP["Pipeline parallel (-pp 2): one model, layers split"]
    direction LR
    P0[":nvidia: GPU 0<br/><small>layers 0–15</small>"] -->|"activations, once per step"| P1[":nvidia: GPU 1<br/><small>layers 16–31</small>"]
  end
  subgraph DP["Data parallel (-dp 2): two models, requests split"]
    direction LR
    LB[":i-route: load balancing"] --> D0[":nvidia: GPU 0<br/><small>full model + own scheduler</small>"]
    LB --> D1[":nvidia: GPU 1<br/><small>full model + own scheduler</small>"]
  end
  class T0,T1,P0,P1,D0,D1 compute
  class LB queue

How it really works#

The five, side by side#

SplitsFlagEngine coresWorkersCommunicationUse it when
TPEach layer’s weights, across GPUs--tensor-parallel-size N (-tp)1NAn all-reduce after every attention and every MLPThe model does not fit on one GPU; GPUs share a node with a fast interconnect
PPThe layer stack, into stages--pipeline-parallel-size N (-pp)1NActivations handed to the next stage once per stepThe model does not fit on one node; uneven splits; no fast interconnect
DPRequests, across full replicas--data-parallel-size N (-dp)NN × TPNone for dense modelsThe model fits; you want throughput
EPA mixture-of-experts model’s experts--enable-expert-parallel (with DP)as DPas DPTokens routed to the GPU holding their expertLarge MoE models
CPOne sequence’s cache or queriescontext-parallel sizes1NShards of KV across GPUsVery long contexts

For the mathematics of each, see Distributed Inference. This lesson is about what vLLM does with them.

Tensor parallelism#

Every layer’s weight matrices are cut by columns or rows across the GPUs, following the Megatron scheme; each GPU computes a partial result and an all-reduce sums them. There are two all-reduces per layer, so a 32-layer model performs 64 collective operations per forward pass.

What that means in vLLM:

  • One engine core, one scheduler, one block table. The workers are replicas of the same execution: every rank receives the same SchedulerOutput and computes the same sampled tokens. Only rank 0 answers.
  • Each GPU holds its share of the KV heads. Cache cost per token per GPU divides by the TP size, so splitting a model over more GPUs increases capacity twice over (Sizing the Cache).
  • Speedup is sub-linear. The arithmetic divides by N; the all-reduces do not. Without a fast link between GPUs (NVLink), communication can dominate. The documentation advises that on GPUs without it — it names the L40S — pipeline parallelism is the better choice.
  • vLLM ships its own all-reduce for small messages, faster than the general library for the sizes inference produces, and can fuse the all-reduce with the following normalisation layer into one kernel (fuse_allreduce_rms, part of -O2).

Pipeline parallelism#

The first GPU runs the first group of layers and passes activations to the second. There is one transfer per stage per step rather than two collectives per layer, so the link matters far less.

Its problem is idleness: while stage 2 computes a batch, stage 1 has nothing to do. vLLM fills the gap with the batch queue from The Engine Core Loop: with -pp 2 the engine keeps two batches in flight, so each stage is always working on a different one. The source calls the empty slots it removes “pipeline bubbles”. A request therefore advances only every other step, which is why per-request latency does not improve; throughput does.

The usual recipe for a model too large for one node is from the documentation: “Set tensor_parallel_size to the number of GPUs per node and pipeline_parallel_size to the number of nodes.”

Data parallelism#

--data-parallel-size 4 starts four complete engine cores, each with its own scheduler, block pool, prefix cache and workers, behind one HTTP endpoint. From Processes and Wires: -tp 2 -dp 4 is 4 API servers, 4 engine cores, 8 workers and a coordinator — 17 processes.

Properties that follow from “N independent engines”:

  • --max-num-seqs is per rank. Four ranks at 256 admit 1,024 running requests in total.
  • Admission limits are per server. --max-num-queued-reqs and --max-num-queued-tokens count across all ranks an API server routes to.
  • Each rank has its own prefix cache. A request benefits from a cached prefix only if it lands on the rank that holds it.

That last point makes load balancing matter. The API server chooses a rank per request using live statistics each engine publishes through the DP coordinator: running and waiting counts, and KV cache usage. For better cache hits, send related requests to the same rank: data_parallel_rank can be set explicitly, and conversations can carry a session identifier that keeps them on one rank.

Two deployment shapes exist:

ModeShapeBalancing
InternalOne vllm serve exposes one endpoint for all ranks. Ranks may be on other nodes, started with --headless.By vLLM, using engine statistics
ExternalEach rank is its own server on its own port.By your router or gateway

The external mode is what Kubernetes-native routers want: they can see every rank and make cache-aware decisions themselves (Deploying on Kubernetes).

Is DP different from just running N servers?#

For a dense model, barely. -dp 4 is four independent engines with a shared front door; you could run four vllm serve processes and a load balancer and get the same throughput.

For a mixture-of-experts model it is very different, and that is what DP was built for.

Expert parallelism#

A mixture-of-experts layer contains many small networks, and each token uses only a few. With --enable-expert-parallel, the experts are distributed across the DP ranks’ GPUs rather than replicated on each. A token whose expert lives elsewhere is sent there and its result sent back.

So the layers of one model use different strategies:

attention layers   data parallel     each rank handles its own requests independently
expert layers      expert parallel   all ranks cooperate on every forward pass

The ranks are no longer independent. From the documentation: “Forward passes must be aligned, and expert layers across all ranks are required to synchronize during every forward pass, even when there are fewer requests to be processed than DP ranks.” A rank with nothing to do must still run an empty dummy forward pass so that its experts are available to the others. Coordinating that is the job of the DP coordinator process, and it is why the engine core’s has_work() includes engines_running: “someone else has work” counts as work.

The token exchange between GPUs is an all-to-all operation with selectable implementations (--all2all-backend), tuned differently for prefill-heavy and decode-heavy workloads.

Because every rank steps together, one rank’s large prefill slows everyone’s step. --prefill-schedule-interval N admits new prefills only every Nth step, aligned across ranks, to keep step times even.

Context parallelism#

For very long sequences. Decode context parallelism shards one request’s KV cache across GPUs, so a single sequence can exceed one GPU’s cache; prefill context parallelism spreads the queries of a long prompt across GPUs to shorten time to first token. The documentation describes the prefill side as “under active development”. It appears in the code as dcp/pcp world sizes, and it is the reason block-size arithmetic contains a dcp_world_size factor.

Choosing#

Does the model fit on one GPU, with room for a useful KV cache?
├── yes → one replica per GPU. Data parallel, or independent servers.
└── no  → Does it fit on one node?
          ├── yes → tensor parallel across the node's GPUs
          │         (pipeline parallel instead if there is no fast interconnect)
          └── no  → tensor parallel within each node, pipeline parallel across nodes
                    (for MoE models: data parallel attention + expert parallel)

One trade-off is worth stating plainly. With eight GPUs and a model that fits on one, -dp 8 gives roughly eight times one GPU’s throughput. -tp 8 gives far less than eight times, but each request finishes faster and the cache is one large shared pool. Use the smallest TP that fits the model and gives acceptable latency, and spend the remaining GPUs on DP.

Sizing reminders#

  • CPU cores. One per process, minimum: A + DP + N (+1) (Processes and Wires).
  • Startup. Every worker loads its shard, compiles and captures graphs. vllm preload and a persisted compile cache matter more as N grows.
  • Failure. By default, any worker dying takes down its engine, and with a single engine the whole server. An optional fault-tolerance mode for multi-engine deployments exists on main.
  • Troubleshooting. For multi-node runs every node needs the same image, model path and package versions. Set VLLM_HOST_IP when a node has several network interfaces.

Code#

Given a model, a GPU type and a GPU count, compare the layouts: does it fit, how many tokens of cache, and a rough relative throughput and latency.

Go
package main

import "fmt"

const GiB = 1 << 30

type layout struct{ tp, dp int }

func main() {
	const (
		gpuGiB     = 80.0
		util       = 0.92
		otherGiB   = 1.0
		gpus       = 8
		tpOverhead = 0.12 // fraction of a step lost to communication per doubling of TP
	)
	type model struct {
		name        string
		weightGiB   float64
		kvPerTokKiB float64 // single-GPU KV bytes per token, in KiB
	}
	models := []model{
		{"8B model", 15, 128},
		{"70B model", 131.5, 320},
		{"235B model", 440, 440},
	}
	layouts := []layout{{1, 8}, {2, 4}, {4, 2}, {8, 1}}

	for _, m := range models {
		fmt.Printf("%s (%.0f GiB of weights) on %d x %.0f GiB GPUs\n", m.name, m.weightGiB, gpus, gpuGiB)
		fmt.Printf("  %-12s %6s %16s %14s %12s\n", "layout", "fits", "KV tokens/replica", "throughput", "step time")
		for _, l := range layouts {
			perGPUWeights := m.weightGiB / float64(l.tp)
			free := gpuGiB*util - perGPUWeights - otherGiB
			name := fmt.Sprintf("-tp %d -dp %d", l.tp, l.dp)
			if free <= 0 {
				fmt.Printf("  %-12s %6s\n", name, "no")
				continue
			}
			tokens := free * GiB / (m.kvPerTokKiB * 1024 / float64(l.tp))

			// One replica on tp GPUs: compute divides by tp, communication grows.
			doublings := 0
			for t := l.tp; t > 1; t /= 2 {
				doublings++
			}
			step := 1/float64(l.tp) + tpOverhead*float64(doublings) // relative to tp=1
			throughput := float64(l.dp) / step                      // replicas / step time

			fmt.Printf("  %-12s %6s %16.0f %13.1fx %11.2fx\n", name, "yes", tokens, throughput, step)
		}
		fmt.Println()
	}
	fmt.Println("throughput and step time are relative to one GPU running the model at -tp 1;")
	fmt.Println("for models that do not fit at -tp 1 they are relative to that hypothetical.")
}

For the model that fits on one GPU, eight replicas win on throughput by a wide margin and higher TP only shortens the step. For the 70B model the single-GPU layout is gone and the two-GPU one leaves almost no cache; for the 235B model only one layout fits at all. Fitting decides first; preference comes second.

Remember this#

  • TP splits each layer (one scheduler, N workers, two all-reduces per layer).
  • PP splits the layer stack; the batch queue keeps every stage busy.
  • DP runs N complete engines behind one endpoint, each with its own scheduler and prefix cache.
  • EP distributes a mixture-of-experts model’s experts across DP ranks, which makes the ranks step together.
  • CP shards one long sequence across GPUs.
  • max_num_seqs is per DP rank; admission limits are per server.
  • Use the smallest TP that fits and meets latency; spend the rest of the GPUs on DP.
  • Without a fast interconnect, prefer PP to TP.

Try it#

  1. Set tpOverhead to 0.02 (an excellent interconnect) and to 0.4 (none). How does the best layout for the 8B model change?
  2. Change the GPUs to 48 GiB. Which layouts fit the 70B model now?
  3. On a two-GPU machine, run the same benchmark against -tp 2, -dp 2 and -pp 2. Record throughput and time per output token for each.

Check yourself#

  1. Under tensor parallelism of 4, how many schedulers and how many block pools exist?
  2. Why must an idle data-parallel rank still run forward passes for a mixture-of-experts model with expert parallelism?
  3. Eight GPUs, a model that fits on one. Which layout maximises throughput, and what do you give up?

Sources#

Checked on 5 October 2026 against vLLM v0.30.0 and main at commit 0c16eee.

↑↓ navigate↵ openesc close

drag to pan · scroll to zoom