Pidoku

Processes and Wires

Basic 50 min Difficulty 2/5 Lesson 01 of 04

Prerequisites The Map

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 compute

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

ProcessTitle in psCountJobWhat it must never wait for
API serverVLLM::APIServer_0A, default 1 (default DP with data parallelism)HTTP, chat templates, tokenising, media loading, detokenising, streamingThe GPU
Engine coreVLLM::EngineCore (or VLLM::EngineCore_DP<rank>)DP, default 1Scheduling, KV block accountingTokenising or network clients
WorkerVLLM::Worker, with suffixes such as _TP0, _PP1, _DP2one per GPU: DP × PP × TPBuild tensors, run the model, sampleAnything
DP coordinatorVLLM::DPCoordinator1 if DP > 1Balance 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      = 17

warning

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:

SocketWho owns itBehaviour used
ROUTERAPI server, bindsSends each message to a specific engine, chosen by the engine’s identity
DEALEREngine core, connectsIts identity is the engine index, 2 bytes little-endian
PUSHEngine core, one per API serverSends outputs to exactly the API server that owns the request
PULLAPI serverReceives 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 frames

The type byte is defined in vllm/v1/engine/__init__.py:

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

Python
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=True encodes the struct as a msgpack array, not a map. Field names are never sent; position identifies the field.
  • omit_defaults=True drops trailing fields that hold their default.
  • gc=False tells 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:

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

Python
# 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’s Request object 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 flag

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

Python
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_FAILED item 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_DEAD to every API server. Each API server marks itself errored; every pending and future request fails with EngineDeadError, and /health starts 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 uses spawn;
  • when used as a library, it uses fork, unless it detects that CUDA is already initialised, in which case it switches to spawn and 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.

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

Go
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: A API servers + DP engine cores + one worker per GPU (+ a coordinator when DP > 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#

  1. On a machine running vLLM, run ps -eo pid,args | grep "VLLM::" and count the processes. (The titles need the setproctitle package, which vLLM installs.) Does the count match the formula?
  2. In the first program, add a case for 16 GPUs with -tp 8 -dp 2. How many vCPUs does that node need at minimum?
  3. 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#

  1. Why does the engine core send the first message to the API server, and not the other way round?
  2. What three options on EngineCoreRequest’s class declaration make it cheaper to send and hold?
  3. With tensor parallelism of 4, how many workers execute execute_model and how many reply?

Sources#

Checked on 5 October 2026 against main at commit 0c16eee.

↑↓ navigate↵ openesc close

drag to pan · scroll to zoom