The idea in one minute#
A running vLLM server is several operating-system processes, not one. Python can execute only one thread of Python code at a time, so vLLM puts each job that must not wait for another into its own process: text handling in the API server, decisions in the engine core, and arithmetic in one worker per GPU. They are joined by two kinds of wire. Between the API server and the engine core runs ZeroMQ carrying msgpack messages. Between the engine core and its workers runs a shared-memory ring buffer. Knowing which wire carries what tells you how many CPU cores to give the server, why it dies as a unit, and where a slow request is spending its time.
A picture#
flowchart LR
subgraph API["API server process"]
direction TB
IN[":i-network: <b>ROUTER socket</b><br/><small>binds; sends requests</small>"]
OUT[":i-network: <b>PULL socket</b><br/><small>receives outputs</small>"]
end
subgraph CORE["Engine core process"]
direction TB
T1[":i-activity: <b>Input thread</b><br/><small>DEALER socket, decode</small>"] --> Q1[("input_queue")]
Q1 --> LOOP[":i-recycle: <b>Busy loop</b><br/><small>schedule, execute, update</small>"]
LOOP --> Q2[("output_queue")]
Q2 --> T2[":i-activity: <b>Output thread</b><br/><small>encode, PUSH socket</small>"]
end
subgraph WK["Worker processes"]
direction TB
W0[":nvidia: <b>Worker rank 0</b><br/><small>GPU 0</small>"]
W1[":nvidia: <b>Worker rank 1</b><br/><small>GPU 1</small>"]
end
IN -->|"ADD / ABORT / UTILITY"| T1
T2 -->|"EngineCoreOutputs"| OUT
LOOP -->|"shared-memory<br/>broadcast queue"| W0
LOOP --> W1
W0 -->|"response queue"| LOOP
class IN,OUT,T1,T2 io
class Q1,Q2 memory
class LOOP queue
class W0,W1 computeHow it really works#
Why separate processes#
CPython has a global interpreter lock: only one thread runs Python bytecode at a time. Suppose everything lived in one process. While the server tokenised a 100,000-token prompt, the scheduler could not schedule and the GPU would sit idle. While the scheduler ran, no response could be streamed.
Splitting by process gives each job its own interpreter and its own lock:
| Process | Title in ps | Count | Job | What it must never wait for |
|---|---|---|---|---|
| API server | VLLM::APIServer_0 | A, default 1 (default DP with data parallelism) | HTTP, chat templates, tokenising, media loading, detokenising, streaming | The GPU |
| Engine core | VLLM::EngineCore (or VLLM::EngineCore_DP<rank>) | DP, default 1 | Scheduling, KV block accounting | Tokenising or network clients |
| Worker | VLLM::Worker, with suffixes such as _TP0, _PP1, _DP2 | one per GPU: DP × PP × TP | Build tensors, run the model, sample | Anything |
| DP coordinator | VLLM::DPCoordinator | 1 if DP > 1 | Balance load across engine cores | — |
TP, PP and DP are the tensor-, pipeline- and data-parallel sizes
(Parallelism). The total is:
processes = A + DP + (DP × PP × TP) + (1 if DP > 1 else 0)
vllm serve model 1 + 1 + 1 = 3
vllm serve model -tp 4 1 + 1 + 4 = 6
vllm serve model -tp 2 -dp 4 4 + 4 + 8 + 1 = 17warning
The engine core runs a busy loop: when there is work it never sleeps. Give the server
fewer physical CPU cores than it has processes and the engine core gets descheduled by the
operating system in the middle of a step, which shows up as throughput loss with the GPU
under-used. The project’s guidance is at least 2 + N physical cores for N GPUs, and
twice that in vCPUs when hyper-threading is on. In containers, check the CPU limit, not the
node size.
Wire 1: API server ↔ engine core#
ZeroMQ is a library that gives sockets a few useful shapes. vLLM uses four of them:
| Socket | Who owns it | Behaviour used |
|---|---|---|
ROUTER | API server, binds | Sends each message to a specific engine, chosen by the engine’s identity |
DEALER | Engine core, connects | Its identity is the engine index, 2 bytes little-endian |
PUSH | Engine core, one per API server | Sends outputs to exactly the API server that owns the request |
PULL | API server | Receives them |
With everything on one machine the address is a Unix-domain path,
ipc://<VLLM_RPC_BASE_PATH>/<uuid>, where the base path defaults to the system temp directory.
When engines live on other nodes it is TCP, and the API server binds port 0 and tells the
engines which port the kernel assigned.
The first message on the wire goes the “wrong” way. A ROUTER socket cannot address a peer it has never heard from, so on start-up each engine core sends a ready message to every API server. It carries facts the API server needs, which may differ from what was configured: the maximum model length after any auto-fitting, the number of GPU blocks, the block size and the KV cache capacity in tokens. Only then can requests flow.
Every request message is two or more frames:
frame 0 the engine identity (ROUTER uses it to route; it is stripped on arrival)
frame 1 one byte: the request type
frame 2+ the msgpack-encoded body; large tensors travel as extra zero-copy framesThe type byte is defined in vllm/v1/engine/__init__.py:
class EngineCoreRequestType(enum.Enum):
ADD = b"\x00" # a new request
ABORT = b"\x01" # a list of request IDs to cancel
START_DP_WAVE = b"\x02" # data-parallel coordination
UTILITY = b"\x03" # call a method on the engine core and return its result
EXECUTOR_FAILED = b"\x04"
WAKEUP = b"\x05"UTILITY is a small remote-procedure-call mechanism: the body is
(client_index, call_id, method_name, args), the engine core calls the named method on
itself, and the result travels back tagged with the same call_id. Loading a LoRA adapter,
resetting the prefix cache, putting the engine to sleep and profiling all use it.
What is in a request#
The body of an ADD is an EngineCoreRequest. Text never crosses this wire; the engine core
is a tokens-in, tokens-out machine.
class EngineCoreRequest(msgspec.Struct, array_like=True, omit_defaults=True, gc=False):
request_id: str
prompt_token_ids: list[int] | None
mm_features: list[MultiModalFeatureSpec] | None # processed images, audio
sampling_params: SamplingParams | None
pooling_params: PoolingParams | None # for embedding models
arrival_time: float
lora_request: LoRARequest | None
cache_salt: str | None # isolates the prefix cache
data_parallel_rank: int | None
prompt_embeds: torch.Tensor | None = None
client_index: int = 0 # which API server sent it
priority: int = 0
...Three details of that class declaration are performance decisions:
array_like=Trueencodes the struct as a msgpack array, not a map. Field names are never sent; position identifies the field.omit_defaults=Truedrops trailing fields that hold their default.gc=Falsetells Python’s garbage collector not to track these objects. A server holds thousands of them, and tracking them makes every collection pause longer.
The engine core does one more thing at start-up for the same reason: it calls
freeze_gc_heap() to mark everything allocated so far as permanent, so the garbage collector
never scans the model’s Python objects again.
What comes back#
One EngineCoreOutputs per engine step, holding one EngineCoreOutput per request that made
progress:
class EngineCoreOutput(msgspec.Struct, array_like=True, omit_defaults=True, gc=False):
request_id: str
new_token_ids: list[int]
new_logprobs: LogprobsLists | None = None
finish_reason: FinishReason | None = None # STOP, LENGTH, ABORT, ERROR, REPETITION
stop_reason: int | str | None = None
events: list[EngineCoreEvent] | None = None # QUEUED, SCHEDULED, PREEMPTED + timestamps
prefill_stats: PrefillStats | None = None # includes cached-token counts
...finish_reason is an integer on the wire and a word in the API: STOP and LENGTH are the
ones clients see every day. REPETITION is newer; it ends a request whose output has fallen
into a repeating loop. The events list is how latency metrics are built: the engine core
stamps each state change with its own monotonic clock, and the API server computes queue time
and prefill time from the differences.
Inside the engine core: two threads and a loop#
Socket I/O and serialisation release the global interpreter lock, so they can overlap with the GPU. The engine core exploits that with two helper threads:
# Background Threads and Queues for IO. These enable us to
# overlap ZMQ socket IO with GPU since they release the GIL,
# and to overlap some serialization/deserialization with the
# model forward pass.
# Threads handle Socket <-> Queues and core_busy_loop uses Queue.- The input thread receives frames, decodes msgpack, and — importantly — does the
per-request preparation there too (
preprocess_add_request): it builds the scheduler’sRequestobject and computes the prompt’s block hashes for prefix caching. That work overlaps with the forward pass of the step in progress. - The output thread encodes results and sends them, reusing its byte buffers between steps instead of allocating new ones.
- The main thread only ever touches two in-memory queues. It never blocks on a socket.
An abort gets special handling. It is placed on both the ordinary input queue and a separate
aborts_queue, and the main loop drains the second one right after each forward pass. A
cancelled request therefore stops consuming GPU one step later, not after everything queued
ahead of it has been handled.
Wire 2: engine core ↔ workers#
Workers are on the same machine as their engine core, and with tensor parallelism every worker
must receive the same schedule at the same moment. That is a broadcast, and vLLM implements
it as a ring buffer in shared memory (ShmRingBuffer in
vllm/distributed/device_communicators/shm_broadcast.py): one writer, N readers, no locks.
+-------------------------------+------------------------------------+
| chunk 0 | chunk 1 | ... | chunk | meta 0 | meta 1 | ... | meta |
+-------------------------------+------------------------------------+
max_chunks × max_chunk_bytes max_chunks × (1 + n_readers) bytes
each meta = [ written | reader0 | reader1 | ... ] one byte per flagThe writer may reuse a chunk only when it is unwritten or every reader has marked it read. A
reader may read a chunk when written is 1 and its own flag is 0. Because there is a single
writer and each reader touches only its own flag byte, no mutex is needed.
Each message is (method_name, args, kwargs, output_rank): an instruction to call a method on
the worker object. The worker’s whole main loop is:
def worker_busy_loop(self):
while True:
self._execute_worker_rpc(self.rpc_broadcast_mq.dequeue(indefinite=True))Every worker executes the call, but only the worker whose rank equals output_rank sends a
reply. All ranks compute the same sampled tokens, so one answer is enough.
A chunk holds 16 MiB by default (VLLM_MQ_MAX_CHUNK_BYTES_MB); anything larger is sent through
a ZeroMQ side channel. If the engine core waits more than 300 seconds for a reply
(VLLM_EXECUTE_MODEL_TIMEOUT_SECONDS), it treats the executor as failed.
They live and die together#
There is no partial failure mode:
- If a worker dies, the executor’s monitor thread notices, an
EXECUTOR_FAILEDitem is put on the engine core’s input queue, and the busy loop raises. - If the engine core dies, its output thread’s last act is to send the literal bytes
ENGINE_CORE_DEADto every API server. Each API server marks itself errored; every pending and future request fails withEngineDeadError, and/healthstarts failing.
The intended recovery is to restart the whole server, which is what a Kubernetes liveness
probe on /health does. (An optional fault-tolerance mode for multi-engine deployments exists
on main; it is not the default.)
Start methods#
Child processes are started with fork by default (VLLM_WORKER_MULTIPROC_METHOD), which is
fast but unsafe once CUDA has been initialised in the parent. So:
- when started as
vllm serve, vLLM controls the main process and usesspawn; - when used as a library, it uses
fork, unless it detects that CUDA is already initialised, in which case it switches tospawnand logs a warning.
This is why a script that uses the LLM class must put its code under
if __name__ == "__main__":. With spawn, the child re-imports your script, and unguarded
code would start a second engine from inside the first.
Code#
Two small programs. The first answers the sizing question for any topology.
package main
import "fmt"
type topology struct {
name string
tp, pp, dp int
apiServers int // 0 means "use the default"
}
func main() {
cases := []topology{
{name: "1 GPU", tp: 1, pp: 1, dp: 1},
{name: "4 GPUs, one big model (-tp 4)", tp: 4, pp: 1, dp: 1},
{name: "8 GPUs, -tp 2 -dp 4", tp: 2, pp: 1, dp: 4},
{name: "8 GPUs, 8 replicas (-dp 8)", tp: 1, pp: 1, dp: 8},
{name: "8 GPUs, -dp 8, 2 API servers", tp: 1, pp: 1, dp: 8, apiServers: 2},
}
fmt.Printf("%-34s %4s %5s %8s %6s %6s %10s\n",
"deployment", "API", "eng", "workers", "coord", "total", "min vCPUs")
for _, c := range cases {
api := c.apiServers
if api == 0 {
api = c.dp // the API server count defaults to the data-parallel size
}
workers := c.dp * c.pp * c.tp
coord := 0
if c.dp > 1 {
coord = 1
}
total := api + c.dp + workers + coord
// One physical core per process; a vCPU is half a core with hyper-threading.
fmt.Printf("%-34s %4d %5d %8d %6d %6d %10d\n",
c.name, api, c.dp, workers, coord, total, 2*total)
}
}The second is the ring buffer’s flag protocol, for one writer and two readers.
package main
import "fmt"
const readers = 2
type chunk struct {
data string
written byte
read [readers]byte
}
type ring struct {
chunks []chunk
writeIdx int
readIdx [readers]int
}
// canWrite: the chunk was never written, or every reader has consumed it.
func (c *chunk) canWrite() bool {
if c.written == 0 {
return true
}
for _, f := range c.read {
if f == 0 {
return false
}
}
return true
}
func (r *ring) enqueue(msg string) bool {
c := &r.chunks[r.writeIdx]
if !c.canWrite() {
return false // a slow reader is holding the writer back
}
c.written = 0 // order matters: clear, write, reset readers, then publish
c.data = msg
c.read = [readers]byte{}
c.written = 1
r.writeIdx = (r.writeIdx + 1) % len(r.chunks)
return true
}
func (r *ring) dequeue(id int) (string, bool) {
c := &r.chunks[r.readIdx[id]]
if c.written == 0 || c.read[id] == 1 {
return "", false
}
c.read[id] = 1 // each reader writes only its own flag: no lock needed
r.readIdx[id] = (r.readIdx[id] + 1) % len(r.chunks)
return c.data, true
}
func main() {
r := &ring{chunks: make([]chunk, 2)}
fmt.Println("enqueue step-1:", r.enqueue("execute_model(step 1)"))
fmt.Println("enqueue step-2:", r.enqueue("execute_model(step 2)"))
fmt.Println("enqueue step-3:", r.enqueue("execute_model(step 3)"), "<- ring is full")
m, _ := r.dequeue(0)
fmt.Println("worker 0 got:", m)
fmt.Println("enqueue step-3:", r.enqueue("execute_model(step 3)"), "<- worker 1 has not read chunk 0")
m, _ = r.dequeue(1)
fmt.Println("worker 1 got:", m)
fmt.Println("enqueue step-3:", r.enqueue("execute_model(step 3)"), "<- now every reader is done with it")
}The last three lines are the important ones: the writer is blocked by the slowest reader. In vLLM the writer then spins, and after 60 seconds logs “No available shared memory broadcast block found”. If you ever see that message, one worker is stuck — usually a GPU that has hung — and the others are waiting for it.
Remember this#
- Processes:
AAPI servers +DPengine cores + one worker per GPU (+ a coordinator whenDP > 1). - Give it at least one physical core per process; the engine core busy-loops.
- API server ↔ engine core: ZeroMQ ROUTER/DEALER in, PUSH/PULL out, msgpack bodies, token IDs only.
- Engine core ↔ workers: a lock-free shared-memory broadcast of method calls; one rank replies.
- The engine core’s main thread only touches in-memory queues; two threads do the socket work.
- Any process dying kills the server on purpose. Restart is the recovery path.
Try it#
- On a machine running vLLM, run
ps -eo pid,args | grep "VLLM::"and count the processes. (The titles need thesetproctitlepackage, which vLLM installs.) Does the count match the formula? - In the first program, add a case for 16 GPUs with
-tp 8 -dp 2. How many vCPUs does that node need at minimum? - In the ring buffer program, change the ring to 4 chunks and enqueue 5 messages before any reader runs. Which enqueue fails, and what does that say about how far ahead of its workers an engine core can get?
Check yourself#
- Why does the engine core send the first message to the API server, and not the other way round?
- What three options on
EngineCoreRequest’s class declaration make it cheaper to send and hold? - With tensor parallelism of 4, how many workers execute
execute_modeland how many reply?
Sources#
Checked on 5 October 2026 against main at commit 0c16eee.
vllm/v1/engine/core.py—EngineCoreProc, the two I/O threadsvllm/v1/engine/core_client.py— the ROUTER and PULL socketsvllm/v1/engine/__init__.py— the wire typesvllm/distributed/device_communicators/shm_broadcast.py— the ring buffer- Architecture overview: V1 process architecture
- Python multiprocessing in vLLM
- CPU resources for GPU deployments