The idea in one minute#
The scheduler produces a plan: “these requests, this many tokens each, these blocks.” Three layers turn that plan into arithmetic. The executor delivers the plan to every GPU that holds a piece of the model. The worker is the process that owns one GPU: it loaded the weights and allocated the KV cache. The model runner inside it does the real work each step: it keeps a table of every active request, assembles this step’s input tensors from that table, calls the model, samples, and reports back. The runner is where most of vLLM’s per-token CPU cost is spent, and its central trick — never rebuild what did not change — is worth understanding before reading any tensor code.
A picture#
flowchart TB
EC[":vllm: <b>EngineCore</b><br/><small>scheduler_output</small>"] --> EX[":i-split: <b>Executor</b><br/><small>broadcast one call to all workers</small>"]
EX --> W0
EX --> W1
subgraph W0["Worker, rank 0 — GPU 0"]
direction TB
MR0[":pytorch: <b>Model runner</b>"] --> PB0[(":i-database: <b>Persistent batch</b><br/><small>one row per active request</small>")]
MR0 --> M0[":nvidia: <b>Model shard</b>"]
MR0 --> KV0[(":i-layers: <b>KV cache tensors</b>")]
end
subgraph W1["Worker, rank 1 — GPU 1"]
direction TB
MR1[":pytorch: <b>Model runner</b>"] --> M1[":nvidia: <b>Model shard</b>"]
end
M0 <-->|"all-reduce<br/>every layer"| M1
MR0 -->|"ModelRunnerOutput<br/>rank 0 only"| EC
class EC queue
class EX io
class MR0,MR1,M0,M1 compute
class PB0,KV0 memoryHow it really works#
The executor#
Executor (vllm/v1/executor/abstract.py) is a small interface. The engine core calls a
handful of methods on it and does not care how many GPUs are behind it:
| Method | Purpose |
|---|---|
execute_model(scheduler_output, non_block=True) | Run the forward pass. Returns a future. |
sample_tokens(grammar_output, non_block=True) | Sample from the logits just computed |
determine_available_memory() | Ask each worker how much memory is left for KV cache |
get_kv_cache_specs() | Ask each worker what its layers need |
initialize_from_config(kv_cache_configs) | Tell each worker to allocate its KV cache |
collective_rpc(method, args) | Call any method on every worker and gather the results |
The last one is the general mechanism; the others are built on it. Loading a LoRA adapter,
profiling, sleeping and checkpointing all go through collective_rpc.
Which implementation is used follows from the configuration:
| Class | When | How workers are reached |
|---|---|---|
UniProcExecutor | One GPU and no multiprocessing needed | A direct method call in the same process |
MultiprocExecutor | Several GPUs on one node (the usual case for -tp N) | The shared-memory broadcast queue from Processes and Wires |
RayDistributedExecutor | Workers on several nodes | Ray actors |
With MultiprocExecutor, execute_model is one line of substance: enqueue
("execute_model", (scheduler_output,), {}, output_rank) on the broadcast queue. Every worker
dequeues it and runs it; only output_rank replies.
The worker#
A worker (vllm/v1/worker/gpu_worker.py) is created once per GPU and does its important work
at startup:
- Claim the device. Set the CUDA device from its
local_rank, take the memory snapshot used for the budget. - Join the parallel groups. Initialise the tensor-, pipeline- and data-parallel communication groups so that layers can exchange activations.
- Load the model. Build the
torch.nn.Module, load only this rank’s shard of the weights. - Report and allocate. Answer
get_kv_cache_specs, run the memory profile, then allocate the KV cache tensors when told how many blocks. - Compile and warm up.
torch.compile, capture CUDA graphs, run a dummy batch.
After that the worker is a thin shell. execute_model forwards to the model runner and, under
pipeline parallelism, passes intermediate activations to the next stage.
Weights are sharded while loading, not after#
A design decision in the model classes makes large multi-GPU deployments possible. Every model in vLLM has the same constructor signature:
def __init__(self, *, vllm_config: VllmConfig, prefix: str = ""):and each layer creates only the slice of each weight matrix that its rank owns. The architecture document explains the alternative that was rejected:
Suppose we want to run a 405B model (with roughly 810GB weights) with 16 H100 80GB GPUs. Ideally, every GPU should only load 50GB weights. If we change the model weights after the model is initialized, we need to load the full 810GB weights to every GPU and then shard the weights, leading to a huge memory overhead.
The same applies to quantisation: a layer that will run quantised is built quantised. Nothing is ever loaded at full size and shrunk afterwards.
The model runner’s job each step#
For one execute_model(scheduler_output) call, in order:
| Step | Method (older runner) | What happens |
|---|---|---|
| 1 | _update_states | Apply the schedule to the persistent batch: add new requests, drop finished ones, append new block IDs |
| 2 | _prepare_inputs | Assemble this step’s flat tensors: token IDs, positions, per-request offsets (Building the Batch) |
| 3 | _build_attention_metadata | Package block tables and slot mappings for the attention backend |
| 4 | _execute_mm_encoder | If any request has an image or audio clip scheduled, run the encoder for it |
| 5 | _determine_batch_execution_and_padding | Choose a CUDA graph for this batch shape, or eager mode |
| 6 | _model_forward | The transformer’s forward pass. Returns hidden states. |
| 7 | (kept for sample_tokens) | Select the hidden state of each request’s last token and compute logits |
Then sample_tokens(grammar_output) applies the grammar mask, runs the sampler
(Sampling and Structured Output), and returns a
ModelRunnerOutput holding one or more token IDs per request.
The persistent batch#
Steps 1 and 2 are where a naive implementation loses. Consider what the model needs for 500 running requests: a block table of 500 rows by thousands of columns, 500 temperatures, 500 top-p values, 500 random-number generators. In consecutive steps these are almost identical; perhaps two requests finished and three joined. Rebuilding everything in Python on every token would cost more than the forward pass for a small model.
So the runner keeps state between steps. The Model Runner V2 design document states the idea:
The persistent batch optimization exploits the fact that request batches in consecutive steps are mostly identical. Only a few requests (if any) join or finish per step. By maintaining persistent state tensors and applying incremental diffs instead of reconstructing inputs from scratch, CPU overhead can be reduced significantly.
This is the reason the scheduler sends only differences for requests the worker already
knows (CachedRequestData): the worker has everything else.
Two runners#
There are currently two implementations of the model runner, and the config chooses one at startup:
| Older runner (“MRV1”) | Model Runner V2 (“MRV2”) | |
|---|---|---|
| Location | vllm/v1/worker/gpu_model_runner.py, a single 7,500-line file | vllm/v1/worker/gpu/, split by feature |
| Persistent state | The persistent tensors are the model inputs, so rows must be kept in a required order and compacted when requests leave | Each request owns a fixed row for its lifetime; inputs are gathered from state each step |
| Input preparation | NumPy on the CPU, then copied to the GPU | Triton kernels on the GPU |
| Sampling | PyTorch operations | Triton kernels, including a Gumbel-max sampler |
| Async scheduling | Retrofitted | Assumed from the start |
VllmConfig.use_v2_model_runner decides: V2 is used unless Triton is missing, unless some
configured feature is not yet supported by it (the log says which: “Model Runner V2 does not
yet support …; using the V1 model runner instead”), or unless VLLM_USE_V2_MODEL_RUNNER=0
forces the old one. A few features require V2.
The V2 design is worth reading for three ideas that apply well beyond vLLM.
Fixed rows. Pre-allocate state for max_num_reqs rows. A request takes a free row when it
starts and keeps it until it ends. A preempted request is simply treated as finished and
re-added later. Nothing is ever shifted.
No shared buffer between CPU and GPU. An asynchronous copy to the GPU reads a pinned CPU buffer later. If the CPU overwrites that buffer for the next step first, the GPU receives corrupt data. The older runner guards such buffers with a synchronisation barrier. V2 removes the hazard by copying from a temporary:
self.states[req_idx] = new_req.data
tmp_states = self.states.pin_memory() # a private pinned copy
states = tmp_states.to("cuda", non_blocking=True)The CPU keeps writing to self.states; the GPU reads tmp_states. No barrier, no race.
Send differences, apply them on the GPU. A block table is too large to copy every step.
StagedWriteTensor keeps the table on the GPU, collects this step’s changes on the CPU, packs
them into one contiguous buffer, copies that, and launches a single kernel to apply them:
state = StagedWriteTensor(size=(1024, 1000), dtype=torch.int32, device="cuda")
state.stage_write(row=2, start=3, value=[3, 1, 2])
state.stage_write(row=0, start=1, value=[-1, -2, -5])
state.apply_write()Why the runner must never synchronise#
Everything above serves one requirement: the Python thread must not wait for the GPU. An operation such as reading a GPU tensor’s value into Python forces all queued GPU work to finish first. With async scheduling the engine has already queued the next step; an accidental synchronisation in the runner stalls the pipeline and silently removes the benefit.
The V2 document describes the target as “a CUDA stream with no CPU synchronization points. CPU entrypoints queue work onto the stream.” The one place a synchronisation must happen is fetching sampled token IDs for the scheduler, and it is done once per step, as late as possible.
Code#
How much work does a persistent batch save? This program processes the same request trace two ways: rebuilding every per-request row each step, and applying only what changed.
package main
import (
"fmt"
"math/rand"
)
const (
maxReqs = 1024
steps = 2000
blockSize = 16
)
type req struct {
row int
tokens int // tokens so far
remaining int // tokens still to generate
}
func main() {
rng := rand.New(rand.NewSource(1))
active := map[int]*req{}
var freeRows []int
for i := maxReqs - 1; i >= 0; i-- {
freeRows = append(freeRows, i)
}
var rebuildWrites, diffWrites, joins, leaves int
nextID := 0
for s := 0; s < steps; s++ {
// A few requests arrive each step, with prompts of a few hundred tokens.
for arrivals := rng.Intn(3); arrivals > 0 && len(freeRows) > 0; arrivals-- {
row := freeRows[len(freeRows)-1]
freeRows = freeRows[:len(freeRows)-1]
r := &req{row: row, tokens: 200 + rng.Intn(800), remaining: 50 + rng.Intn(400)}
active[nextID] = r
nextID++
joins++
// Diff: write this request's whole row once (its block IDs so far).
diffWrites += r.tokens/blockSize + 1
}
for id, r := range active {
// Rebuild: copy this request's entire block-table row, every step.
rebuildWrites += r.tokens/blockSize + 1
// Diff: one new block ID, and only when a block boundary is crossed.
r.tokens++
if r.tokens%blockSize == 1 {
diffWrites++
}
r.remaining--
if r.remaining == 0 {
freeRows = append(freeRows, r.row) // the row is reused; nothing shifts
delete(active, id)
leaves++
}
}
}
fmt.Printf("%d steps, %d requests joined, %d finished, %d still active\n",
steps, joins, leaves, len(active))
fmt.Printf("block-table entries written, rebuilding each step: %d\n", rebuildWrites)
fmt.Printf("block-table entries written, applying diffs: %d\n", diffWrites)
fmt.Printf("ratio: %.0fx fewer writes\n", float64(rebuildWrites)/float64(diffWrites))
}The ratio is two orders of magnitude, and block tables are only one of a dozen per-request tensors. This is the saving that makes a Python model runner viable at hundreds of requests per step.
Remember this#
- Executor: deliver one call to all workers. Worker: own one GPU. Model runner: do the step.
collective_rpcis the general mechanism;execute_modelandsample_tokensare the hot path.- Weights are sharded and quantised as they are loaded, never afterwards.
- The runner keeps a persistent batch and applies differences, because consecutive steps are nearly identical.
- Two runners exist. V2 gives each request a fixed row, prepares inputs on the GPU and avoids shared CPU/GPU buffers.
- The runner must never force a GPU synchronisation except to return the sampled tokens.
Try it#
- In the program, make every request very short (
remaining: 2 + rng.Intn(4)), so that the batch turns over almost completely each step. What happens to the ratio? When does a persistent batch stop paying for itself? - Start a server and look in the startup log for which model runner was selected. If it is the older one, the log says which feature caused that.
- Call
collective_rpcyourself from the offline API:llm.collective_rpc("get_kv_cache_specs"). How many entries are returned for-tp 2?
Check yourself#
- Why does every vLLM model class take the same constructor arguments?
- What is wrong with copying a pinned CPU buffer to the GPU asynchronously and then reusing that buffer?
- Under tensor parallelism, how many workers run the forward pass and how many return the sampled tokens?
Sources#
Checked on 5 October 2026 against main at commit 0c16eee.
- Model Runner V2 design
- Architecture overview — the class hierarchy and the 405B sharding example
vllm/v1/executor/multiproc_executor.pyvllm/v1/worker/gpu_worker.pyvllm/v1/worker/gpu_model_runner.pyvllm/config/vllm.py—use_v2_model_runner