Pidoku

Executor, Worker, Model Runner

Advanced 45 min Difficulty 3/5 Lesson 01 of 05

Prerequisites The Engine Core Loop, One Scheduling Step

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 memory

How 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:

MethodPurpose
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:

ClassWhenHow workers are reached
UniProcExecutorOne GPU and no multiprocessing neededA direct method call in the same process
MultiprocExecutorSeveral GPUs on one node (the usual case for -tp N)The shared-memory broadcast queue from Processes and Wires
RayDistributedExecutorWorkers on several nodesRay 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:

  1. Claim the device. Set the CUDA device from its local_rank, take the memory snapshot used for the budget.
  2. Join the parallel groups. Initialise the tensor-, pipeline- and data-parallel communication groups so that layers can exchange activations.
  3. Load the model. Build the torch.nn.Module, load only this rank’s shard of the weights.
  4. Report and allocate. Answer get_kv_cache_specs, run the memory profile, then allocate the KV cache tensors when told how many blocks.
  5. 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:

Python
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:

StepMethod (older runner)What happens
1_update_statesApply the schedule to the persistent batch: add new requests, drop finished ones, append new block IDs
2_prepare_inputsAssemble this step’s flat tensors: token IDs, positions, per-request offsets (Building the Batch)
3_build_attention_metadataPackage block tables and slot mappings for the attention backend
4_execute_mm_encoderIf any request has an image or audio clip scheduled, run the encoder for it
5_determine_batch_execution_and_paddingChoose a CUDA graph for this batch shape, or eager mode
6_model_forwardThe 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”)
Locationvllm/v1/worker/gpu_model_runner.py, a single 7,500-line filevllm/v1/worker/gpu/, split by feature
Persistent stateThe persistent tensors are the model inputs, so rows must be kept in a required order and compacted when requests leaveEach request owns a fixed row for its lifetime; inputs are gathered from state each step
Input preparationNumPy on the CPU, then copied to the GPUTriton kernels on the GPU
SamplingPyTorch operationsTriton kernels, including a Gumbel-max sampler
Async schedulingRetrofittedAssumed 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:

Python
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:

Python
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.

Go
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_rpc is the general mechanism; execute_model and sample_tokens are 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#

  1. 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?
  2. 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.
  3. Call collective_rpc yourself from the offline API: llm.collective_rpc("get_kv_cache_specs"). How many entries are returned for -tp 2?

Check yourself#

  1. Why does every vLLM model class take the same constructor arguments?
  2. What is wrong with copying a pinned CPU buffer to the GPU asynchronously and then reusing that buffer?
  3. 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.

↑↓ navigate↵ openesc close

drag to pan · scroll to zoom