LLM Inference Technology — Study Notes
Reference for engineers studying inference-engine systems
0
Key terms
Glossary
⌃

Quick reference for the acronyms you'll hear in every inference conversation. Read this once, then come back when something below feels unfamiliar.

Latency & throughput
TTFT
Time To First Token. Wall-clock from request arrival to the first output token reaching the user. Dominated by prefill + queueing.
ITL
Inter-Token Latency. Time between consecutive output tokens during streaming. Also called TPOT (Time Per Output Token).
TPOT
Same as ITL. Some vendors prefer this name.
TTLT
Time To Last Token. Total time for the full response. TTFT + (num_output_tokens - 1) × ITL.
Goodput
Tokens/sec actually delivered within SLO. A server doing 5000 tok/s but missing TTFT SLO on 50% of requests has 2500 goodput.
QPS
Queries per second. Less useful than goodput for LLMs since requests vary 100× in size.
Hardware & memory
HBM
High Bandwidth Memory. The GPU's on-package DRAM. H100 has 80 GB at 3.35 TB/s; B200 has 192 GB at 8 TB/s. Decode performance is bounded by HBM bandwidth.
NVLink
NVIDIA's GPU-to-GPU interconnect inside a server. H100: 900 GB/s per GPU across 8 GPUs. Fast enough for tensor parallelism.
NVSwitch
The NVLink fabric switch — gives every GPU in a node full-bandwidth all-to-all.
NCCL
NVIDIA Collective Communications Library. The standard for all-reduce, all-gather, broadcast in distributed training and inference.
RDMA
Remote Direct Memory Access. Kernel-bypass network reads/writes. Latency ~1 µs, throughput tens of GB/s.
GPUDirect
NIC writes directly to GPU HBM without going through host RAM. Essential for fast KV transfer in disaggregated serving.
InfiniBand
RDMA-capable network fabric. Common speeds: 200G NDR, 400G XDR. The substrate for cross-node inference fleets.
RoCE
RDMA over Converged Ethernet. RDMA on commodity Ethernet, cheaper than IB. Used at hyperscalers.
MFU
Model FLOPs Utilization. Achieved FLOPs / peak FLOPs. Prefill MFU on H100 can hit 50%+; decode MFU is much lower because it's memory-bound.
MBU
Memory Bandwidth Utilization. Achieved HBM bandwidth / peak. The relevant efficiency metric for decode.
Parallelism
TP
Tensor Parallelism. Split each matmul across N GPUs (Megatron-style column & row split). Requires all-reduce per layer → needs NVLink.
PP
Pipeline Parallelism. Split layers across GPUs. Micro-batched. Tolerates cross-node links but adds pipeline-bubble latency.
EP
Expert Parallelism. Shard MoE experts across GPUs. Requires all-to-all per MoE layer to route tokens to experts.
SP / CP
Sequence / Context Parallelism. Split the sequence dimension across GPUs for very long contexts. Ring-attention is the classic implementation.
DP
Data Parallelism. Replicate the whole model; each replica handles different requests. The most common form of scale-out for inference once one replica fits on a node.
Numerics
FP16
Half-precision float, 16 bits. Standard inference dtype before FP8 era.
BF16
Brain-float-16. Same range as FP32, less mantissa precision than FP16. Now the default training dtype; many inference engines use it.
FP8
8-bit float. Two variants: E4M3 (4 exponent / 3 mantissa, better precision, used for weights/activations) and E5M2 (5/2, wider range, used for gradients). Native tensor-core support on H100 Hopper and later. Roughly 2× throughput vs FP16, ~1−2% quality cost.
FP4
4-bit float. NVFP4 (NVIDIA Blackwell) and MXFP4 (Microscaling open standard). Native on B200. Another ~2× over FP8 at some quality cost.
INT8
8-bit integer. Older but still common (SmoothQuant). Per-tensor or per-channel scales.
INT4
4-bit integer. Typically weight-only (AWQ, GPTQ). Halves HBM footprint and doubles decode throughput on memory-bound layers.
PTQ / QAT
Post-Training Quantization (calibrate after training, fast) vs. Quantization-Aware Training (train with fake-quant ops, higher quality, slow).
Inference-specific
OSS
Open Source Software. Used throughout this doc to mean engines whose source code is publicly available and can be self-hosted (vLLM, SGLang, TGI, TensorRT-LLM, llama.cpp, etc.), as opposed to closed commercial APIs (OpenAI, Anthropic, Google). Many inference providers build on OSS engines and add proprietary kernels, schedulers, and routing on top.
KV cache
The cached Key/Value tensors of attention for every token already in the context. Grows with sequence length; dominates HBM use on long requests.
Prefill
The first forward pass over all input tokens of a request. Compute-bound, parallel over the sequence.
Decode
Subsequent forward passes producing one token at a time. Memory-bound, naturally batchable across requests.
MoE
Mixture of Experts. Sparse model where each token activates only K of N expert FFNs (e.g., DeepSeek-V3: 256 experts, 8 active). Lower FLOPs per token, larger total parameters.
LoRA
Low-Rank Adaptation. Tiny additive matrices (rank 8-64) applied per layer to specialize a base model. Multi-LoRA serving lets one base model serve many fine-tunes simultaneously.
Draft model
A small fast model that proposes the next K tokens in speculative decoding; the large model verifies them in one forward pass.
Acceptance rate
Fraction of speculative tokens the verifier accepts. The key efficiency metric of spec decoding. Modern EAGLE-2 hits 60-80% on typical workloads.
Block / page
A fixed-size chunk (e.g., 16 tokens) of KV cache. The unit of allocation in PagedAttention-style engines.
Prefix cache
Cached KV state for a shared prompt prefix (system prompt, few-shot examples) that can be reused across requests.
→
Primer
Chapter 0.5 · before the techniques
⌃

Before the techniques, two things help the rest of the document land. First, a worked example — one specific request walked end-to-end so the lifecycle and where time goes feel concrete. Second, the engine-stack choice — how common inference-engine stacks differ and why that matters.

0.5.1 · Worked example

Request lifecycle — following one request end-to-end

Every chapter in this document is a technique. This section is one request — walked from network arrival to final token — so you can see those techniques composed in their natural order. Each stage links back to the technique chapter that owns it; read this primer first if you want a system view, or skim it last as a recap.

The scenario

A user in San Francisco submits a 1,536-token RAG prompt to Llama-3.1-8B-Instruct (FP8) running on a single H100 80GB. The engine is vLLM-style with continuous batching, PagedAttention, chunked prefill, and n-gram speculative decoding enabled. The replica is mid-batch with 32 other in-flight users. The user asks for a 256-token response.

Approximate budgets for this load (typical published vLLM/SGLang benchmarks): TTFT ≈ 150 ms, total wall-clock ≈ 1,700 ms, average ITL ≈ 6 ms/token with speculation, ~25% spec acceptance on RAG prompts.

The flame chart — where each component spends its time
REQUEST LIFECYCLE FLAME CHART · LOG-SCALED TIME Network Gateway / Router Engine CPU GPU compute KV cache TTFT 150ms 1 2a 2b 3 4 5 6a 6b 6c Chunked prefill: 3 chunks of ~512 tokens KV grows: 0 → 96MB (1536 tokens) 7 8 · decode loop (256 tokens, ~205 iterations with spec) streaming tokens to client (SSE) KV grows: 96MB → 112MB (+16MB over 256 tokens) 9 · batch dynamics throughout decode → 10 10 0 10ms 20ms 60ms 100ms 150ms 150ms 500ms 1000ms 1700ms ← SETUP + PREFILL zone (0−150 ms, expanded scale) → ← STEADY-STATE DECODE zone (150−1700 ms, compressed scale) → Network / GPU activity Gateway / Router Engine CPU KV cache state Key markers: • t = 150 ms — TTFT achieved (first token reaches client) • t = 1700 ms — EOS, full response delivered (256 tokens streamed)
Time axis is two-zoned: 0–150ms expanded to show pre-decode work; 150–1700ms compressed since decode is repetitive. Numbered cells map to the stage cards below.
Latency budget breakdown — where the 1,700 ms go
LATENCY BUDGET · PROPORTIONAL STACKED BAR prefill 130ms 7.6% decode loop 1,545 ms 90.9% of total wall-clock 0 ms 1,700 ms ↓ TTFT 150 ms Pre-decode breakdown (the 8.8% before TTFT): • Network & TLS: 5 ms   ·   Gateway (auth + Redis lookups): 3 ms   ·   Gateway (tokenize 1,536 tok): 3 ms   ·   Router (radix-tree walk): 2 ms • Engine admission: 1 ms   ·   Prefix-cache lookup + scheduler: 1 ms   ·   Prefill (chunked, 1,536 tok): 130 ms   ·   First-token sample + stream: 5 ms The single biggest lever for cost-per-token is decode speed: spec decoding here saves ~700 ms (45%) vs. no-spec.
Decode dominates (~91% of wall-clock). Network and engine setup are rounding error. This is the typical shape of any non-tiny LLM request.
GPU resource utilization — what's actually binding
GPU RESOURCE UTILIZATION DURING THE REQUEST Tensor cores (FP8 FLOPs) 100% 50% 0% ~75% (prefill: compute-bound) ~15% (decode: idle) HBM bandwidth (3.35 TB/s peak) 100% 50% 0% ~40% ~85% (decode: memory-bound) PREFILL: compute is the bottleneck (HBM idle) DECODE: HBM is the bottleneck (tensor cores idle) 0 20 150 1,700 ms
The two phases stress completely different resources — this is the reason P/D disaggregation (item 1.2) exists, and the reason decode is dominated by HBM-bandwidth optimizations (quantization, KV-cache layout, FlashAttention) rather than peak FLOPs.

This chart is worth memorizing. Almost every single-node optimization in Chapters 2 is justified by which of these two utilization lines it pushes up:

  • Prefill optimizations (FlashAttention kernel work, FP8 tensor cores, chunked prefill) target the compute-utilization curve — turning that 75% into 90%.
  • Decode optimizations (PagedAttention, KV quantization, larger batches) target the HBM-utilization curve — making each byte of bandwidth produce more output tokens.
  • Speculative decoding is unique: it produces extra tokens with spare compute during the otherwise-idle 85% tensor-core slack of the decode phase. That's why it works.
The 10 stages, in detail
Stage 1 · t = 0 — 5 ms · Network

Edge arrival (TLS + GeoDNS)

User (SF) browser / SDK HTTPS / SSE ~3 ms RTT (same metro) Edge POP (anycast IP) backbone Regional GW us-west-2

The request leaves the user's browser as a regular HTTPS POST with SSE response-type headers. Anycast routing via BGP gets the connection to the nearest edge POP in tens of milliseconds best case — for an SF user, probably 2−3 ms to the closest CDN POP. The edge POP terminates TLS (saving the regional gateway from doing it), inspects the request URL to determine the target region, and forwards over the provider's backbone to the regional gateway. For our scenario this is all same-coast, so cumulatively ~5 ms.

What can go wrong here is exactly what you'd expect: an EU user hitting a US-only deployment pays ~150 ms of speed-of-light overhead before any compute happens. The answer is multi-region routing (item 1.6) — route each user to the engine closest to them. For an SF user hitting a us-west-2 region, this stage is essentially free.

Resource cost: client-side network stack; CDN compute. No engine state changes. ↳ See item 1.6 Multi-region routing.

Stage 2 · t = 5 — 11 ms · Gateway CPU

Auth, quota, tenant lookup — and tokenization

Regional GW (Envoy / Rust) JWT verify (in-mem) Quota check (Redis) Tenant config (Redis) GW tokenize model-specific BPE → router with token IDs ~3 ms · lookups in parallel ~3 ms · 1,536 tok @ ~500K tok/s PHASE A · t = 5 → 8 ms PHASE B · t = 8 → 11 ms

The gateway does two phases of work back-to-back. Phase A (t = 5–8 ms) runs three lookups in parallel: JWT signature verification (usually in-memory since the public key is cached), per-tenant rate-limit check against Redis, and tenant-config fetch (model whitelist, priority tier, BYOC routing rules). On a warm path with all three caches hot, ~3 ms; on a cold-cache miss it can spike to 10 ms+.

Phase B (t = 8–11 ms) tokenizes the prompt. This is the non-obvious step: tokenization happens at the gateway, not the engine, because the next stage (cache-aware routing) needs token IDs to walk the radix tree. The gateway hosts a model-specific tokenizer per supported model (Llama, Qwen, DeepSeek, etc.) and produces the full token sequence in ~2−3 ms using a Rust/C++ tokenizer (HuggingFace tokenizers crate or equivalent). Token IDs are then forwarded to the router; the engine will receive pre-tokenized input and skip re-tokenization entirely.

Phase A is also where multi-tenancy enforcement happens (item 1.7). A free-tier user gets routed to the shared serverless pool; an enterprise customer with a dedicated SLA gets routed to their reserved engine pool; a BYOC customer gets routed to engines running in their VPC. The gateway is the single source of truth for which "tier" each request lives in.

One subtle point: the gateway is not on the GPU's critical path. It runs on commodity CPU boxes and can scale horizontally. If tokenization or lookups become a bottleneck, add more replicas. The GPU never waits for it.

Resource cost: CPU + Redis lookups (~3 ms) + tokenize (~3 ms). No GPU activity. ↳ See item 1.7 Multi-tenancy & isolation.

Stage 3 · t = 11 — 13 ms · Router

Cache-aware router picks the replica

PREFIX MATCH · ROUTER PICKS REPLICA WITH LONGEST CACHED PREFIX Incoming request tokens: [A B C D E F G ...] Router radix tree lookup R1 cached: [A B] — matches 2 tokens load: 70% R2 ✓ cached: [A B C D E] — matches 5 tokens load: 60% — chosen R3 cached: [X Y Z] — no match load: 30% Decision: longest prefix match, subject to load cap. R2 wins.

The router receives the request already tokenized from the gateway (Stage 2, Phase B). It holds a radix tree of token-prefix → replica mappings, and walks the tree with the incoming request's token blocks (16 tokens per edge by default). Finding the longest matching prefix tells it which replicas have that prefix cached; load-scoring picks among them. This entire decision is well under a millisecond in well-implemented routers (SGLang router, NVIDIA Dynamo's KV router). The 2 ms in the budget mostly covers the network hop from gateway to router process; the tree walk itself is hundreds of microseconds.

Why this matters: this request is a RAG prompt, which means its first ~1,000 tokens are likely a shared template ("You are a helpful assistant...", followed by the same set of retrieved chunks if this is a hot query). Routing to a replica that already has those tokens cached saves the full prefill cost for that prefix — potentially tens of milliseconds. Across the fleet, this is the single biggest cost-per-token lever you have.

The router also handles graceful fallback: if all replicas with the matching prefix are overloaded, it falls back to the least-loaded replica that doesn't have the prefix, accepting the cache miss to preserve admission latency.

Resource cost: ~2 ms router CPU. ↳ See item 1.1 Cache-aware prefix routing.

Stage 4 · t = 13 — 14 ms · Engine CPU

Engine admission (pre-tokenized request)

Token IDs from GW [128000, 5234, ... ×1536] Parse sampling params temp, top_p, max_tokens... Admission queue request_id assigned ~1 ms total · no re-tokenize (token IDs already provided by gateway)

The chosen replica's HTTP/gRPC frontend accepts the request — with token IDs already pre-computed by the gateway in Stage 2 Phase B. Because tokenization happened upstream, the engine skips it entirely and goes straight to parsing sampling parameters (temperature, top_p, top_k, max_tokens). The request gets its request_id assigned and is placed in the admission queue.

The admission queue is decoupled from the active batch: requests can pile up here for a few ms during heavy traffic without blocking the scheduler loop. If the queue exceeds a backpressure threshold, the engine returns 503 rather than letting the queue grow unbounded — this is also the autoscaling signal (queue depth, not CPU).

Why no tokenization at the engine? Cache-aware routing (Stage 3) needs token IDs to walk its radix tree, so tokenization has already happened upstream at the gateway. Re-tokenizing here would be wasted work and could even produce slightly different IDs if the tokenizer versions drifted — which would silently break prefix-cache lookups. The engine treats the gateway-provided token IDs as authoritative.

Resource cost: ~1 ms CPU (admission queue + sampling-params parse). No tokenizer call. ↳ See item 2.1 Continuous batching for the scheduler loop that pulls from this queue.

Stage 5 · t = 14 — 15 ms · Engine CPU

Prefix-cache lookup & scheduler admits the request

RADIX TREE WALK ON THIS REPLICA'S KV BLOCK INDEX root [128K..] b0 [5234..] b1 [9182..] b62 ━━━ [7821..] ? ← first miss 63 blocks (1,008 tokens) cached — first 65.6% of prompt Will only prefill remaining 528 tokens — saved ~62 ms of compute Scheduler: admit request, allocate 33 KV blocks for new tokens, refcount-bump on shared blocks, add to batch

The replica maintains a radix tree indexing its currently-cached blocks. The engine walks the tree with the tokenized prompt, matching it block by block (16 tokens per block). Because the router (Stage 3) already picked this replica for its match, we expect a deep hit — in this case 63 blocks (1,008 tokens) of the 1,536-token prompt are already cached. The first 65.6% of the prompt is free.

The block manager refcount-bumps each shared block: any block reused from cache must not be evicted while in flight. Then the scheduler checks two admission gates: (1) does the request fit under max_num_seqs (the concurrency cap), (2) is there room for the remaining ~33 KV blocks the new tokens will need. Both pass; the request is admitted into the active batch.

Critical to notice: the actual prefill on the GPU only runs the tail of the prompt — tokens 1008−1535 (528 tokens). The first 1,008 tokens contribute their cached KV directly via PagedAttention's block-table indirection. This is what makes prefix caching so impactful at scale.

Resource cost: ~1 ms tree walk + ~1 ms block allocation. Block manager state changes; no GPU activity yet. ↳ See items 1.3 Hierarchical KV cache, 2.2 PagedAttention.

Stage 6 · t = 20 — 145 ms · GPU

Prefill (chunked, interleaved with in-flight decodes)

CHUNKED PREFILL: 528 NEW TOKENS × 3 CHUNKS INTERLEAVED WITH 32 DECODES Iter A (43 ms): prefill chunk 1 (512 tok) +32 dec Iter B (43 ms): prefill chunk 2 (16 tok) +32 dec (small chunk: rest of cache miss) Iter C (39 ms): (no new prefill needed) +33 dec (our request joins decode pool) KV cache state: Start: 1,008 prefix tokens shared (cached blocks reused via block-table indirection) End: full 1,536 prompt tokens KV-resident on this replica · 33 new blocks allocated PagedAttention kernel gathers KV via per-request block table — non-contiguous K/V reads handled by kernel-level gather

The GPU now does the actual work. Because of chunked prefill, the engine doesn't run the 528-token prefill as one monolithic forward pass; instead it splits it into chunks (default ~512 tokens) and treats each chunk as one scheduler iteration that interleaves with the 32 ongoing decodes. This keeps decode tokens flowing to other users instead of pausing them for hundreds of ms.

Each iteration is a single forward pass that handles both kinds of work:

  • For the new request: 512 (or 16) query positions running through all 32 layers, attending to the request's full cached + just-computed KV. This is the compute-bound piece (~75% tensor-core utilization, as seen in the chart above).
  • For the 32 existing requests: one query position each (their newest token), attending to their own KV blocks. This is the memory-bound piece (HBM bandwidth saturated reading their KVs).

Both pieces use the same kernel — FlashAttention-3 with the PagedAttention block-table extension. The kernel takes a per-request block table and does a gather of the K/V tensors from non-contiguous HBM blocks before computing attention. Quantization (FP8 weights, FP8 KV) cuts the bytes read by 50% vs FP16, doubling the effective HBM bandwidth.

By iteration C, the entire 1,536-token prompt has been processed; our request can now join the decode pool alongside the 32 others. TTFT is t = 145 ms.

Resource cost: ~125 ms GPU compute. 33 new KV blocks allocated (~33 MB at FP8 KV: 528 tokens × ~64 KB/token). ↳ See items 2.6 Chunked prefill, 2.2 PagedAttention, 2.3 FlashAttention, 2.5 Quantization.

Stage 7 · t = 145 — 150 ms · GPU + Network

First token sampled & streamed → TTFT achieved

logits [vocab_size] Sample top-p / temp Detokenize "The " SSE: data: {"t": "The "} → user's browser ~5 ms end-to-end · TTFT marker fires at t=150ms

At the end of iteration C the model produces a logit vector over the vocabulary (~128K entries for Llama-3.1) for our request. The sampler applies temperature, top-p, top-k, and any user-provided sampling parameters to pick the next token (token ID 5234, say). The detokenizer converts it back to a string fragment ("The "). The engine emits an SSE event with the token fragment. The event travels back through the gateway and edge to the user's browser, where the streaming UI renders the first visible character.

This is the TTFT moment — the most-watched metric in interactive LLM serving. Hitting 150 ms means a user perceives the response as "instant." Sub-100 ms is what the best vendors target; 200−500 ms is "feels normal"; >1 s is "feels broken."

One subtle but important detail: detokenization can be tricky. BPE tokens often span partial UTF-8 sequences (e.g., emoji), so the engine has to buffer a token or two before emitting bytes the client can actually render. Production engines have careful "incremental detokenization" code paths to avoid mid-character SSE chunks.

Resource cost: <1 ms GPU (sampler kernel) + ~4 ms network back to user. ↳ Sampling-related techniques aren't covered as separate items; see 2.1 Continuous batching for context.

Stage 8 · t = 150 — 1,695 ms · GPU (1,545 ms = 91% of wall clock)

Steady-state decode loop (256 tokens, with speculative decoding)

ONE DECODE ITERATION (REPEATED ~205 TIMES) 1. n-gram lookup match last 3 generated tokens against prompt & produced text → K=4 candidate tokens 2. Verify (1 forward pass) target model evaluates 5 positions (1 + 4 spec) ~6 ms (batch of 33) 3. Commit accept matching prefix ~2 tokens / iter avg SSE: 2 tokens streamed → user's browser Over 256 output tokens: • ~205 iterations (each ~7.5 ms, including verifier overhead) • ~25% spec acceptance on RAG content (template tokens, retrieved-chunk reuse, common phrasing) • 1.25 tokens / iter average (1 baseline + 4 spec × 25% acceptance) • ~700 ms saved vs no-spec (~2,240 ms total without speculation)

This stage is 91% of total wall-clock and almost everything that matters for cost-per-token. The engine is now in the decode loop — one forward pass per iteration, producing one token per request in the batch. With our request and 32 others, the batch size is 33.

Speculative decoding sits inside this loop. For RAG workloads, n-gram lookup speculation is the winner: tokens being generated frequently appear verbatim somewhere in the prompt (e.g., the model is paraphrasing a retrieved chunk). The engine maintains a hash index over the prompt's n-grams; when the last 3 generated tokens match, it proposes the continuation as candidate tokens.

Each iteration:

  1. Lookup: cheap — ~0.1 ms on CPU, run by the scheduler.
  2. Verify: the target model runs a single forward pass with K+1 positions in its sequence (instead of 1). For our batch of 33, that's a slightly larger matmul — ~7.5 ms instead of ~6 ms for plain decode.
  3. Commit: the verifier emits its true next-token distribution for each position. The first speculated token that doesn't match the verifier's choice ends the accepted prefix. On average 1.25 tokens commit per iteration.

Continuous batching is what keeps this efficient: any request that hits its max_tokens or emits <eos> leaves the batch instantly; any new request whose prefill just finished joins immediately. The batch composition shifts every iteration without GPU idleness.

HBM bandwidth is the binding resource here — tensor cores are at ~15%, HBM at ~85% (see chart above). This is why quantization and KV cache layout matter so much: every byte you don't have to read is a faster decode.

Resource cost: ~1,545 ms GPU. ~16 new KV blocks (256 tokens / 16 = 16 blocks × ~1 MB/block) = ~16 MB KV growth. ↳ See items 2.1 Continuous batching, 2.2 PagedAttention, 2.3 FlashAttention, 2.4 Speculative decoding, 2.5 Quantization.

Stage 9 · t = 150 — 1,695 ms · Scheduler (concurrent with Stage 8)

Batch dynamics — other requests join & leave during our decode

BATCH MEMBERSHIP TIMELINE (OUR REQUEST = ⓡ) t=150ms ~t=900ms t=1695ms ⓡ our request: 256 tokens streaming R_a (50 tok) ✓ R_new (prefill + 100 tok decode) ✓ R_b (200 tok) ✓ R_c (very long: 1000+ tok, still running) R_late (joined at t=400, finished t=1500) R_very_late (joined at t=1100) Batch slot: slot 1 (us) slot 2 slot 3 slot 4 slot 5 slot 6

While our request is decoding its 256 tokens, the scheduler is continuously reshaping the batch. Some neighbors finish faster than us (short requests, or ones that hit <eos>). Slots they vacate get reused by new requests whose prefill just completed. Very-long-running requests stay alongside us. Every decode iteration is a new composition.

This is what makes continuous batching dramatically better than static batching: the GPU is never idle waiting for the slowest request to finish — finished requests depart instantly. From our request's perspective, the batch size we share the GPU with fluctuates between maybe 25 and 40 depending on neighbor traffic, which slightly changes our per-token decode time.

Two consequences worth understanding:

  • Per-token latency variance is normal. If you have a strict ITL SLO and one request next to a 4000-token monster, your decode steps run slightly slower during that overlap. Engineering for "p99 ITL" is mostly about controlling batch contention.
  • Throughput scales sublinearly with concurrency. Going from batch 1 to batch 16 multiplies aggregate throughput ~10× (HBM amortization). Going from 16 to 64 multiplies it maybe 2× further. Beyond ~128 the GPU is just memory-limited.

Resource cost: scheduler CPU work, sub-ms per iteration. ↳ See item 2.1 Continuous batching.

Stage 10 · t = 1,695 — 1,700 ms · All

EOS detected — cleanup & response finalized

EOS sampled (or max_tokens hit) Scheduler drop from batch Block mgr refcount-dec, free SSE: [DONE] stream closed ~5 ms total · request finalizes, slot freed for next admission

The model samples its end-of-sequence token (or hits the user's max_tokens limit). The scheduler removes the request from the active batch on the next iteration. The block manager walks the request's block table and decrements the refcount on each block; blocks with refcount=0 return to the free pool. Blocks shared with other requests (e.g., the cached prefix prefix of the original prompt) stay live with their higher refcount.

The engine emits the final SSE event (data: [DONE]) and closes the streaming connection. The gateway logs the request metadata (TTFT, total tokens, ITL, tenant ID) for billing, autoscaling signal, and feedback into Layer-2 re-tuning (item 3.1's "periodic re-tuning" loop).

On the user's side, the streaming UI receives the close event and finalizes the rendering. The conversation context is now ready for the next user turn — and if this is a multi-turn chat, the assistant's response just got added to the system prompt for the next round, where it'll be a hot prefix cache entry the next time the user sends a follow-up.

Resource cost: ~5 ms cleanup + signal emission. KV blocks return to free pool. ↳ Cleanup pattern not covered as a separate item; the autoscaling signal feeds 3.1 Workload-feedback optimization.

Reading this as a primer

The point of walking through one request is to see how the techniques in Chapters 1−3 compose. A few patterns are worth carrying forward:

  • Decode is everything. 91% of wall-clock; nearly all the optimization literature in Section 2 targets this phase. If you remember one thing, remember that inference perf is mostly about making decode faster.
  • Prefill and decode bottleneck on different resources. Prefill is compute-bound (tensor cores), decode is HBM-bound (memory bandwidth). The whole P/D disaggregation story (item 1.2) and most of the per-phase tuning falls out of this single fact.
  • Cache reuse pays everywhere. Stage 3 (router) and Stage 5 (radix-tree lookup) cooperated to skip ~65% of the prefill compute on this request. At scale, fleet-wide cache hit rate is the dominant cost lever.
  • Many decisions happen at many cadences simultaneously. While the decode loop runs at 6 ms per iteration, the scheduler is reshaping the batch every iteration, the router is updating its prefix-cache table every few seconds, and the workload-feedback re-tuning runs hourly. Item 3.1 ties this together.

When an interviewer asks "walk me through what happens when a request hits your inference service," this is what they're testing for: the ability to traverse the stack, point at the right technique at each stage, and articulate where time and resources actually go.

0.5.2 · Context

Engine stacks

The OSS inference world has split into two camps based on how the model graph is executed.

Two camps: PyTorch-native vs. AOT-compiled

Style Engines How the forward pass runs Trade-off
PyTorch-native vLLM, SGLang, TGI, Together, DeepInfra Python scheduler dispatches PyTorch ops + custom CUDA/Triton kernels at runtime. CUDA graphs capture the decode step to remove launch overhead. Fast iteration on new model architectures, weeks not months to add a model. Slightly higher per-op overhead than fully compiled.
AOT-compiled TensorRT-LLM, llama.cpp (different reason), Groq, Cerebras Model graph is compiled ahead of time into a fused engine binary; runtime just streams inputs/outputs. Lowest per-op overhead, best peak throughput on supported models. Slow to support new architectures — weeks of engineering per model.

For PyTorch-native systems, this choice is strategic. They can serve dozens of open-weight models (Llama, DeepSeek-V3, Qwen, Mixtral, Whisper, Stable Diffusion, etc.) and need to onboard new models the day they're released. A new HuggingFace Transformers modeling_*.py file can be wired into a PyTorch-native engine in hours.

Two PyTorch-ecosystem perf tools worth knowing

CUDA Graphs

Each kernel launch in PyTorch costs ~5−15 µs of host-to-device overhead. A decode step calls hundreds of kernels (attention, layer norm, matmul, etc.) per layer per token. On small batches this launch overhead, not GPU compute, becomes the bottleneck.

CUDA Graphs record an entire decode step's kernel sequence once, then replay it as a single GPU command with no host involvement. Typical win: 10−30% faster decode at small batch sizes. Used by vLLM and SGLang. The trade-off: the captured graph is shape-specialized, so engines maintain a graph for each (batch_size, seq_len) bucket.

Triton

A Python-embedded GPU kernel language (from OpenAI, 2021). Lets you write custom CUDA-like kernels with much less ceremony than CUDA C++. FlashAttention has both CUDA and Triton implementations; PagedAttention's gather kernel is often Triton. vLLM ships a large Triton kernel library; SGLang adds more.

For an engineer interviewing at a PyTorch-native shop, "I read/write Triton" is a stronger signal than "I know CUDA," because that's the actual kernel-development surface for most modern inference engines.

With that context in hand, the rest of the document covers techniques engine-agnostically — without vendor-specific callouts.

1
Distributed infra
Section 1 · multi-node techniques
⌃

These techniques live above the single-GPU engine. They decide where requests go, how KV state is shared across machines, and how the cluster makes the most of every GPU it owns. For an inference startup, this is where the largest cost-per-token wins still live.

1.1 · Tier 0

Cache-aware prefix routing

Importance: HIGH
Origin
SGLang Router (LMSys, 2024) · Mooncake Conductor (Moonshot, 2024) · NVIDIA Dynamo KV-router (2024)
OSS engines
SGLang Router vLLM Production Stack NVIDIA Dynamo AIBrix (ByteDance) llm-d (Red Hat)
What it does & why it matters

A standard load balancer (round-robin, least-connections) sends requests to a random GPU. That's terrible for LLMs because each GPU caches the KV state of every prompt it recently processed (a "prefix cache" or "radix tree" of KV blocks). If you route a request with a system prompt to a GPU that has never seen that prompt before, you pay the full prefill cost — even though another GPU in the fleet has the answer in its HBM already.

Cache-aware routing picks the GPU that already holds the longest matching prefix of the incoming request's tokens. On workloads with shared system prompts (which is essentially every chat product, RAG app, and agent system), this typically converts a 10−30% single-node cache-hit rate into a 60−90% fleet-wide hit rate.

The pain point it solves: prefill is the dominant cost on long-prompt workloads. A 4K-token prompt + 200-token response spends ~95% of its compute in prefill. Cutting prefill via cache reuse cuts the cost-per-token of the entire request by a similar factor.

Key implementation idea

The router maintains a distributed index of which GPU has which prefix cached. For each incoming request, it scores every replica by "longest matching prefix length" and picks the highest scorer, subject to load-balancing constraints (don't overload one hot replica).

Two design patterns are common:

  • Hash-based (consistent hashing on prefix chunks). Tokenize the prompt, chunk into 16-token blocks, hash each prefix. The router maps each hash to a replica via a consistent hash ring. Same prefix → same replica, deterministically. Implementation is simple, no central index. Used by Mooncake.
  • Index-based (centralized prefix tree / RadixAttention). Each replica periodically reports its cached prefix tree to a central router; the router does longest-prefix-match on the union tree. More accurate (knows actual cache state, not just hashes), but the router becomes a fan-in point. Used by SGLang Router.

Both must handle cache eviction events — when a replica evicts a prefix under memory pressure, routing decisions get stale. The fix is either short-TTL heartbeats (index style) or letting consistent hashing's "second-choice replica" handle misses gracefully.

A subtle detail: routing is also load-aware. If the best-matching replica is overloaded (queue depth high, KV cache near full), the router falls back to the second-best replica and accepts the cache miss. Pure cache affinity without load balancing creates hot spots.

Hash-based vs prefix-tree-based: pros & cons

The two approaches have opposite trade-offs — one wins on operational simplicity, the other wins on the metric that drives cost (cache hit rate). Knowing which dimension you're optimizing for is half the answer.

Dimension Hash-based (consistent hashing) Prefix-tree-based (radix tree)
Routing decision latency < 1 µs — hash + ring lookup 10–100 µs — tree walk on global state
Router throughput Millions of QPS (stateless) ~10K QPS per instance; horizontal scaling adds sync cost
Router state None — just the hash ring Global radix tree, can be GBs at fleet scale
Operational complexity Low — no sync, no consistency, no leader election High — replicas must report cache changes; sync bugs are a real category
Determinism Same prompt → same replica, always State-dependent — same prompt might route differently as cache evolves
Single point of failure None — any router instance can decide independently Yes — the router holds canonical state (typically replicated for HA)
Cache hit rate (the cost lever) 50−70% fleet hit rate typical 70−90% fleet hit rate typical
Hot-prefix handling One home replica → hot prefix becomes hot replica Hot prefix naturally replicated; router picks least-loaded copy
Cache-eviction awareness None — ring assumes the deterministic mapping is the cache state Real-time — replicas notify on every eviction
Cold start (new replica) Ring shifts; ~1/N of prefixes migrate → uniform cache miss until warm New replica reports empty cache → router never sends it prefix-match traffic → needs explicit warm-up policy
Load balancing within a prefix group Hard — needs an extension like consistent hashing with bounded loads (CHWBL) Natural — pick the least-loaded replica that has the prefix
Debugging Easy (deterministic, reproducible) Harder (state-dependent decisions)

So which has better performance? Depends entirely on which metric:

  • Router throughput & decision latency: hash wins ~100×.
  • Fleet-wide cache hit rate: prefix tree wins by 15−25 percentage points on diverse workloads — which translates to roughly 15−25% lower GPU cost per token.
  • Operational simplicity: hash wins by a large margin.

Because cache hit rate is the dominant cost lever for a serious inference vendor, most production systems pay the operational tax for the tree. The exception is small deployments (< 20 GPUs) where simplicity wins.

The production hybrid that combines both

The sharpest production systems use both, in layers:

  • Layer 1 — Coarse placement (hash-based): consistent-hashing on prefix shards determines which subset of replicas could ever cache this prefix. Bounds the global state and creates natural failure isolation.
  • Layer 2 — Fine matching (prefix tree): among that small subset, walk a per-shard radix tree to pick the best-matching, least-loaded replica.

Mooncake (FAST '25) describes essentially this layout for their KVCache pool. NVIDIA Dynamo's KV router does something similar. The key insight: the radix tree's complexity is bounded by the per-shard subset, not the whole fleet, so it scales horizontally.

How the prefix tree actually works — a worked example

The foundational property: chained hashes (Merkle-style)

Before the worked example, one property has to be clear because everything else builds on it. A natural first guess is that each block has its own independent hash: hash(B1), hash(B2), hash(B3) — each computed only from its own 16 tokens. That's wrong, and getting it right is the foundation.

The block hashes are chained — each block's hash recursively incorporates its parent's hash:

block_hash(B_n) = hash(block_hash(B_{n-1}), token_ids[16n .. 16n+16]) # starting condition for the root: block_hash(root) = 0 (or any fixed sentinel)

So block_hash(B_n) is a fingerprint of the entire prefix from position 0 through position 16(n+1), not just the most recent 16 tokens. Two requests producing the same hash at depth n have by construction followed the identical token path from root. The hash is the path identity.

This is exactly the Merkle-tree pattern Git uses for commit history, IPFS uses for content addressing, and Bitcoin/Ethereum use for block linkage. In the inference context it gives the tree a powerful invariant:

  • Same hash → same tree node, always. No path ambiguity, no need to disambiguate via parent context at lookup time.
  • Engines can compute hashes independently and agree. No coordination protocol needed for cache-identity consensus across the fleet.
  • Eviction events need only the hash. The router can locate the tree node in O(1) from a flat hash → node map.
  • Why attention requires this: KV state for token N depends on every prior token (attention is contextual). Two prefixes that share trailing tokens but differ earlier have completely different K and V values for the trailing block. Independent per-block hashing would falsely collapse these as "the same cache" and produce wrong outputs. Chaining makes them correctly distinct.

So when our worked example below shows blocks B1, B2, B3 shared across requests, what's really happening is: the chained hashes for those positions collide because the underlying tokens are identical. The tree structure isn't a separate indexing layer on top of independent hashes — the chained hash is the tree's identity scheme.

The setup

Imagine three requests have arrived at the fleet, each tokenized into 16-token blocks. (Block size is the granularity at which KV cache is allocated — same idea as PagedAttention, item 2.2.)

  • Request A: blocks [B1, B2, B3, B4, B5] — served by replica R1
  • Request B: blocks [B1, B2, B3, B6, B7] — served by replica R2
  • Request C: blocks [B1, B2, B8, B9] — served by replica R3

All three start with B1, B2 (think: shared system prompt). A and B continue with B3 (think: shared retrieved chunk). They then diverge.

The tree that results

PREFIX TREE AFTER REQUESTS A, B, C · SHARED NODES = SHARED PREFIX root (empty prefix) B1 refcount: 3 (A,B,C) replicas: R1, R2, R3 B2 refcount: 3 (A,B,C) replicas: R1, R2, R3 B3 refcount: 2 (A,B) replicas: R1, R2 B8 refcount: 1 (C) replicas: R3 B4 refcount: 1 (A) replicas: R1 B6 refcount: 1 (B) replicas: R2 B5 R1 (request A end) B7 R2 (request B end) B9 refcount: 1 (C) R3 (request C end) shared by all 3 shared by 2 unique to one
The tree shape captures prefix sharing structurally. Greener nodes are deeper in the sharing — the “cache hit at this prefix length” bar is exactly the depth at which the green-zone ends for a new request.

What's actually stored at each node

Each tree node holds a small struct:

PrefixTreeNode { block_hash: u64 // CHAINED: hash(parent.block_hash, token_ids[16]) // ↳ uniquely identifies the entire path from root parent: NodePtr children: HashMap<u64, NodePtr> // child-block-hash → node replica_set: BitMap // which replicas have this prefix cached refcount: u32 // # in-flight requests using this prefix last_access: timestamp // for LRU eviction block_size: u8 // typically 16 tokens } // Insertion (engine processing a new request): parent_h = 0 // root sentinel for chunk of 16 tokens in prefill_tokens: h = hash(parent_h, chunk) // ← the chained step node = tree.lookup_or_insert(h) // O(1) via flat hashmap node.replica_set.set_bit(this_replica_id) parent_h = h // ← next iteration uses this as parent

Two flavors of hashing appear here and they do different jobs:

  • Chained content hash (the block_hash field itself): the Merkle-style fingerprint that makes "same hash → same tree node" true. This is what makes prefix sharing structurally automatic — two requests with identical first-N tokens compute identical first-N chained hashes, so they walk into the same tree nodes.
  • Children hashmap (HashMap<u64, NodePtr>): a per-node lookup table for O(1) child traversal. When walking the tree, you compute the next chained hash and check whether the current node has a child with that hash. Found → descend. Not found → prefix walk ends.

So tree walking isn't "compute a separate hash and pray for collision" — it's "compute the chained hash, which by construction matches if and only if the prefix actually matches, and use the children hashmap for fast descent."

Walking the tree for an incoming request

Now suppose Request D arrives with blocks [B1, B2, B3, B6, B8, B9]. The router walks the tree:

  1. From root, hash B1, find it in root's children → descend to B1 node. Note replicas: {R1, R2, R3}.
  2. From B1, hash B2, find it → descend to B2. Replicas still {R1, R2, R3}.
  3. From B2, hash B3, find it → descend to B3. Replicas narrow to {R1, R2}.
  4. From B3, hash B6, find it → descend to B6. Replicas narrow to {R2}.
  5. From B6, hash B8... not found in B6's children. Walk ends.

Result: longest match is 4 blocks (B1, B2, B3, B6) and the only replica holding that full prefix is R2. The router sends Request D to R2, subject to load checking. R2 reuses 64 cached tokens of KV state and only has to prefill the tail (blocks B8, B9 — 32 new tokens).

The operations the tree must support

  • Lookup (per-request, fast path): walk from root, descend on hash match, return deepest node. O(prefix length / block size). For a 1,000-token prompt with 16-token blocks: 63 hash lookups, each O(1) — well under 100 µs.
  • Insert (when a new request gets admitted to a replica): walk from root, create new child nodes for any uncached suffix. Increment refcount on every node touched. Update replica_set to include this new replica.
  • Release (when a request finishes or gets evicted): decrement refcount on every node it touched. If a node reaches refcount 0, it's eligible for eviction but not deleted yet (LRU policy).
  • Evict (when a replica reports a block was evicted under memory pressure): remove that replica from the node's replica_set. If the set becomes empty, the prefix is gone from the fleet — delete the node and all its descendants.

The refcount is the magic that prevents bad eviction: a prefix in active use by multiple requests stays alive even when the LRU policy would evict it. PagedAttention (item 2.2) uses the same refcounting idea within a single engine's block manager; the cluster-level prefix tree is the same pattern, one level up.

The three roles hashing plays — don't conflate them

"Is this just hashing?" is a fair question, but the answer is "three different uses of hashing, each doing a different job":

Where hashing appears What it does Where in this section
Chained content hash (the block_hash field) Merkle-style fingerprint of the entire prefix path. Same hash ⇒ same tree node, by construction. This is what makes prefix sharing automatic. The foundational property above
Children hashmap (per-node HashMap<u64, NodePtr>) O(1) lookup of "does this node have a child with hash X?" during tree walks. Fast descent without scanning siblings. The node struct above
Consistent-hash ring (only in hash-based routing variant, not in prefix-tree routing) Maps a prefix hash directly to a replica. The hash is the routing decision. No tree, no state. The hash-based routing variant earlier in this section

Hash-based routing uses #3 alone. Prefix-tree routing uses #1 and #2 together: chaining gives the tree its identity invariant, the children hashmap gives it fast descent. The combination is what lets prefix-tree routing capture sharing relationships across requests, support partial-prefix matches, and track which specific replicas hold which sub-trees in real time — capabilities hash-based routing structurally cannot offer.

Implementation note — how the replica_set BitMap actually addresses machines

The BitMap doesn't contain IP addresses; it contains bit positions. Translating a bit position to an actual replica needs a separate, small lookup table that the cluster control plane maintains. Concretely:

Cluster membership table (one global table, ~bytes per replica): index IP:port status AZ joined ──────────────────────────────────────────────────────────── 0 10.0.1.5:8080 healthy us-east-1a 2026-04-01 1 10.0.1.6:8080 healthy us-east-1a 2026-04-01 2 10.0.2.5:8080 healthy us-east-1b 2026-04-12 3 10.0.2.6:8080 draining us-east-1b 2026-04-20 4 10.0.3.5:8080 healthy us-east-1c 2026-05-02 ... Prefix-tree node "B3": replica_set: 0 0 0 1 0 1 0 1 ... 0 (1024 bits = 128 bytes) ↑ ↑ ↑ R3 R5 R7 ← these replicas have this prefix cached Routing decision after walking the tree: 1. Final candidate bitmap (after tree walk) = 0 0 0 1 0 1 0 1 ... 2. Iterate set bits → indices [3, 5, 7] 3. Score each by load → R5 wins 4. cluster_table[5] → 10.0.2.5:8080 5. Forward HTTPS request → 10.0.2.5:8080

The membership table is consulted exactly once per request — at the very end, to translate the chosen index into an IP:port. Every operation in the hot path (tree walk, set intersection, popcount, load comparison) operates on integer indices. That's why bit-level operations are so cheap even on millions of tree nodes.

Yes, BitMap size caps the cluster size

The bitmap dimension is fixed at engine startup, so it does cap the number of replicas the router can address. Production choices:

  • 1024 bits (128 bytes per node) — fine for fleets up to ~1000 replicas. A 5M-node tree uses ~640 MB just for replica-set fields. Suitable for typical mid-size inference deployments.
  • 4096 bits (512 bytes per node) — for large fleets. Same 5M-node tree uses ~2.5 GB. Plenty of headroom.
  • Roaring BitMap — auto-adjusts representation: a tiny sorted array for sparse nodes (most leaves), full BitMap for dense ones (high-up shared prefixes). Used by Mooncake, Lucene, ClickHouse. Removes the static cap entirely while keeping the speed of BitMap operations.

The conservative move is to pre-size for 4−10× current fleet so you never hit a resize event (which would require rebuilding the whole tree). Roaring sidesteps this problem if you don't mind the slight implementation complexity.

How does a replica get assigned to a specific index?

The cluster control plane owns the membership table and assigns indices. The mechanism varies by orchestrator, but the contract is the same: stable IDs across reboots. Typical pattern when a new replica joins:

  1. Replica boots, registers with the control plane (Kubernetes controller, etcd-backed registry, or vendor-specific service). Sends its address, AZ, and capabilities.
  2. Control plane assigns the lowest unused index, or returns the replica's previously-assigned index if it's a known node restarting (looked up by pod name, instance ID, or persistent UUID).
  3. The new (index, IP:port) pair is broadcast to all routers via the membership-table update channel.
  4. The replica starts serving; its index begins appearing in replica_set bitmaps as it caches blocks and reports the events upstream.

When a replica leaves (planned drain or crash):

  1. Control plane emits a departure event.
  2. Each router clears that bit from all tree nodes in a single sweep — one AND-NOT mask, cheap because the bit position is constant.
  3. The index returns to the free pool (or is left gapped, depending on policy — gaps are harmless; they're just unused bits).

One thing to be careful about: don't use the IP:port itself as the index (e.g., hash(ip) % 1024). That breaks index stability when nodes change addresses (k8s reschedules a pod to a new IP) and forces re-bitmapping of the entire tree on every IP change. Indices must be stable across address changes.

For Kubernetes deployments, the cleanest answer is StatefulSets with a persistent name + small index registry in etcd or a sidecar service — pod name engine-7 maps to index 7 every time, regardless of which node the pod runs on or what IP it gets.

References
1.2 · Tier 0

Prefill / Decode disaggregation

Importance: HIGH
Origin
Splitwise (Microsoft, ISCA '24) · DistServe (UCSD/SJTU, OSDI '24) · Mooncake (Moonshot, FAST '25)
OSS engines
vLLM (PD mode, 0.6+) SGLang TensorRT-LLM NVIDIA Dynamo Mooncake llm-d
What it does & why it matters

Prefill and decode have opposite hardware appetites:

  • Prefill processes all input tokens in one forward pass. It's compute-bound — the GPU's tensor cores are saturated. A 4K-token prompt on Llama-70B costs ~30 TFLOPs of compute; HBM bandwidth is barely touched.
  • Decode processes one token per step, reading the entire KV cache each time. It's HBM-bandwidth-bound — tensor cores sit idle waiting for memory. Larger batch sizes help (amortize the KV read), but more batches means more KV pressure.

When you run both phases on the same GPU (the default), they interfere. A new request's prefill freezes all in-flight decodes for hundreds of milliseconds — users see TTFT spikes and decode tokens-per-second going to zero. Continuous batching helps interleave them, but the fundamental hardware mismatch remains.

Disaggregation runs the two phases on separate GPU pools: "prefill GPUs" that do nothing but absorb new prompts, and "decode GPUs" that stream tokens for established sessions. After prefill finishes, the KV cache is shipped over NVLink/RDMA to a decode GPU, which takes over the stream.

DistServe reports 4−7× goodput improvement at the same TTFT/TPOT SLO. The win is largest on long-prompt workloads (RAG, code, agents) where prefill dominates.

Key implementation idea

Three components make up a disaggregated cluster:

  1. Prefill pool. GPUs optimized for compute: high tensor-core utilization, large batch sizes (latency-tolerant), often higher TP degree to fit long-prompt activations. A prefill request runs once, produces a KV cache, then the GPU is free for the next prefill.
  2. Decode pool. GPUs optimized for HBM throughput and concurrent sessions: small TP degree (decode all-reduces are tiny per token), large per-replica batch size, often FP8 weights for memory savings. Each decode GPU holds dozens to hundreds of in-flight sessions.
  3. KV transport. Between the two pools. This is the hard part. The KV cache for a 4K-token request on Llama-70B is ~1 GB (FP16) or 500 MB (FP8). It must move from prefill GPU to decode GPU in milliseconds, not seconds. Solutions: NVLink for same-node transfer (900 GB/s — basically free); GPUDirect RDMA over InfiniBand for cross-node (50−200 GB/s). See item 1.4.

The scheduler decides how many prefill vs. decode GPUs to provision. A naive static split is wasteful: prefill-heavy traffic needs more prefill GPUs and vice versa. Modern systems (Splitwise's "mixed pool", Mooncake's elastic ratio) dynamically rebalance: track the prefill/decode queue depths, repurpose GPUs by paging the model in/out when ratios drift.

Chunked prefill (item 2.6) is a related but different idea — it interleaves prefill chunks with decodes on the same GPU. Disaggregation is the architectural answer; chunked prefill is the single-node mitigation. Modern systems often use both: disaggregation for the cluster, chunked prefill on each prefill GPU to keep background decode tasks flowing.

References
1.3 · Tier 0

Hierarchical / pooled KV cache

Importance: MEDIUM-HIGH
Origin
Mooncake Store (Moonshot, 2024) · LMCache (UChicago, 2024) · NIXL (NVIDIA, 2024)
OSS engines
LMCache Mooncake Store NIXL NVIDIA Dynamo vLLM + LMCache integration
What it does & why it matters

HBM is the most expensive RAM in the data center. At a typical H100 rental of $2.5/hr, amortizing across its 80 GB of HBM works out to ~$275 per GB-year ($2.5/hr × 8,760 hr ÷ 80 GB). CPU DRAM is ~100× cheaper per GB-year, NVMe is another ~20× cheaper than DRAM. Yet by default, KV cache lives only in HBM — the instant a request finishes or gets evicted under pressure, its KV blocks are gone.

Hierarchical KV cache extends the cache outward into cheaper memory tiers:

  • L1 — HBM (TBs of bandwidth, ~80 GB capacity). Hot, in-use blocks.
  • L2 — Host DRAM (PCIe 5.0: ~64 GB/s read, 1−2 TB capacity). Recently evicted prefixes.
  • L3 — Local NVMe (~10 GB/s, 10−100 TB capacity). Long-tail prefixes from yesterday's traffic.
  • L4 — Remote DRAM/NVMe pool (RDMA: 10−50 GB/s, petabytes capacity). Fleet-wide shared cache.

On a chat workload with a 4K-token system prompt shared across users, this extends the effective KV cache from minutes (in HBM) to days (on NVMe). The arithmetic that matters: even at NVMe's 10 GB/s, restoring a cached prefix is still 10× faster than recomputing prefill. So unless your hit rate drops below ~10%, paging from disk is a win.

The pain point: fleet-wide cache hit rate is the single biggest cost lever. On a large-scale serverless inference workload, every 10 percentage points of additional prefix-cache hit rate translates to roughly 10% less compute and proportionally less GPU spend.

Key implementation idea

The cache is structured as a tiered key-value store keyed by token-hash prefixes, where the value is a chunk of KV-cache blocks. A typical lookup goes:

  1. Tokenize the incoming prompt, chunk into N-token blocks (e.g., 16 tokens), compute a rolling hash for each prefix.
  2. Walk the hierarchy: check HBM (the engine's own paged-attention table) first; on miss, check host DRAM; on miss, check NVMe; on miss, check the remote pool.
  3. If found in a lower tier, promote: copy the blocks back into HBM (DMA from DRAM, or RDMA from remote) before kicking off prefill for the suffix.
  4. If not found anywhere, do a full prefill, then persist the new blocks back into the hierarchy (likely L2/L3, since L1 is space-constrained).

Two implementation details matter:

  • Block format compatibility. The KV blocks stored at L2/L3/L4 must be in a format that can be DMA'd back into HBM and used directly by PagedAttention. This usually means storing them in the engine's native block layout, with the model version baked into the key. A change in model version invalidates the whole cache.
  • Eviction policy. LRU on access time is the baseline, but more sophisticated systems use cost-aware eviction: weight by (prefix length × recency × access frequency). A long prefix that gets hit twice a day is more valuable to keep than a short prefix hit ten times in an hour.

LMCache is the easiest OSS option to plug into vLLM — it adds a CPU DRAM tier and an optional disk tier transparently. Mooncake Store is the most ambitious: a global key-value store backed by RDMA, designed for petabyte-scale fleet-wide caches. NIXL is the transfer-engine layer (see item 1.4) that makes the cross-tier transfers fast.

Implementation note — tier-aware prefix matching at the router

Once cache lives in multiple tiers, "longest prefix wins" (the rule from item 1.1) is no longer correct. Two nodes can match the same prompt to different depths in different tiers, and the cheapest node to first-token isn't always the one with the longest match. Three approaches are used in production.

The cost model the router is implicitly (or explicitly) minimizing is:

$$\text{TTFT}(\text{node}) = \sum_{i \in \text{tiers}} m_i \cdot c_i \;+\; \left(N - \textstyle\sum_i m_i\right) \cdot r \;+\; q$$

where $m_i$ = matched tokens in tier $i$ (L1/L2/L3/L4), $c_i$ = µs-per-token load cost at tier $i$, $N$ = prompt length, $r$ = µs-per-token prefill (recompute) cost, $q$ = queueing delay at the node.

Typical per-token numbers for Llama-3-70B on TP=8 H100s (Grouped-Query Attention / GQA, ~320 KB KV/token):

Source µs / token vs. recompute (~40 µs)
L1 (HBM)~0essentially free
L2 (host DRAM, PCIe 5.0)5–104−8× faster
L4 (RDMA, 25 GB/s)~15~2.5× faster
L3 (local NVMe, 10 GB/s)~30~1.3× faster (marginal)
Recompute (prefill)~40baseline

Two non-obvious consequences fall out of this table:

  • L4 (RDMA) often beats L3 (NVMe) per byte. Modern data-center fabric outpaces a single NVMe device, so the remote pool can be the second-fastest tier after DRAM.
  • L3 is barely faster than recompute at TP=8. A "long match in L3" is not the strong signal it would have been at TP=1 (where recompute is ~280 µs/token and every tier wins decisively).

Three approaches in production

  1. Longest-prefix-wins with tier as tiebreaker (legacy / simple). Correct only when all candidates are in the same tier — the original 1.1 world before tiered KV existed. Becomes wrong as soon as L2/L3/L4 are in play. Most projects have moved off this once they integrate a multi-tier KV store.
  2. Tier-weighted prefix length (most common). The router scores each candidate as a weighted sum of matched tokens per tier:
    score(node) = w_L1 × matched_L1 + w_L2 × matched_L2 + w_L3 × matched_L3 + w_L4 × matched_L4 − k × queue_depth Typical weights (TP=8, modern Grouped-Query Attention / GQA model): w_L1 = 1.0 (HBM is free) w_L2 = 0.8 (DRAM is fast) w_L4 = 0.6 (RDMA is fast; weighted above L3 on purpose) w_L3 = 0.3 (NVMe barely beats recompute; small credit) k = tuned per fleet
    Used by SGLang router with tier extension and vLLM + LMCache in recent versions. A soft approximation of the cost model that's cheap to compute and easy to reason about.
  3. Explicit TTFT cost model (most accurate). The router evaluates the full formula above for each candidate and picks the lowest. More expensive per routing decision (lookup of per-tier match lengths, arithmetic, sometimes a queueing model), but it naturally handles edge cases like "longer match in L3 vs shorter match in L1" without needing hand-tuned weights. Tunable for different SLO targets (minimize TTFT, minimize p99, maximize goodput). Used by Mooncake Conductor, NVIDIA Dynamo's smart router, and most internal systems at the big labs.

Worked example — when longest match loses

Prompt is 4,096 tokens. Two candidate nodes:

Node A: 4,096 tokens matched, all in L3 (local NVMe) Node B: 2,048 tokens matched, all in L1 (HBM) cost(A) = 4,096 × 30 µs (L3 load) ≈ 123 ms cost(B) = 2,048 × 0 (L1) + 2,048 × 40 µs (recompute second half) ≈ 82 ms Picking by longest match → A → 41 ms slower TTFT. Picking by cost model → B → wins.

The break-even rule that captures this without arithmetic: treat cached tokens as real "savings" only if their tier is meaningfully faster than recompute. L1/L2/L4 tokens are real assets; L3 tokens are close to free recompute equivalents and shouldn't tip a routing decision on their own. Once a tier's per-token load cost exceeds the per-token recompute cost (e.g., L4 under heavy contention dropping to ~50 µs/token), every additional cached token at that tier is a liability, not an asset — the router should ignore that tier entirely for that decision.

References
1.4 · Tier 1

RDMA / GPUDirect KV transport

Importance: MEDIUM (foundational for 1.2 and 1.3)
Origin
GPUDirect RDMA (NVIDIA, ~2013) · NCCL (NVIDIA, ~2016) · NIXL (NVIDIA, 2024) · UCX (open consortium)
OSS engines
NCCL NIXL UCX Mooncake Transfer Engine vLLM / SGLang (use NCCL or NIXL)
What it does & why it matters

Disaggregated serving (item 1.2) and hierarchical KV cache (item 1.3) both rely on moving KV-cache blocks fast, sometimes across the cluster. The default Linux networking stack is far too slow: a TCP socket round-trip is ~30 µs, throughput tops out at single-digit GB/s, and every byte gets copied through kernel buffers and host RAM. For a Llama-70B KV cache of ~1 GB, that's 100+ ms of transfer time — longer than the prefill it's supposed to replace.

The solution stack is the same one HPC has used for 20 years, repackaged for AI:

  • RDMA — the NIC reads/writes remote memory without involving the remote CPU or kernel. Round-trip latency ~1 µs; throughput approaches the NIC's line rate.
  • GPUDirect RDMA — the NIC reads/writes GPU HBM directly, with no detour through host DRAM. On a Llama-70B KV transfer, this saves two memory copies (HBM→DRAM and DRAM→NIC) plus their associated latency.
  • NCCL — collective ops (all-reduce, all-gather) tuned for ML workloads, runs over NVLink (intra-node) or InfiniBand/RoCE with GPUDirect (inter-node).
  • NIXL — NVIDIA's new (2024) unified transfer-engine API. Abstracts NVLink, IB, RoCE, and CPU memory behind one async send/recv interface. Designed specifically for KV transfer in disaggregated inference.

For practical numbers: an H100 node with 8×400G ConnectX-7 NICs and GPUDirect can push KV at ~300 GB/s aggregate over IB. A 1 GB KV cache crosses the network in ~3 ms — comparable to one decode step.

Key implementation idea

Three layers stack to make this work:

  1. Hardware: NICs that support RDMA verbs (Mellanox/NVIDIA ConnectX-6/7, Broadcom Thor2). InfiniBand fabric (200G NDR or 400G XDR is current SOTA) or RoCEv2 on lossless Ethernet.
  2. Driver/firmware: NVIDIA's nvidia-peermem kernel module (or its successor nvidia.ko) exposes GPU memory as an RDMA-addressable region. Mellanox OFED drivers terminate the RDMA verbs in NIC firmware.
  3. Application API: Engines don't call libibverbs directly — they call NCCL, UCX, or NIXL. The library handles connection setup, memory registration, and queue-pair management; the engine just posts a transfer descriptor and gets a completion callback.

A typical KV transfer in a disaggregated cluster:

  • Prefill GPU finishes prefill. Its KV cache for request R is at HBM addresses {p_1, p_2, ..., p_n} (PagedAttention block pointers).
  • The scheduler picks a decode GPU. Decode GPU pre-allocates blocks at {d_1, ..., d_n}.
  • The transfer engine issues N RDMA_WRITEs, source=p_i, dest=d_i. With GPUDirect, the NIC reads from prefill HBM and writes to decode HBM. No CPU touches the bytes.
  • Completion callback fires on the decode GPU. Decode begins.

The performance tricks: pipeline the transfer with the rest of prefill (start sending layer-1 KV while the model is still computing layer-N), use multiple queue pairs in parallel across the NIC's lanes, and register large memory regions once at startup rather than per-transfer (registration is the slow operation).

References
1.5 · Tier 1

Topology-aware parallelism placement (TP / PP / EP)

Importance: HIGH
Origin
Megatron-LM (Shoeybi et al., NVIDIA, 2019; arXiv 1909.08053) · GPipe (Google, NeurIPS '19) · Alpa (Zheng et al., OSDI '22)
OSS engines
vLLM SGLang TensorRT-LLM Megatron-LM DeepSpeed-Inference
What it does & why it matters

A 405B-parameter model (Llama-3.1) at FP16 is ~810 GB — ten H100s' worth of HBM. It cannot run on a single GPU; the model weights must be split across many. How you split them dictates how fast it runs — getting the placement wrong slows you down 5−10×.

Three parallelism dimensions are commonly combined:

  • Tensor Parallelism (TP). Each individual matmul is split across N GPUs. After every layer, an all-reduce sums partial results back to the full activation. Per-layer, per-token. Extremely bandwidth-hungry.
  • Pipeline Parallelism (PP). Different layers are on different GPUs. The activation tensor moves from GPU-1 (layers 1−10) to GPU-2 (layers 11−20). Communication is small (one activation per micro-batch boundary) but adds pipeline bubble latency.
  • Expert Parallelism (EP). For MoE models, different experts live on different GPUs. Per MoE layer, an all-to-all shuffles tokens to their assigned experts and back. See item 2.7.

The pain point: each dimension has very different bandwidth requirements, but the hardware has a sharp bandwidth cliff at the NVLink boundary. Within a node, GPUs talk over NVLink at 900 GB/s. Across nodes, even with the fastest IB (400G XDR), you get ~50 GB/s — an 18× drop. Place TP across that boundary and you've crippled performance.

Key implementation idea: match the topology

The rule that gets tattooed on every inference engineer's forehead:

Parallelism Bandwidth needed Goes on
TPHighest (all-reduce every layer)Inside one NVLink domain (1 node, ≤8 GPUs)
EPHigh (all-to-all every MoE layer)Inside one node, sometimes 2 if linked by NVLink (NVL72)
PPModest (point-to-point per micro-batch)Across nodes via InfiniBand
DP (data parallelism / replication)None during inferenceAcross regions, freely

For a typical 70B-class model on H100s: TP=8 fits inside one node, no PP needed. For Llama-405B or DeepSeek-V3 (671B): TP=8 within a node + PP=2−8 across nodes, with MoE-specific layers using EP within the TP group.

Within those rules, several second-order knobs matter:

  • NCCL algorithm selection. ring vs. tree vs. NVLS (NVLink SHARP, in-network reduction on H100+). Engines often expose this via env vars; the right choice depends on message size.
  • Overlapping communication with compute. Async all-reduces that run concurrently with the next layer's matmul. PyTorch's torch.distributed + CUDA streams handle this; engines like Megatron also do sequence-parallel comm fusion.
  • Avoid TP across PCIe. If a node has 8 GPUs but only 4 are on each PCIe root, splitting TP=8 across the PCIe boundary is catastrophic. The fix: NVLink topology check at engine startup, refuse to start with a bad placement.
  • Decode vs. prefill TP. Decode can often run with smaller TP (TP=2 or TP=4) than prefill, because per-token compute is small. Disaggregation (item 1.2) makes this an explicit choice.
References
1.6 · Production practice (Tier 1)

Multi-region / geo-distributed inference routing

Importance: MEDIUM (HIGH at global scale)
Origin
Operational discipline borrowed from gaming and CDN infrastructure. Academic touchpoint: SkyServe (Yang et al., SoCC '24) on cross-region serving with spot instances.
OSS / commercial
AWS Bedrock cross-region Together AI multi-region SkyPilot / SkyServe Cloudflare Workers AI
What it does & why it matters

Speed of light bounds every cross-continent request. A user in Sydney hitting a US-East GPU pays ~180 ms RTT just for packets, before the model does any compute. For interactive products (chat, code completion, voice), that's enough to make the system feel slow even with perfect inference throughput.

Multi-region inference solves three operational pains at once:

  • Latency. Route each request to a GPU in the same continent (ideally same country) as the user. Brings TTFT floor from "speed of light + compute" down to "compute only."
  • Capacity. No single cloud or region has unlimited GPUs — H100/B200 capacity is rationed. Spreading across providers and regions accesses more total supply, and lets you fall back when one region runs out.
  • Compliance & data residency. EU customers need EU-only processing under GDPR; healthcare data may not leave the country; FedRAMP and similar regimes need defined data paths. Multi-region routing makes these enforceable.

The pain point: a single-region inference service caps your reachable market at customers who can tolerate inter-continental latency. For a serious inference vendor, multi-region is table-stakes the moment you have international users.

Key implementation idea

Four pieces stack to make multi-region inference work:

  1. Edge routing layer. A geographically distributed ingress (anycast IP via BGP, GeoDNS, or per-region edge POPs) terminates the user's TCP/TLS connection at the nearest edge. The edge then forwards to the nearest backend that has capacity.
  2. Per-region engine clusters. Each region runs its own full inference stack — engines, KV caches, prefix routers (item 1.1), schedulers. Regions are mostly independent: a request lives entirely in one region from prefill through last token.
  3. Model weight distribution. Weights are 100s of GB. You can't pull them on-demand on first request — the cold-start latency would be minutes. Strategies: prewarmed images per region, peer-to-peer replication (BitTorrent-style), shared NVMe pools per region. regional container orchestrators handle this as part of "place container in region X."
  4. Failover ladder. When a region is overloaded or down, requests fall through to the next-best region. The fallback is bounded — you'd rather a Sydney user pay +180 ms RTT to US-East than 503 their request.

A subtle but important design choice: KV cache is generally NOT replicated across regions. The RTT to fetch a cached prefix from another continent is longer than recomputing prefill. Cross-region locality matters for routing, but the cache pool is per-region.

What about BYOC (Bring-Your-Own-Cloud)? Some enterprise customers want the inference engine running inside their own VPC for compliance. This is multi-region taken to its extreme — one "region" per customer. The control plane (routing, scheduling, autoscaling) still lives on the vendor's infrastructure; only the engine data path is in the customer's VPC. Several vendors offer this.

References
1.7 · Production practice (Tier 1)

Multi-tenancy & isolation (dedicated vs. shared pools)

Importance: MEDIUM-HIGH
Origin
Operational practice from classical cloud multi-tenancy, adapted to LLM serving. Tooling: AIBrix (ByteDance, 2024) · llm-d (Red Hat, 2024) · KServe.
OSS / commercial
BYOC vendor deployment AIBrix llm-d KServe vLLM Production Stack AWS Bedrock provisioned throughput
What it does & why it matters

Inference services serve many customers from one cluster. The two ends of the spectrum:

  • Serverless / shared pool. All tenants share one big GPU fleet. Pay only for what you use. Cheap floor cost. The risk: a noisy neighbor — a tenant suddenly bursting 10K req/s — can spike everyone else's latency for tens of seconds.
  • Dedicated pool. A tenant gets their own GPUs, isolated from everyone else's traffic. Predictable latency, hard performance guarantees. The cost: you pay for the GPUs whether you use them or not.
  • BYOC (dedicated, in customer's cloud). The dedicated pool runs inside the customer's own VPC. Adds data-residency and compliance guarantees that even “dedicated in vendor's cloud” can't match.

For a commercial inference service, which pool a tenant is in determines their p99 more than almost any other engineering decision. A well-tuned engine in a noisy shared pool can have worse user-visible latency than a poorly-tuned engine on dedicated hardware.

Three concrete pain points multi-tenancy must solve:

  1. Noisy neighbor. One tenant's burst eating GPU time, KV-cache space, or scheduler attention from others.
  2. Cache pollution. If tenant A's system prompts evict tenant B's from the shared prefix cache, B's effective cache hit rate collapses. The fix is per-tenant cache namespacing, but it costs aggregate hit rate.
  3. Data isolation & compliance. Tenant A's prompt or KV state must never be visible to tenant B. For PHI / FedRAMP / sovereign-data workloads, "trust us" is insufficient — you need architectural separation.
Key implementation idea

Modern multi-tenant inference platforms have three operational tiers a tenant can choose between:

Tier Resource model Pros Cons
Serverless (shared) One big pool serves everyone; per-token billing. Cheapest. No capacity planning. Variable p99. Noisy neighbor risk. Cache shared.
On-demand / reserved Tenant reserves N GPUs of capacity; per-hour billing. Predictable latency. Cache is tenant-private. Pay even when idle. Capacity planning needed.
BYOC / private Engine runs inside customer's VPC on their compute. Data never leaves customer's cloud. Compliance. Higher operational burden. Per-customer ops.

Inside the shared (serverless) pool, several mechanisms enforce fairness even without dedicated hardware:

  • Per-tenant token-rate quotas. A tenant submitting more than their share gets queued or rate-limited. Cheap to implement, prevents complete dominance.
  • Fair queueing on the scheduler. The scheduler's admission loop (item 2.1) picks the next request based on each tenant's recent compute usage, not strict FIFO. Weighted-fair-queueing keeps no single tenant's bursty traffic from monopolizing the GPU.
  • Per-tenant prefix-cache namespaces. Cache keys are (tenant_id, prompt_hash) not just prompt_hash. Prevents cache poisoning and information leakage, at the cost of lower aggregate cache hit rate (no cross-tenant sharing of, say, common system prompts).
  • Priority lanes. Paid tenants get a separate higher-priority queue; free-tier traffic only fills the GPU's idle headroom. Common pattern at every inference vendor.
  • KV-cache pressure isolation. Per-tenant caps on how many KV blocks a tenant can hold — prevents one tenant with 10K-token contexts from starving everyone else.

On the data-isolation side: serverless tiers typically run untrusted prompt content in the same engine process, relying on the engine's correctness to keep data separated. Dedicated/BYOC tiers add process- or VM-level isolation — the engine never sees another tenant's data because there are no other tenants on that GPU.

References
2
Single-node technology
Section 2 · what runs on one GPU server
⌃

These are the techniques inside the inference engine itself — the kernels, the scheduler, and the numerical tricks that decide how fast one GPU (or one tensor-parallel group) can serve a request. Most engines (vLLM, SGLang, TensorRT-LLM, llama.cpp, and managed inference stacks) implement variants of every item below.

2.1 · Tier 0

Continuous batching (iteration-level scheduling)

Importance: HIGH
Origin
Orca — Yu, Jeong, Lee, Chun (Seoul National University & FriendliAI), OSDI '22
OSS engines
vLLM SGLang TGI TensorRT-LLM LMDeploy llama.cpp
What it does & why it matters

The old way to batch LLM inference (pre-2022): group N requests into a fixed batch, run all of them to completion, then take the next batch. Two problems:

  • Padding waste. Different requests generate different output lengths. A batch of 16 requests where one needs 500 tokens and the rest need 50 will spend most of its time computing dummy padding for the 15 short requests.
  • Head-of-line blocking. A new request that arrives 1 ms after a batch started can't join — it waits seconds for the batch to finish.

Continuous batching (also called in-flight batching, iteration-level scheduling) makes scheduling decisions at every decode step, not every batch:

  • At each iteration, the scheduler picks which requests are in the active batch.
  • Newly arrived requests that have finished prefill can join the batch.
  • Requests that just emitted their <eos> token leave the batch immediately.
  • Long-running requests continue alongside short ones; no padding waste.

This was the single biggest scheduler innovation in LLM serving. Orca reported ~36× throughput improvement over static batching at the same latency. Today every modern inference engine implements some variant.

Key implementation idea

The engine runs a scheduler loop that, on every decode iteration:

  1. Pulls completed requests out (decoded <eos>, hit max-tokens, or got cancelled).
  2. Examines the admission queue for new requests that have finished prefill.
  3. Admits as many new requests as fit under the max_num_seqs budget and the KV-cache space budget. KV space is the binding constraint — if there's no room for another request's KV blocks (see PagedAttention, item 2.2), the request waits.
  4. Builds the active batch's input tensor: one new token per request, gathered KV blocks per request.
  5. Runs one forward pass; samples the next token per request; streams those tokens to clients.

Two important real-world details:

  • Prefill and decode get interleaved. A new request's prefill is itself a forward pass — either run alone (one big request) or batched with the ongoing decodes (chunked prefill, see item 2.6). Modern engines treat both as "iterations" the scheduler dispatches.
  • CUDA Graphs for decode. Each decode iteration looks identical except for the input tensor. Engines capture the iteration as a CUDA Graph (one graph per batch-size bucket: 1, 2, 4, 8, 16, 32, ...) and replay it. This removes kernel-launch overhead at small batch sizes and is essential for fast decode on small batches.

The scheduler is the engine's hottest loop — it runs once per token output. Most are written in Python (vLLM, SGLang) and optimized to be cheap; some have moved to Rust or C++ (vLLM v1 has a substantially faster scheduler than v0).

References
2.2 · Tier 0

PagedAttention

Importance: HIGH
Origin
vLLM — Kwon et al. (UC Berkeley), SOSP '23. The technique that founded the vLLM project.
OSS engines
vLLM (original) SGLang TGI TensorRT-LLM (variant) LMDeploy (TurboMind)
What it does & why it matters

Before PagedAttention, KV-cache memory was allocated as a single contiguous tensor per request, pre-sized to max_model_len. If your model supports 32K tokens but the average request only uses 512, you've reserved 64× the memory you actually need. Engines either had to allocate the worst case (wasting most of the HBM) or carefully re-pack tensors when sequences shrank (slow and complex).

PagedAttention applies the OS-virtual-memory trick to KV cache:

  • KV memory is divided into fixed-size blocks (typically 16 tokens worth of KV per block).
  • Each request gets a block table — a list of pointers to the blocks that hold its KV state, in token order.
  • Blocks are allocated on demand, one block per ~16 generated tokens.
  • Different requests can occupy different physical blocks; they don't need contiguous memory.
  • The attention kernel reads KV through the block table (a gather, not a strided load) — this is the “PagedAttention kernel,” a variant of FlashAttention that handles non-contiguous KV.

Effective KV utilization goes from ~20−40% (worst-case allocation) to ~95% (block-granularity allocation). On a fixed HBM budget, you can run 2−4× more concurrent requests — which directly multiplies throughput.

PagedAttention also unlocks two other capabilities that depend on it:

  • Prefix sharing. Two requests with the same system prompt can share the same physical blocks for the shared prefix. The block table for each request just points to the same blocks. Reference-counted; freed when the last user finishes.
  • Beam search / parallel sampling. Forking a sequence is just copying the block table, not copying the KV data. Practically free.
Key implementation idea

Three pieces work together:

  1. Block manager. A free-list of fixed-size KV blocks in HBM. Tracks reference counts (a block shared by 5 prefix-matching requests has refcount 5). Allocates blocks on each append_token when the current last block is full.
  2. Block tables. Per-request mapping from logical position (token index / 16) to physical block ID. Stored in CPU and copied to GPU once per scheduler iteration.
  3. Paged attention kernel. A custom CUDA/Triton kernel that, for each query token, gathers its K/V values from the blocks listed in the block table and computes attention on the gathered tensor. Has all the FlashAttention tiling tricks (item 2.3) on top.

The block size is a tunable. Larger blocks (32, 64) reduce block-table overhead but increase internal fragmentation (the last block of a sequence is partially used). Smaller blocks (8) reduce fragmentation but make the kernel's gather more expensive. 16 is the empirical sweet spot for most engines.

A subtle bit: PagedAttention doesn't change the math of attention. It changes only how K/V are stored and fetched. So accuracy is bit-identical to naive contiguous attention; the whole win comes from memory layout.

Extensions: SGLang's RadixAttention generalizes PagedAttention's block sharing into a full prefix tree, sharing not just system prompts but any common prefix that emerges naturally (e.g., few-shot examples, common code snippets). vLLM has since adopted similar prefix caching by default.

References
2.3 · Tier 0

FlashAttention (1 / 2 / 3)

Importance: HIGH
Origin
Tri Dao et al. (Stanford / Princeton / Together AI) — FA1 (NeurIPS '22), FA2 (2023), FA3 (2024).
OSS engines
Every modern engine vLLM SGLang TensorRT-LLM HuggingFace Transformers PyTorch SDPA backend
What it does & why it matters

Standard attention computes the matrix $\mathbf{P} = \text{softmax}(\mathbf{Q}\mathbf{K}^\top / \sqrt{d})$, then $\mathbf{O} = \mathbf{P}\mathbf{V}$. For a sequence of length $N$, this materializes an $N \times N$ attention matrix. At $N = 32\text{K}$, that's a 1 GB intermediate tensor — and it has to live in HBM, the slowest GPU memory.

The bottleneck of attention isn't compute. It's HBM I/O: streaming the $N \times N$ matrix in and out of HBM for the matmuls and softmax dominates. GPUs idle while waiting for memory.

FlashAttention rearranges the computation so the $N \times N$ matrix never gets materialized. $\mathbf{Q}, \mathbf{K}, \mathbf{V}$ are streamed in tiles; softmax is computed incrementally using an online algorithm; the output $\mathbf{O}$ is accumulated tile by tile. All intermediates stay in SRAM (on-chip cache), never touching HBM until the final $\mathbf{O}$ is written out.

Result: attention goes from $\mathcal{O}(N^2)$ HBM I/O to $\mathcal{O}(N)$ HBM I/O. On modern GPUs this is a 2−9× speedup for typical sequence lengths, much more for long contexts. Plus the memory savings — you can now train and serve much longer sequences without OOM.

This single technique made long-context LLMs practical. Almost every model trained or served after 2022 uses FlashAttention or one of its derivatives. It's table-stakes.

Key implementation idea: tiling + online softmax

The classic softmax formula is:

$$\text{softmax}(x_i) = \frac{\exp(x_i - \max(\mathbf{x}))}{\sum_j \exp(x_j - \max(\mathbf{x}))}$$

You need the global max and the global sum to normalize. Naively, you'd have to see all of $\mathbf{x}$ before producing any output — which is exactly the problem.

The online softmax algorithm (Milakov & Gimelshein, 2018) computes this incrementally. As you see tile $\mathbf{x}_{[i:j]}$, you maintain a running max $m$ and a running sum $\ell$. When you see a new tile with a larger max $m'$, you rescale the prior sum:

$$\ell \;\leftarrow\; \ell \cdot \exp(m - m') \;+\; \sum \exp(\mathbf{x}' - m')$$

The same trick lets you accumulate the output $\mathbf{O}$ tile by tile.

So the FA-1 inner loop, per query tile, is:

  1. Load $\mathbf{Q}$-tile into SRAM (sticks around for the inner loop).
  2. For each $\mathbf{K}$-tile, $\mathbf{V}$-tile:
    • Load $\mathbf{K}$-tile, $\mathbf{V}$-tile into SRAM.
    • Compute scores $\mathbf{S} = \mathbf{Q}\text{-tile} \cdot \mathbf{K}\text{-tile}^\top$.
    • Update running max/sum; compute output contribution $\mathbf{O} \mathrel{+}= \text{rescale}(\mathbf{O}) + \mathbf{P} \cdot \mathbf{V}\text{-tile}$.
    • Discard $\mathbf{K}/\mathbf{V}$ tile.
  3. Write final $\mathbf{O}$-tile to HBM.

FA-2 (2023): better work partitioning

FA-1 parallelized over batch and heads only. FA-2 also parallelizes over the sequence dimension within each head, much better utilizing the GPU's SMs (Streaming Multiprocessors). It also rearranges the inner loops so K/V (the larger tensors) are the outer loop, reducing redundant SRAM loads. Result: about 2× faster than FA-1.

FA-3 (2024): exploits Hopper async + FP8

H100 (Hopper) introduced:

  • WGMMA — warp-group matrix-multiply-accumulate. A whole warpgroup (128 threads) issues one async matmul; the next instruction doesn't have to wait. Lets you overlap matmul with softmax.
  • TMA — Tensor Memory Accelerator. Async HBM→SRAM transfers, so memory loads overlap with compute.
  • FP8 tensor cores — double the throughput of BF16.

FA-3 uses all three. The result is ~75% of H100 peak (FP16) and is also the basis for efficient FP8 attention. A typical decode step on Llama-70B at FP8 sees a further 1.5−2× speedup over FA-2 on the same H100.

PagedAttention kernel = FA with a gather

What inference engines actually run is a PagedAttention variant of FlashAttention — the K/V tiles are gathered from the non-contiguous block table (item 2.2) rather than loaded as a strided tensor. Same tiling, same online softmax, plus an indirection through the block pointers.

References
2.4 · Tier 0

Speculative decoding

Importance: HIGH
Origin
Leviathan, Kalman, Matias (Google), ICML '23 · concurrently Chen et al. (DeepMind, 2023). Modern variants: Medusa (Cai et al., 2024), EAGLE (Li et al., 2024), Lookahead (LMSys, 2023), n-gram / prompt lookup.
OSS engines
vLLM (all four variants) SGLang TensorRT-LLM TGI
What it does & why it matters

Decode is fundamentally sequential: token n+1 can only be sampled after token n is known. Each step is one full forward pass through the (large) model. At small batch size, the GPU is mostly idle — HBM-bandwidth-bound, with massive unused compute capacity.

Speculative decoding exploits the spare compute. The insight:

  1. Use a cheap predictor (a small "draft" model, or a head, or even an n-gram lookup) to guess the next K tokens.
  2. Run the large "target" model on all K+1 positions in one forward pass — same wall-clock cost as one normal decode step, since the model is HBM-bound and adding more compute is nearly free.
  3. The target model produces its true probabilities for each position. Verify: for each speculated token, sample from the target's distribution and check if it matches the draft. The longest matching prefix is committed; the first mismatch becomes the resampled true token; everything after is discarded.

If the draft is right about any of the K speculated tokens, you've gotten extra tokens "for free." Modern systems achieve 1.5−3× decode speedup with minimal quality impact — the verification step ensures the sampled distribution is mathematically identical to greedy/non-speculative decoding (for temperature 0; for T>0 there's a rejection-sampling correction that preserves the distribution).

The pain point it solves: decode underutilizes the GPU at small batch sizes. Spec decoding converts spare FLOPs into output tokens.

Key implementation idea — four variants

1. Classic draft-model speculation (Leviathan, ICML '23)

A small model (e.g., Llama-3.2-1B) drafts K tokens by running K sequential forward passes. The large model (Llama-70B) verifies them in one pass via a longer input sequence. Acceptance rate depends on how well-aligned the small model's distribution is to the target's — typically 50−70% per token, yielding ~2× speedup.

Implementation cost: you need to train or distill a draft model. For popular target models there are pretrained drafts; for custom fine-tunes, you have to make your own.

2. Medusa (Cai et al., 2024)

Add multiple lightweight "decoding heads" to the target model itself. Each head predicts one position into the future from the same hidden state. No separate draft model; the heads share the target's representations. Trained on top of a frozen target with a small dataset.

Pros: no separate model to host, simpler infrastructure. Cons: heads see only one hidden state, so they don't model future tokens autoregressively — acceptance rate is lower than EAGLE.

3. EAGLE / EAGLE-2 (Li et al., 2024) — current SOTA

A small auto-regressive draft module operates on the target's hidden states (not on token embeddings). One layer, very cheap. Produces K tokens by re-running on the target's last layer features. Acceptance rate 60−80% on most workloads, leading to 2.5−3× speedup. Currently the best published method by acceptance rate.

EAGLE-2 adds dynamic draft trees: instead of a linear K-token draft, propose a small tree of plausible alternatives at each step, and let the verifier pick the longest accepted branch.

4. n-gram / prompt lookup (no model needed)

Many tokens the model is about to generate already appear verbatim in the prompt — obvious for code completion, RAG (where retrieved chunks reappear in output), JSON extraction, and structured outputs. Build a hash index over the prompt's n-grams; when the most recent K generated tokens match an n-gram in the prompt, propose the continuation as the draft.

Trivial to implement, zero training cost, huge wins on RAG/code workloads (3−5× speedup). Useless on creative-writing workloads where output doesn't mirror input.

Verification: tree attention

For multi-branch drafts (EAGLE-2, Medusa with multiple heads), the verifier needs to evaluate a small tree of candidate continuations in one forward pass. This requires a custom attention mask (each draft token only attends to its ancestor draft tokens, not siblings). Implemented via a "tree mask" passed to the FlashAttention kernel.

References
2.5 · Tier 1

Quantization (FP8 / INT4 / KV-quant)

Importance: HIGH
Origin
Many: GPTQ (Frantar et al., ICLR '23), AWQ (Lin et al., MLSys '24), SmoothQuant (Xiao et al., ICML '23), FP8 (NVIDIA Hopper, 2022), FP4 / NVFP4 (NVIDIA Blackwell, 2024).
OSS engines
vLLM (FP8 / AWQ / GPTQ / FP4) SGLang TensorRT-LLM (FP8 / FP4) llama.cpp (GGUF) FP8 / FP4 managed serving
What it does & why it matters

Decode is HBM-bandwidth-bound. The GPU spends most of its time reading model weights from HBM (for each token, every weight is read once). Halve the size of the weights and you halve the bytes that need to flow over HBM — decode roughly doubles. Quantization is the technique for halving (or quartering) those weight sizes with minimal quality loss.

Three flavors, with very different trade-offs:

  • Weight-only quantization. Weights are stored in INT4 (or INT8, FP8); activations stay in FP16/BF16. Computation requires dequantizing each weight tile back to FP16 before matmul. Wins on memory bandwidth (less to read) but doesn't speed up the matmul itself. Examples: AWQ, GPTQ.
  • Weight + activation quantization. Both weights and activations are quantized; matmul runs natively at the low precision (FP8 tensor cores, or INT8). Wins on memory and compute — FP8 tensor cores are 2× faster than FP16 on H100. Examples: SmoothQuant, FP8 PTQ.
  • KV-cache quantization. A separate axis. Quantize only the K/V tensors stored in HBM (not the model weights). Cuts KV-cache size in half (or quarter), letting you keep more concurrent requests in memory. Examples: KIVI, vLLM's --kv-cache-dtype fp8.

For inference, the killer combo today is FP8 weights + FP8 activations + FP8 KV cache on H100. ~2× throughput vs. BF16 at ~1−2% quality loss. On B200 (Blackwell), NVFP4 pushes another ~2× on top of that, again with modest quality hit.

Key implementation idea

The basic operation: a scale + offset

Quantize: q = round((x - z) / s); dequantize: x ≈ q * s + z. Where s is the scale and z is the zero point. For symmetric quantization (typical for weights), z = 0. Scales are stored alongside the quantized weights.

The art is in choosing where to apply scales. Common granularities:

  • Per-tensor: one scale for the whole weight matrix. Crude, cheap, lossy on outliers.
  • Per-channel: one scale per output channel of the matmul. Better.
  • Per-group: one scale per N consecutive weights (typically N=64 or 128). The default for INT4 weight-only.
  • Per-token (activations only): one scale per token of activation. Common for FP8 activation quantization.

GPTQ (ICLR '23) — INT4 weight-only via second-order analysis

Frames quantization as: pick the INT4 weights that minimize the error in the next layer's activations. Uses the Hessian (computed from a small calibration set) to decide which weights are most sensitive and quantize them first. Other weights are then adjusted to compensate. Result: 4-bit weights with quality close to FP16.

AWQ (MLSys '24) — activation-aware weight quantization

Observation: a small fraction of channels (~1%) carries the biggest activations and is therefore most sensitive. AWQ keeps those channels at higher precision (or scales them up before quantizing). Faster to calibrate than GPTQ, similar or better accuracy. Now the default INT4 method in vLLM.

SmoothQuant (ICML '23) — for activation quantization

Activations have outlier channels that prevent INT8 from working. SmoothQuant migrates the outlier magnitude from activations to weights via a per-channel rescaling (math: X · W = (X / s) · (s · W)), so activations become smooth and quantizable.

FP8 (Hopper, 2022) and FP4 (Blackwell, 2024)

Hardware-native low-precision floats. FP8 has two formats: E4M3 (4 exp, 3 mantissa) for forward pass / weights / activations, and E5M2 (5 exp, 2 mantissa) for gradients. H100's FP8 tensor cores are 2× faster than BF16. The conversion is typically a per-tensor or per-row scale.

FP4 (NVFP4 on Blackwell, MXFP4 as the open-standard variant) extends this another step. Native FP4 tensor cores on B200, 2× faster than FP8. Calibration is harder (less dynamic range), so PTQ techniques like NVIDIA's ModelOpt with double quantization (group scales themselves quantized to FP8) become important.

KV-cache quantization

Store K and V tensors at FP8 (E5M2 usually, because K values have wider range). Halves HBM use for KV; lets you serve ~2× the concurrent requests at the same memory budget. The cost is a small accuracy hit on long contexts (KV from far-back tokens dequantizes to slightly different values). vLLM supports this with --kv-cache-dtype fp8.

References
2.6 · Tier 1

Chunked prefill

Importance: MEDIUM-HIGH
Origin
Sarathi-Serve — Agrawal et al. (Microsoft Research India / Georgia Tech), OSDI '24. Earlier: Sarathi (2023) introduced the idea.
OSS engines
vLLM (0.5+, enabled by default) SGLang TensorRT-LLM DeepSpeed-FastGen
What it does & why it matters

With plain continuous batching (item 2.1), a new request's prefill is one big monolithic forward pass. If a user submits a 4K-token prompt and the engine is currently decoding 100 in-flight requests, the engine has two bad choices:

  • Run the 4K prefill alone: the 100 ongoing decodes stop streaming for hundreds of milliseconds. Existing users see ITL spikes.
  • Defer the prefill: the new user's TTFT balloons while waiting for an idle iteration.

Chunked prefill resolves this by splitting the prefill into chunks (e.g., 512 tokens) and treating each chunk as a normal iteration that the scheduler interleaves with decodes:

  • Iteration 1: chunk 1 of new prefill (tokens 0−511) + 100 decodes → one forward pass.
  • Iteration 2: chunk 2 (tokens 512−1023) + 100 decodes → one forward pass.
  • ... and so on, 8 iterations to absorb the 4K-token prompt.
  • After the last chunk, the new request joins the decode pool.

Decodes never pause. TTFT for the new user is slightly higher than a monolithic prefill (the chunks are interleaved with other work), but ITL for existing users stays flat. The tail latency (p99 ITL, p99 TTFT) improves dramatically — which is what matters in production.

Key implementation idea

Two design choices:

  1. Chunk size. Big enough to amortize per-iteration overhead, small enough to not starve decodes. Sarathi-Serve recommends sizing the chunk to keep total tokens per iteration (chunk + sum of decode tokens) below a hardware-specific threshold. Typical default: 512 or 1024 tokens. vLLM exposes this as max_num_batched_tokens.
  2. Stitching the chunks together. The attention kernel must compute attention for the new chunk's tokens against (a) the prior chunks' KV cache, plus (b) the same chunk's own positions. The KV cache for prior chunks is already stored after each chunked iteration. The PagedAttention kernel handles this naturally via the block-table indirection.

The kernel has to support a slightly more general attention pattern than pure decode: within one forward pass, some sequences contribute multiple query tokens (the prefill chunk, length K) and others contribute one (the decode tokens). This is called a "prefill-decode mixed batch" attention kernel. Both vLLM and SGLang have optimized variants.

Interplay with disaggregation (item 1.2): If you're doing P/D disaggregation, you might think chunked prefill is unnecessary — the prefill pool only does prefill, no decodes to interfere with. In practice, even prefill pools run chunked prefill so that multiple prefills can interleave (a 32K-token request shouldn't completely block a 1K-token one). It's still useful for fairness within the prefill pool.

Interplay with speculative decoding: Tricky. Speculation's verification step puts K+1 tokens per decode, increasing the per-iteration token count. Engines have to budget chunk size to leave room for verifier work, or alternate which iterations do speculation.

References
2.7 · Tier 1

MoE serving (expert parallelism + fused gating)

Importance: HIGH
Origin
Sparsely-Gated MoE — Shazeer et al. (Google, ICLR '17). Modern open MoEs: Mixtral 8x7B (Mistral, 2023), DBRX (Databricks, 2024), DeepSeek-V2/V3 (DeepSeek, 2024), Qwen2.5-MoE (Alibaba, 2024). Inference techniques: DeepSpeed-MoE, Tutel, MegaBlocks.
OSS engines
vLLM (DeepSeek-V3, Mixtral, Qwen-MoE) SGLang (best DeepSeek-V3 perf) TensorRT-LLM DeepSpeed-MoE
What it does & why it matters

A standard dense model activates every parameter for every token. A Mixture-of-Experts model has a much larger total parameter count (e.g., DeepSeek-V3: 671B), but per token activates only a small subset (e.g., 37B for DeepSeek-V3 — 8 of 256 experts). The result: much higher capacity at lower compute cost.

The MoE FFN layer structure:

  • A small "gating" network looks at each token's hidden state and picks the top-K experts (typically K=2 or K=8).
  • Each expert is a regular FFN, but only the top-K run for any given token.
  • The outputs of the K chosen experts are weighted (by the gate's softmax) and summed.

Why this matters for inference:

  • Memory. The full 256 experts must live in HBM, even though only 8 run per token. This dominates the model's HBM footprint — for DeepSeek-V3, the MoE layers are ~600 GB of FP8 weights.
  • Compute. Per token, MoE compute is the same as a 37B dense model. Decode is faster than a 671B dense.
  • Sharding. Experts are the natural unit to shard. Different experts on different GPUs → expert parallelism (EP).
  • Load imbalance. The gating distribution is not uniform — some experts get 10× more tokens than others. This causes hot-spot GPUs and idle ones.

Serving MoE well is a different engineering problem than serving dense. Open-weight MoE deployments made those differences highly visible.

Key implementation idea

Expert parallelism (EP)

Split the N experts across G GPUs. Each GPU holds N/G experts and their weights. Per MoE layer, the forward pass is:

  1. Compute gating. The gating network is small; run it on every GPU (or replicate the result). Output: per-token expert assignments.
  2. All-to-all #1 (dispatch). Send each token's hidden state to the GPU that holds the experts it was assigned to. Each GPU now has the tokens for its experts.
  3. Run the experts. Each GPU computes the FFN forward pass for the tokens it received, on its local experts.
  4. All-to-all #2 (combine). Send the expert outputs back to the original GPUs, where they're weighted and summed per token.

The two all-to-alls are the bottleneck. They run every MoE layer. Their performance depends entirely on the GPU-to-GPU bandwidth — another argument for keeping EP inside the NVLink domain (item 1.5).

Fused gating + dispatch kernels

Naive implementation does several discrete steps (compute gate, sort tokens by expert, permute, run FFN, un-permute, weighted sum). Each step launches kernels, reads/writes HBM.

Production MoE kernels (Tutel, vLLM's fused MoE, SGLang's DeepGEMM for DeepSeek) fuse these into one or two big kernels: a grouped GEMM that processes all experts' work in a single launch, with the token-to-expert routing handled as a permutation inside the kernel.

Load balancing

Train-time techniques (auxiliary load-balance losses) help but don't fully fix runtime imbalance. Inference techniques:

  • Drop-tokens. If too many tokens get routed to one expert (overflowing a per-expert capacity), drop the overflow (their representation just gets the residual). Quality cost is small if capacity is tuned right.
  • Replicating hot experts. Put copies of the most-loaded experts on multiple GPUs. Routing then picks the least-loaded copy.
  • Expert-Choice routing (Zhou et al., NeurIPS '22). Instead of "each token picks K experts," have "each expert picks the top-N tokens it wants." Naturally balanced, slightly different math.

DeepSeek-V3 specifics

DeepSeek-V3 has its own twists worth knowing:

  • 256 routed + 1 shared expert per layer. The shared expert runs for every token; the 256 routed experts pick 8.
  • Auxiliary-loss-free balancing. Uses a per-expert bias term updated by an EMA to balance load without an auxiliary training loss.
  • MLA (Multi-head Latent Attention). Their attention variant compresses KV cache by ~5×. Not strictly MoE-related, but ships with their MoE models.
  • FP8 training and inference in the same model. Their public weights are FP8. This makes them ideal for FP8 inference engines.
References
3
Cross-cutting themes
Chapter 3 · patterns that span both sections
⌃

Sections 1 and 2 covered 14 discrete techniques. But several of them share a single underlying pattern. Naming that pattern lets you reason about inference performance as a system rather than a bag of tricks — which is exactly the framing an inference-startup interviewer is testing for.

3.1 · The PGO analogue

Workload-feedback optimization

Importance: HIGH (as a mental model)
Classical analogue
Compiler Profile-Guided Optimization (PGO) · AutoFDO (Chen et al., CGO '16) · BOLT (Panchenko et al., CGO '19) · Propeller
In inference
Workload optimizer Mooncake elastic ratio SGLang cache-aware router vLLM Production Stack autotuner Custom draft-model training
The pattern (and the analogy that names it)

Classical compilers face an information gap: the source code is static, but real programs run on real input distributions that the compiler can't see at compile time. Profile-Guided Optimization (PGO) closes that gap. Instrument the binary, run it on representative workloads, harvest the runtime profile, feed it back into the compiler. Now the compiler knows which branches are hot, which functions are cold, which loops dominate. It uses that knowledge to pick better inlining decisions, branch placement, function layout, and register allocation. Typical real-world wins: 5−15%. AutoFDO (Google, 2016) removed the instrumentation requirement by using sampled perf profiles from production — the technique can now run continuously, not just at build time.

Modern inference serving has exactly the same information gap. The engine, the kernels, and the model are all static artifacts. But real users have specific prompt distributions, system-prompt reuse patterns, latency requirements, and traffic shapes that no static config can match optimally. Workload-Feedback Optimization — the inference analogue of PGO — observes production, feeds it back, and specializes. The wins are larger than PGO's: workload specialization in inference routinely produces 2−5× cost-per-token improvements over a static-config deployment.

The framework: (Signal, Decision, Cadence)

Any workload-feedback loop is fully described by three properties:

  • Signal — what gets observed from production. Examples: which prefixes are cached on which replica; the prefill : decode queue ratio; the acceptance rate of speculative draft tokens.
  • Decision — what knob the signal turns. Examples: which replica receives a request; whether to repurpose a GPU between prefill and decode pools; whether to enable speculation for a tenant.
  • Cadence — how often the loop closes. Examples: milliseconds (routing decisions), hours (workload-feedback re-tuning), weeks (draft model retraining).

The cadence axis is the most important one and the most underappreciated. Loops that close at different cadences nest: fast loops react to transients within the operating point set by slower loops. Get the cadence wrong — e.g., trying to re-tune quantization on the fly — and the loop either oscillates or never converges.

The cadence ladder — four layers, nested
PRODUCTION TRAFFIC LAYER 1 · ONLINE (ms — sec) Routing tables · scheduler decisions · autoscaling Rebuilds nothing — only in-flight decisions change. The fastest loop, closes every request. aggregated traces LAYER 2 · PERIODIC RE-TUNING (hours — days) Workload optimizer · P/D ratio rebalance · chunk size tuning Rebuilds config (numerical knobs). Same binary, new params. Hot-reloaded or restart-applied. accumulated profile LAYER 3 · RECOMPILATION (per-deploy) CUDA Graphs captured · kernel autotune · quantization bake Rebuilds GPU artifacts. New binary. Released through normal deploy pipeline. distilled trace dataset LAYER 4 · RETRAINING (days — weeks) Draft model training · Medusa/EAGLE heads · quant calibration weights
Cadence ladder: each layer feeds the layers above it as new artifacts. Fast loops above operate on outputs from slow loops below.
The matrix — one row per feedback loop

Each row is one clean (Signal → Decision) pair. Cadence determines the layer.

Layer Signal observed → Decision adjusted Doc anchor
1. Online
ms–sec
Per-replica cached prefixes (the global radix tree) → Which replica receives this incoming request 1.1
1. Online Observed prefill : decode queue ratio → Move a GPU between prefill and decode pools (Mooncake elastic) 1.2
1. Online Per-tenant request rate & KV-cache pressure → Autoscaling decisions; priority-lane admission control 1.7
2. Periodic re-tuning
hours–days
Per-tenant acceptance rate from prior speculation attempts → Spec strategy choice (n-gram vs EAGLE vs off); draft tree depth 2.4
2. Periodic re-tuning Per-tenant prompt-length histogram & system-prompt reuse rate → Chunk size, prefix-cache promotion policy, hierarchical-cache tier weights 2.6, 1.3
2. Periodic re-tuning Per-tenant quality vs. cost SLA observed against held-out evals → Quantization precision per layer (FP16 / FP8 / FP4 mix) 2.5
3. Recompilation
per-deploy
Hot (batch_size, seq_len) buckets seen by the scheduler → Which CUDA Graphs get captured at engine startup 2.1
3. Recompilation Representative tensor shapes (M×N×K from production) → Kernel autotuning choices baked into the deploy (Triton / CUTLASS) 2.3
4. Retraining
days–weeks
Production traces — token sequences & hidden-state distributions → Draft-model weights / Medusa heads / EAGLE module fitted to this workload 2.4
4. Retraining Production prompt distribution used as calibration corpus → Calibrated quantization scales & outlier channels → new weight checkpoint 2.5
Why the cadence axis is the deepest part

The naive way to look at this matrix is "ten optimization techniques." The systems-thinking way is to see that each loop's correct cadence is determined by two things: how fast the signal changes, and how expensive the decision is to apply.

  • Fast-changing signals → fast loops. Replica cache state shifts every request; routing decisions must close in milliseconds. Putting them on a daily cadence would mean routing to GPUs whose caches have completely turned over.
  • Expensive decisions → slow loops. Retraining a draft model is a multi-GPU-day operation; doing it hourly would consume more compute than it saves. Putting it on a weekly cadence lets you amortize the cost across enormous serving volume.
  • Mismatched cadences cause oscillation or staleness. A quant level that's tuned weekly but evaluated against a workload that shifts daily will be perpetually behind. A routing table updated every deploy (hours) will miss the prefix-cache state that's actually evolving every minute.

The art of inference perf engineering is putting each decision at its right cadence. Most production failures aren't from using the wrong technique — they're from running the right technique on the wrong loop.

Interview talking points

Three sentences worth being able to deliver:

  1. “Inference performance is largely a workload-shape problem, not a kernel problem.” — Most of the techniques in this document only deliver their full perf when matched to a specific traffic distribution. The kernel is the floor; workload matching is the ceiling.
  2. “PGO is the right mental model.” — Classical compilers learned this 25 years ago: the compiler can't know what the input distribution is, so harvest it from production and feed it back. Modern inference serving has the same pattern, at multiple cadences.
  3. “Different loops close at different timescales, and that's a design choice, not an accident.” — Online loops handle transients; recompile loops set artifacts; retraining loops shift the operating point. Putting each decision at its right cadence is where production systems live or die.
References
3.2 · Deployment model

BYOC (Bring Your Own Cloud)

Importance: MEDIUM-HIGH
Origin
SaaS deployment pattern from data infrastructure — Snowflake (~2018), Databricks, MongoDB Atlas with private endpoints. Adopted by AI-inference vendors from 2023 onward.
Vendors offering it
BYOC managed serving Together AI private compute Anyscale AWS Bedrock dedicated Mistral private deployments
What it does & why it matters

BYOC is a deployment model where the vendor's software (engine, router, control logic) runs inside the customer's cloud account — in the customer's VPC, on GPUs the customer pays for, governed by the customer's IAM. Compare to the alternatives:

Model Where the engine runs Whose GPUs Whose VPC
SaaS / Serverless Vendor's cloud (shared) Vendor's Vendor's
Dedicated (in vendor cloud) Vendor's cloud (reserved for one tenant) Vendor's, reserved Vendor's
BYOC Customer's cloud Customer's (billed to their cloud account) Customer's

Three reasons customers ask for BYOC, in roughly the order of frequency:

  1. Data residency & sovereignty. "Our prompts contain PHI, financial data, or EU-citizen data. It cannot leave our cloud account, full stop." If the data path never crosses the boundary into the vendor's cloud, regulators and auditors are dramatically easier to satisfy.
  2. Compliance inheritance. Enterprises are already audited (FedRAMP, HIPAA, SOC2 Type II) for their own cloud account. Running the inference workload inside that audited environment inherits the compliance perimeter. Running in the vendor's cloud forces re-auditing of a new boundary.
  3. Cost arbitrage with committed cloud spend. Large enterprises have multi-year committed-spend contracts with AWS, GCP, or Azure (often $100M+/year). They've already paid for those credits. BYOC lets them apply that committed spend to inference compute — much cheaper than paying the vendor's all-in margin on the same GPUs.
The architecture — data plane in customer VPC, control plane in vendor cloud
BYOC ARCHITECTURE · PATTERN A · DATA PATH STAYS IN CUSTOMER VPC CUSTOMER'S CLOUD ACCOUNT e.g., prod, AWS us-east-1 VPC vpc-08a3... Customer app internal service (EKS / Cloud Run / VMs) HTTPS Router (managed) cache-aware load balancing, multi-model dispatch, failover Engine R1 [ GPU ] [ KV cache ] vLLM-style + custom kernels Engine R2 [ GPU ] [ KV cache ] vLLM-style + custom kernels Engine RN [ GPU ] [ KV cache ] vLLM-style + custom kernels streamed tokens (responses) — never leave customer VPC ━━━━━ ACCOUNT BOUNDARY ━━━━━ telemetry ↑ (metadata only) config ↓ (rollouts, tuning) VENDOR CONTROL PLANE runs in vendor's cloud · sees only metadata, never prompts or KV Model Registry model versions, FP8 calibration sets Config & tuning per-tenant spec strategy, cache & chunk knobs Autoscaler queue depth, KV pressure → scale up/down signals Telemetry & Billing latency, hit rate, tokens (no prompts!)
Pattern A BYOC: customer app, router, and engines all live inside the customer's cloud account. The vendor control plane (in vendor's cloud) sees only metadata — never prompts, never KV cache contents.
Why customers still use the vendor's routing

A natural objection: if the customer is bringing their own cloud, GPUs, and VPC, why pay the vendor for routing? Why not just call the engine directly?

The honest meta-answer is that the routing is part of what the vendor sells on top of open-source engine internals. If a customer needed none of the routing's value-add, they should just self-host vLLM and skip the vendor entirely. The specific features that justify keeping the vendor's routing inside the customer's own VPC:

Feature provided by the router What the customer would build themselves without it
Cache-aware prefix routing (item 1.1) A radix-tree-based load balancer that tracks which replica has which prompt prefixes cached. Non-trivial — the SGLang / Mooncake router code, with eviction sync, refcounting, fallback ladders.
Cross-region routing (item 1.6) If the customer has engine clusters in 3+ regions of their cloud, pick the lowest-latency healthy one. Without it, build geo-routing + health checks per region.
Failover & circuit breaking When a pod OOMs or a node degrades, route around it. Requires health-check-aware LB + retry semantics.
Multi-model dispatch "This request is for Llama-3-8B, send to those replicas; this one is DeepSeek-V3, send to those." Routing rules per model, per version.
Per-tenant fairness (item 1.7) If the customer hosts a multi-tenant service on top of a managed inference provider (i.e., they have their own customers), the router can enforce per-sub-tenant quotas.
Autoscaling integration Router emits per-replica load and queue-depth signals that drive the autoscaler. Without it, build your own metrics pipeline.
Spillover to vendor's shared pool "Use my dedicated capacity first, spill to the vendor's serverless if I overflow." The router negotiates this hybrid mode.
A/B & canary rollouts "Send 10% of traffic to the new model version." Routing rules in one place, automated by the control plane.
Workload-feedback tuning (item 3.1) Per-tenant optimization config (spec strategy, quant, cache policy). The router is what receives the tuned config and applies it to traffic decisions.

Take all nine of these away and BYOC reduces to "I'm running open-source vLLM in my own AWS account" — a valid choice, but not what the customer signed a contract for. The vendor value-add is the managed routing + autotuning + ops layer on top of engine internals.

The clean architectural split: control plane vs. data plane

The technical contribution that makes BYOC work as a product is the control-plane / data-plane separation. Each half has a clear job:

  • Data plane (customer's VPC): handles prompts, KV cache, engine processes, GPU work, and response streaming. Nothing in this plane ever talks to the vendor's cloud about prompt contents.
  • Control plane (vendor's cloud): handles model registry, version rollouts, configuration management, autoscaling decisions, telemetry aggregation, billing. Sees only metadata (counts, latencies, queue depths) — never prompt or KV contents.

The bidirectional traffic between them is tiny: kilobytes per minute of telemetry going up, kilobytes per minute of config-change commands coming down. The actual inference traffic (megabytes per request, gigabytes per minute) stays entirely inside the customer's VPC.

Because of this split, the customer's compliance officer can answer "where does our data live?" with a single answer: "In our cloud account. Always. The vendor's infrastructure never sees a prompt or a token." That sentence is the entire sales pitch.

How identification works (and why it's almost trivial)

A subtle but important point: the router in Pattern A doesn't need to identify the customer. It only exists inside that customer's VPC, so by definition every request it receives is from that customer. Identification is implicit in the network topology.

What the router does identify per request:

  • Which sub-tenant (if the customer has their own users/orgs): API key → sub-tenant ID → quota and routing rules.
  • Which model (URL path or header): "this request is for Llama-3-8B" → forward to replicas serving that model.
  • Which replica has the best cache match (radix tree lookup).

The customer's app is typically configured to call something like inference.example.com — a domain owned by the customer, resolvable only inside their VPC. No public exposure of the inference endpoint at all.

References
3.3 · Operations

New pod cold start

Importance: HIGH (gates elastic capacity)
The problem — a naive timeline

Imagine the simplest possible cold start. The autoscaler decides it needs one more replica. The cloud hands you a fresh 8×H100 node. You do the obvious thing end-to-end: pull the container image, download the weights from object storage, page them through host RAM into HBM, start the engine, register with the router. If you do nothing clever, how long until that pod serves its first token?

Work through it stage by stage for a representative model — call it a 200 GB checkpoint (Llama-3 70B in FP16, or a DeepSeek-V3-class model in FP8) on a single 8×H100 SXM node with a 100 Gbps NIC and PCIe Gen4 x16 between host and each GPU:

Stage What's happening Naive time
1. Node provision Cloud allocates the GPU instance; kubelet pulls credentials, attaches volumes, joins the cluster, the scheduler binds the pod. 30 s – 2 min
2. Container image pull Pull the inference image (CUDA + cuDNN + NCCL + Python + vLLM/SGLang/TRT-LLM). Even a slim image runs 5–15 GB; ECR/GCR may throttle. 1 – 5 min
3. Container & driver init Runtime starts, NVIDIA driver and CUDA context init, GPU enumeration, NCCL/NVLink topology probe. 10 – 30 s
4. Weight download (S3 → local NVMe) 200 GB pulled from object storage over a single TCP stream at ~80 MB/s — what you get if you just aws s3 cp in a startup script. ~40 min
5. NVMe → host RAM Engine reads the shards off local disk into pinned host memory. NVMe sustains ~3 GB/s for sequential reads of this size. 60 – 90 s
6. Host RAM → HBM PCIe Gen4 x16 copies the weights into each GPU's HBM. ~25 GB/s effective per link; the 8 links fan out in parallel under tensor parallelism, but bring-up overhead and shard placement keep it from being free. 8 – 15 s
7. Engine warm-up Allocate the KV-cache pool, JIT/compile fused kernels, capture CUDA graphs for the expected batch shapes, run a synthetic forward pass to prime everything. 30 – 90 s
8. Health check + LB registration Router probes the replica, marks it READY, starts shifting traffic to it. 5 – 15 s
Total Fresh GPU node → first token served. ~45 min (typical) · > 1 hr if S3 is slow

The dominant term — by an order of magnitude — is stage 4, the weight download. Forty minutes of a $30/hr H100 node sitting idle while it slurps bytes from S3 is roughly $20 of burned compute per cold start, and that's before the obvious follow-up question: what if the autoscaler needs to add fifty replicas because a customer just turned on a campaign?

Forty minutes is also far longer than any reasonable autoscale latency target. By the time the new pod is ready, the traffic spike that triggered it is over and the autoscaler is permanently chasing a workload it can't keep up with. So the naive path doesn't just waste money — it makes elastic capacity structurally impossible.

Every technique in the rest of this section exists to compress one of these stages, usually by orders of magnitude. The interesting question isn't “can we make cold start faster?” — it's which stage do we attack, and what does the architecture have to look like to make the attack cheap?

Why local NVMe forces a re-download on every fresh node

The first instinct on reading the naive timeline is “why not just pre-stage a persistent volume with the weights, then attach it on cold start?” That's almost right, but there's an architectural constraint that rules out the most obvious version of it.

The local NVMe on a GPU node isn't just for weights — it's the L3 tier of the KV-cache hierarchy (item 1.3). HBM is L1, host DRAM is L2, the on-node NVMe is L3 where prefix-cache pages spill when DRAM fills up. That tier has to sustain 6–12 GB/s at sub-millisecond latency, which only physically-attached instance-store NVMe delivers. Network-attached block storage tops out at ~1–4 GB/s (EBS gp3 / io2) with network jitter on top — wrong shape for KV spillover.

And instance-store NVMe is lifecycle-bound to the node. It's ephemeral, can't be detached, can't be reattached elsewhere, gone when the instance terminates. So there's no “pre-populate this NVMe in a separate workflow and then attach it to a fresh pod” option. On a freshly-allocated GPU node, the local NVMe is empty, period.

Two narrow escape hatches exist, neither of them “attach a pre-warmed volume”:

  • DaemonSet host-path prewarming. A node-level daemon pulls weights to host NVMe before a serving pod schedules. The NVMe is still ephemeral and node-bound, but “cold pod on a warm node” (pod restart, rolling update, model swap) becomes free. Doesn't help on a truly fresh node.
  • Bake weights into the node image. Stage 1 implicitly carries the weights with it. Works, but AMIs over ~100 GB get painful, you re-bake on every weight version, and you lose the flexibility to land any model on any node — node identity becomes model-coupled. Mostly avoided except for a “platinum” tier of always-on models.

So the design space splits cleanly: weights must come from somewhere external on every fresh node. The rest of this section is about where “external” lives, and how cheap we can make the transfer.

From PVCs to shared filesystems — a quick primer

If your mental model of cloud storage is “PVC,” you're probably thinking of block storage: a virtual disk that one VM mounts, presented as a block device, formatted with ext4/xfs. EBS, GCP Persistent Disk, Azure Managed Disk. These are single-attach (one mount at a time), zonal (live in one AZ), and the backing media is a network-replicated block service.

A shared filesystem is a different shape of cloud storage: multi-attach native (many nodes mount the same FS simultaneously), regional (cross-AZ), and presents as a POSIX filesystem tree rather than a block device. It's also designed for substantially higher per-client bandwidth than block storage.

The naming conventions across clouds:

Storage shape AWS GCP Azure
Block PV (single-attach, zonal) EBS Persistent Disk (PD), Hyperdisk Managed Disk
Shared FS, general NFS
~1–3 GB/s per client
EFS Filestore Azure Files / NetApp Files
Shared FS, parallel / AI-purposed
~5–15 GB/s per client
FSx for Lustre (PERSISTENT_2) Hyperdisk ML, Parallelstore Azure Managed Lustre
Self-managed parallel FS on cloud VMs WekaFS, VAST Data, BeeGFS, self-hosted Lustre — deployed on instances with local NVMe; you operate the cluster, the cloud provides the VMs
Object store mounted as POSIX
cheap, slow
Mountpoint for S3, S3 + Alluxio/JuiceFS GCS FUSE BlobFuse

For inference weight delivery, the row that matters is the AI-purposed parallel FS. FSx Lustre on AWS, Hyperdisk ML on GCP, or third-party WekaFS / VAST running on instances. These deliver 5–15 GB/s per client, support GPUDirect Storage (DMA from FS to HBM bypassing host RAM), and handle the multi-attach fanout that block PVs can't.

The conceptual jump from PVC to shared FS is exactly the jump from “one volume, one mounter” to “one filesystem, many mounters.” In Kubernetes both still show up as PVCs, but the underlying CSI driver and capabilities (in particular, ReadWriteMany vs ReadWriteOnce access modes) are different.

Shared FS pattern: one mount, directory-per-model

With a regional shared FS in place, the architecture collapses into something elegant:

  • One FS volume per region, mounted on every GPU node — not one PV per model.
  • Directory-per-model layout: /weights/llama-3-70b/v2.1/, /weights/deepseek-v3/v3.2/, ...
  • Control-plane job populates the FS from S3 once per model version (off the hot path; minutes is fine).
  • Pod boot: mount is already there (node-level, managed by a DaemonSet or node-mode CSI); engine reads the right directory into HBM; pod registers with the router.
  • Serving registry tells each pod which model(s) to keep resident, tracks placement across the fleet, drives autoscaling and model promotion/demotion.

Time-to-first-token: ~30–60 s on a fresh pod where the node is already warm (FS already mounted). On a freshly-allocated node you still pay for node provision + image pull, but stage 4 of the naive timeline (the 40-minute weight download) collapses to a streaming FS read at parallel-FS bandwidth.

The cost shape is the unintuitive part: provision the FS for peak fanout bandwidth, not steady-state capacity. 50 pods cold-starting same-model at 7 GB/s each is 350 GB/s aggregate — on FSx Lustre PERSISTENT_2 that requires roughly 350 TB of provisioned capacity just to hit the bandwidth ceiling, regardless of how much you actually store. WekaFS scales differently (size the cluster for the burst rather than the storage), but the principle holds: plan for the bandwidth burst, not the gigabytes.

Wait — doesn't this still go through the NIC?

A careful reader hits two questions at this point:

  1. If the FS is remote, doesn't the data still have to land on local NVMe first before the engine can use it? — No.
  2. If the data traverses the same NIC as an S3 download, how does it hit 7 GB/s when S3 hits 80 MB/s? Aren't they wire-bound by the same NIC? — The NIC is the same; the protocol stack on top of it isn't.

Where the bytes actually go. The shared FS path skips local NVMe entirely:

S3 path:                                  Shared FS path:
─────────                                 ────────────────
S3 object store                           FSx Lustre / WekaFS
   │                                          │
   │ NIC  (network read)                      │ NIC  (network read)
   ▼                                          ▼
local NVMe  (write)  ← extra hop          host DRAM (pinned)
   │                                          │
   │ NVMe read                                │ PCIe
   ▼                                          ▼
host DRAM                                  HBM
   │
   │ PCIe
   ▼
HBM

The engine mmap()s the file on the FS; when it touches a page, the kernel pulls that page over the network straight into the OS page cache (in DRAM, not on disk). Then cudaMemcpy ships it to HBM. No local file is ever created. "Loading weights" and "transferring weights over the network" are the same syscall — whereas S3 needs an explicit aws s3 cp that lands a file on disk, then a separate engine read off disk. Saves a whole NVMe write+read round-trip.

Why the bandwidth differs despite sharing the NIC. A 100 Gbps NIC delivers ~12.5 GB/s in either direction. S3's 80 MB/s is two orders of magnitude below NIC capacity, so the NIC clearly isn't the bottleneck. The throttling lives in the protocol stack:

Layer S3 (HTTPS) Parallel FS (Lustre / Weka)
App aws s3 cp, parses HTTP responses mmap() / read() direct from user space
Protocol HTTP/1.1 framing, chunked-transfer encoding Native FS protocol (LNet for Lustre, Weka's proprietary)
Encryption HTTPS/TLS — CPU-heavy, ~1–3 GB/s per core Usually disabled in-VPC; or hardware-offloaded
Transport TCP, one connection per GET, congestion-controlled RDMA (kernel-bypass) or many parallel TCP
Backend topology One S3 frontend per request — implicit per-key throughput cap (~80 MB/s) One file striped across N storage servers (OSTs); they all serve in parallel
NIC 100 Gbps 100 Gbps — same hardware

The four sources of the gap, in order of impact:

  1. Striped backend. A single S3 GET hits one frontend that serves one object from one backend shard — bandwidth is capped per-key. Parallel FS stripes one file across 32–64 storage servers; reading it opens parallel RPCs to all of them. Aggregate bandwidth scales with stripe count, not per-request.
  2. RDMA kernel-bypass. RDMA DMAs data straight from the remote NIC into local pinned memory, skipping the kernel TCP/IP stack. One CPU core can push 10+ GB/s because the CPU isn't actually moving bytes — the NIC and DMA engines are. TCP+TLS over the same NIC caps around 1–3 GB/s per core because of kernel copies and packet processing.
  3. No HTTPS/TLS tax. S3 is HTTPS; every byte must be encrypted + authenticated. At 10 GB/s that's ~80 Gbps of TLS — CPU-bound on any reasonable core count. In-VPC FS traffic typically skips encryption (network is already isolated) or uses hardware offload.
  4. Parallelism is native, not bolted-on. To approach FS speeds with S3, you have to manually issue 32–256 parallel HTTP range requests (s5cmd, aria2). It works up to ~3 GB/s before TLS CPU and connection management overhead make further scaling pointless. The FS client does this fanout natively, with RDMA, with no TLS.

The realistic bandwidth ladder:

Path Streams Per-stream bottleneck Aggregate
aws s3 cp (default) 1 S3 per-key + TLS ~80 MB/s
s5cmd, 32 ranges 32 TLS CPU + S3 frontend ~1.5 GB/s
s5cmd, 128+ ranges, tuned 128+ TLS CPU saturates cores ~3 GB/s (asymptote)
Lustre client over TCP, 32 OSTs 32 Kernel copies ~5 GB/s
Lustre / Weka over RDMA, 32 OSTs 32 NIC ~10 GB/s
100 Gbps NIC ceiling — wire ~12.5 GB/s

So the answer to “why does FS go faster on the same NIC?” is: it doesn't, at the NIC level — the NIC can hit ~12 GB/s on either path. The difference is that FS gets close to the NIC ceiling while S3 leaves ~99% of it unused on the floor, because of HTTPS overhead, single-flow TCP, and per-key backend throttling. S3 is built for many small objects with strong durability guarantees; parallel FS is built for few clients reading huge files fast. They're making different trade-offs at every layer above the wire.

Snapshot + clone: the BYOC alternative

What if you can't run a shared FS — for example, you're deploying into a customer's VPC for compliance (BYOC, item 3.2) and there isn't a WekaFS cluster sitting there waiting for you? The fallback is a different pattern: golden snapshot + per-pod clone, riding on block-storage PVs.

The flow:

  1. Downloader — a one-time CPU pod with a fresh empty block PV attached. Pulls weights from S3 in parallel (s5cmd, aria2), saturates its NIC. Maybe 5 minutes for 200 GB. Off the hot path.
  2. Snapshot — create a volume snapshot of the populated PV. This becomes the golden snapshot for the model version.
  3. Detach — the downloader exits; the master PV can be deleted or kept for re-snapshotting on future weight updates.
  4. Clone per pod — on every cold start, the CSI driver creates a new PV from the snapshot, attaches to the GPU node, mounts it. Pod reads weights from the (now-feels-local) block PV.

Kubernetes has first-class support for this via VolumeSnapshot resources and the dataSource field on PVCSpec. AWS EBS, GCP PD, and Azure Disk CSI drivers all implement it.

The catch is that “warm PVC” is a misleading label. The snapshot is warm; the cloned PVC isn't. Each clone pays:

  • CSI provisioning — CreateVolume from snapshot + attach + mount, ~20–40 s of K8s and cloud-API round-trips before any data moves.
  • Lazy snapshot restore — EBS volumes from snapshot are “available” instantly, but blocks are read-through-S3 on first touch. Without Fast Snapshot Restore (~$0.75/AZ/hr per snapshot), initial reads run at near-S3 latency.
  • Block-storage bandwidth ceiling — even with FSR enabled, EBS gp3 is ~1 GB/s and io2 Block Express ~4 GB/s. That's 2–3× slower than a shared FS read at the same step.
  • Per-AZ duplication — block snapshots are zonal in practice; cross-AZ requires snapshot copies.

Realistic time-to-first-token: 2–5 min per cold-started pod, vs ~30–60 s for shared FS. The pattern's value is operational simplicity (no FS cluster to run) and BYOC-compatibility (works in any customer VPC), not low latency. If cold-start latency is the priority, this is a Tier-2 design, not Tier-1.

The three methods, side by side — latency breakdown for a new pod

Putting the paths next to each other with concrete numbers for a 200 GB model on an 8×H100 node with a 100 Gbps NIC. The first table is the warm-node case (existing fleet, just adding a pod) — the scenario where the three methods actually diverge. The second adds node-level cold start on top, since stages 1–3 are paid regardless.

Cold pod on a warm node — this is the case that matters for elastic autoscale:

Stage Naive S3 → local NVMe Snapshot + clone (block PV) Shared FS (Lustre / Weka)
Storage setup ~0 s (NVMe is just empty disk) 20–40 s (CSI CreateVolume from snapshot + AttachVolume + mount) ~0 s (FS already mounted at node level by DaemonSet/CSI)
Data over the network 2–40 min
200 GB ÷ 80 MB/s (single-stream) = ~40 min
200 GB ÷ 1.5 GB/s (s5cmd parallel) = ~2.2 min
50 s – 25 min
With Fast Snapshot Restore: 200 GB ÷ 4 GB/s (io2 BX) = 50 s
Without FSR (lazy restore): ~150 MB/s sustained over EBS's internal fetch-from-snapshot-store path → ~20–25 min for 200 GB. Notably slower than parallel s5cmd directly from S3 to a fresh EBS, so FSR is effectively mandatory for this pattern to pencil out.
~30 s
200 GB ÷ 7 GB/s (FS → DRAM directly, no NVMe hop)
Local NVMe write Bundled with download above (curl/s5cmd writes file) Bundled (lazy-restored blocks land on the cloned PV) skipped entirely
Disk → host DRAM ~25 s (200 GB ÷ 8 GB/s NVMe) ~50–200 s (200 GB ÷ 1–4 GB/s EBS gp3/io2) fused with the network read above
DRAM → HBM (PCIe) 8–15 s 8–15 s 8–15 s
Engine warm-up (KV pool, kernels, CUDA graphs) 30–90 s 30–90 s 30–90 s
Health check + LB registration 5–15 s 5–15 s 5–15 s
Total (warm node) ~3–42 min
(~3 min realistic with parallel s5cmd)
~2–5 min (with FSR + io2)
~22–28 min without FSR — worse than just doing parallel s5cmd from S3
~75–150 s

Adding node cold start — what you pay when the autoscaler provisions a fresh GPU instance (these stages are identical across all three methods):

Pre-pod stage All three methods
Node provision (cloud allocates instance, kubelet joins cluster) 30 s – 2 min
Container image pull (CUDA + engine, 5–15 GB) 1 – 5 min
Container + driver init (CUDA context, NCCL probe) 10 – 30 s
FS mount (shared FS only) or block-PV attach (snapshot only) ~5 s / handled in storage setup above
Cold-node overhead ~2 – 7 min added to all three columns above

The big picture in three numbers (cold pod on warm node):

  • Shared FS: ~1–2.5 min — the floor for cold-start without a pre-warmed pool.
  • Snapshot + clone (with FSR): ~2–5 min — the BYOC fallback, ~2× slower than shared FS, mostly from CSI provisioning + block-storage bandwidth ceiling.
  • Naive S3 download: 3–42 min — only viable if you tune aggressively (s5cmd parallel); the “default” behavior is the 40-minute disaster from the opening timeline.

Three things stand out from the breakdown:

  1. The 30–90 s engine warm-up + 8–15 s PCIe transfer are inescapable across all methods. Even with infinite network bandwidth, you can't get a cold pod below ~50 s. That's the floor — to beat it you need DaemonSet host-prewarm or HBM-warm standby pods (next subsection).
  2. Snapshot+clone's overhead isn't the snapshot — it's the CSI provisioning ceremony plus the EBS bandwidth ceiling. The 20–40 s of K8s and cloud-API round-trips before any data moves, combined with EBS being 2–3× slower than parallel FS at the read step, is where the gap with shared FS comes from.
  3. Shared FS wins primarily because it skips local NVMe entirely. No write-then-read round-trip; the network read is the load step. That ~50–80 s of NVMe round-trip the other paths pay just disappears.

Lustre and WekaFS — what they actually are

Lustre is an open-source parallel filesystem that originated in the HPC world around 1999 (Carnegie Mellon, then Sun, then various stewards; currently maintained by the OpenSFS / EOFS community). The classic use case is national-lab supercomputers — tens of thousands of compute nodes reading and writing checkpoints to a shared scratch volume at terabytes per second. Architecture: a metadata server (MDS) holds the namespace, and object storage targets (OSTs) hold file data striped across them. A single 200 GB file is sharded into stripes (default 1 MB or 4 MB) spread across N OSTs; a client reading it opens parallel RPCs to all N and aggregates the bandwidth. Lustre speaks its own wire protocol (LNet), supports RDMA over InfiniBand/RoCE, and gives you a POSIX mount point. In the cloud you don't operate Lustre yourself — you use AWS FSx for Lustre, Azure Managed Lustre, or GCP Parallelstore, all of which are managed Lustre under the hood.

WekaFS (Weka.IO, founded 2013) is a commercial software-defined parallel filesystem built from scratch for the NVMe + RDMA era. Instead of being open-source you operate it yourself or buy a managed offering; you deploy Weka as a cluster of VMs/bare-metal nodes that each contribute local NVMe and a NIC, and Weka stripes data across them. The pitch vs Lustre: better small-file and metadata performance, native support for GPUDirect Storage (DMA from FS to GPU HBM, skipping host RAM), unified POSIX + S3 + NFS gateways, and a tighter operational story for AI workloads specifically. It costs more, but it's the parallel FS most modern AI clouds run on — CoreWeave, Lambda, and several inference vendors deploy Weka as their primary weight store. VAST Data and DDN's EXAScaler (a managed Lustre distribution) are the main competitors in the same niche.

What they share: striped file layout across many backend nodes, RDMA-capable wire protocols, POSIX mount semantics, designed for few clients reading huge files fast rather than many clients reading many small objects durably (which is what S3 is for). That's why the bandwidth numbers in the previous subsection look the way they do — the filesystem itself is architecturally a fanout-bandwidth engine, not a single-server NFS box.

4
Interview questions
Section 4 · operational scenarios, ordered by what matters most
⌃

These are operations-flavored prompts — the kind of question that follows “walk me through your serving stack” in a senior interview. Each item explains what the problem actually is (the failure mode you're being asked to defend against), then lists production strategies in priority order: the things you must do, then the things that help, then the nice-to-haves.

4.1 · Operations

Traffic-spike handling

Importance: HIGH
What it is

Traffic is bursty by nature: a customer turns on a campaign, a viral post hits, a cron job kicks off across a fleet of agents, or a recovering region's traffic failover lands on you. The aggregate QPS to one model can double or 10× in seconds. The serving system has to absorb that without (a) blowing up the p95 of in-flight requests, (b) sending the GPU into thrash, or (c) silently dropping traffic.

The defining constraint is that you cannot scale fast enough to meet the spike. New pod cold start (item 3.3) is on the order of minutes even with optimized weight loading — a spike that doubles in 30 seconds will have already inflicted its damage before any new replica comes online. So “just autoscale” is not an answer. Spike handling is buying time for autoscaling to catch up, and degrading gracefully if it can't.

Three failure modes the strategy must defend against, in priority order:

  1. Cascading p95 blow-up. The admission queue grows; queued requests miss TTFT SLO; clients retry; the queue grows faster. This is the dangerous one because it self-amplifies.
  2. KV-cache exhaustion. Too many concurrent sessions admitted; the engine OOMs and crashes, taking all in-flight traffic with it. Catastrophic failure mode.
  3. Cross-tenant collateral damage. One tenant's burst eats GPU time that paying SLA tenants are entitled to.
Production strategy — high-priority items first

Tier 0 — you must have these or you don't have a serving system:

  1. Admission control with backpressure 503s. When the admission queue (see stage 4 of the lifecycle, line ~1397) exceeds threshold, return 503 Service Unavailable with a Retry-After header rather than letting the queue grow unbounded. This is non-negotiable. Quoting the doc: “If the queue exceeds a backpressure threshold, the engine returns 503 rather than letting the queue grow unbounded.”
    Why it's Tier 0: without this, a spike turns into a death spiral. With it, you cap worst-case queue delay $q$ in the TTFT formula $\text{TTFT} = \sum_i m_i c_i + (N - \sum_i m_i) r + q$ and protect every request already in flight.
  2. Hard caps on per-tenant concurrency & KV-cache footprint. Cross-link to item 1.7: per-tenant token-rate quotas, per-tenant KV-block caps, weighted fair queueing in the scheduler. Without these, one tenant's spike starves everyone.
  3. Autoscale on KV-cache occupancy and queue depth — never on CPU. CPU utilization is meaningless for GPU work. The two signals that actually matter form a tier:
    • KV-cache occupancy is the leading indicator. Saturation always shows up here first — the scheduler can't admit a new request until KV blocks are free, so a filling cache predicts queueing before any queue forms. Scale up when KV occupancy crosses ~80−85%, which buys ~30 s of warning to bring a new replica online.
    • Queue depth is the confirming symptom. By the time the admission queue is growing, KV is already pinned and you're already late — it's a useful tripwire for the case where KV is fine but the scheduler can't keep up (rare, but real: per-tenant rate limits, scheduler CPU bottlenecks). Treat queue depth as the backup trigger, not the primary.
    • p95 latency is the damage report. By the time it moves, users are already hurting. Use it for alerting, not for autoscaling decisions — it's a lagging indicator of failures the first two should have caught.
    A healthy autoscaler triggers on whichever fires first: KV occupancy > 85% for T seconds OR queue depth growing for T seconds. KV gets you the early warning; queue depth catches the edge cases. Scale-down requires both to be low for a much longer window, to avoid flapping.

Tier 1 — the buffer that buys time for Tier 0:

  1. Standby / warm-pool capacity. Keep N replicas above current demand so the next spike has somewhere to land before cold-start completes. Size the standby pool from the math:
    $$\text{standby replicas} \;\ge\; \text{peak ramp rate (req/s)} \;\times\; \text{cold-start time (s)} \;/\; \text{per-replica capacity (req/s)}$$
    If cold-start is 5 min and traffic can double in 2 min, standby must equal current capacity. If you can get cold-start to 30 s with snapshot-based provisioning (item 3.3), standby shrinks proportionally — this is why cold-start latency is load-bearing for cost.
  2. Priority-lane shedding. Free-tier and batch-priority traffic gets shed first under load; paid tenants and SLA tenants keep their lane (item 1.7). Failure mode without this: a free-tier burst takes down your SLA tenants' latency — the worst possible business outcome.
  3. Client-side retry with jittered backoff. Make the 503s effective. Retry storms without jitter are how a recoverable spike becomes an outage. SDKs should default to exponential backoff with full jitter; document this for any tenant building their own client.

Tier 2 — helpful but not always justified:

  1. Predictive pre-warm. Cron-driven warm-ups for known traffic shapes (Monday 9am, end-of-month batch jobs, scheduled product launches). Cheap if you can identify the pattern; useless if traffic is genuinely unpredictable.
  2. Cross-region overflow. Route spike traffic to a peer region with spare capacity (item 1.6). Works for stateless prefill but pays an inter-region latency tax; mostly used for disaster scenarios, not garden-variety spikes.
  3. Quality-of-service degradation. Under extreme load, downgrade long-prompt traffic to a smaller model variant or truncate prompts. Controversial: changes the answer, not just the latency. Use only with explicit tenant opt-in.
What to say in an interview

A senior-level answer leads with the architectural reality: “You can't scale fast enough to catch a real spike — cold start is minutes, spikes are seconds. So spike handling is about (a) shedding load cleanly with 503s before the queue self-amplifies, (b) per-tenant caps so one customer can't take down others, and (c) standby capacity sized to cold-start latency. Predictive pre-warm and cross-region overflow are nice-to-haves on top.”

The trap to avoid: jumping straight to “autoscaling.” Autoscaling is necessary but it's the last line of defense, not the first — because it's the slowest. Interviewers are watching for whether you grasp that ordering.

4.2 · Operations

Hot / cold model switching

Importance: HIGH (multi-model serving)
What it is

Terminology, since this trips people up:

  • Hot model — weights are loaded in HBM on at least one replica, scheduler is running, can serve in milliseconds. TTFT is in the “normal” 100−200 ms range from your worked example.
  • Cold model — weights need to be paged in from object store / shared FS / EBS into HBM. First-request TTFT is in the seconds-to-minutes range. A 200 GB checkpoint takes ~45 min naively (item 3.3), or 2−5 min with snapshots, or 30−60 s with shared-FS warm-restore on a warm node.
  • Hot/cold switching — the operation of changing which model a GPU is serving: evict one model's weights, load another's. Forced on you whenever you can't keep every model permanently resident.

You hit this the moment you serve more distinct models than fit in your aggregate HBM. Concrete example: 200 fine-tunes (200 GB each at FP8) on 20 H100 nodes (8 × 80 GB = 640 GB per node, ~500 GB usable for weights after KV reserve, so 2 models per node max). You can keep 40 models hot. The other 160 must multiplex — loaded on demand, evicted when idle.

The failure mode this defends against: a request lands on a model that's cold, the user waits 30 seconds for first token, and your interactive SLO is destroyed for that tenant. The deeper failure mode is the “thrash” case — if every request hits a different model and evictions chase admissions, you spend all your HBM bandwidth swapping weights and never actually serve any tokens.

Why this is hard

The arithmetic is brutal. Loading a 200 GB checkpoint from object storage at a realistic 1 GB/s parallel throughput is ~200 seconds just for the S3→NVMe stage. From a warm node's local NVMe to HBM is another ~75 seconds. That is the absolute floor of how long a cold load takes — nothing else you can do in software changes those numbers. Every strategy below is about either avoiding the cold load, hiding it, or pre-paying it.

Key cost-vs-latency tradeoff:

$$\text{HBM cost} \;\propto\; \text{models kept hot}$$

Every model you pin uses HBM that could serve KV cache for other requests. Pinning everything is impossibly expensive on the long tail. The whole problem is deciding which models to keep hot, and how to handle the rest gracefully.

Production strategy — high-priority items first

Tier 0 — required for any multi-model platform:

  1. Keep the head of the distribution permanently hot. Track per-model RPS. The top-N models that fit in your aggregate HBM are pinned — never evicted. This handles 80−95% of traffic with zero cold-start exposure. The hot set is your single biggest lever.
  2. Model-residency aware routing. The router has to know which replicas have model X loaded. A request for X is routed to a replica where X is hot, even if a closer replica is cold. This is a small addition to the prefix-aware router from item 1.1 — same data plane, extra dimension. Without this, you can have replicas idle with the right model while busy replicas thrash on the wrong one.
  3. Eviction on idle, not on admission. Page a model out only after it's been idle for a long-enough window (minutes, not seconds). Evicting a model that gets a request 100 ms later is the worst case — you paid to evict and now you pay to reload. LFU with a long decay window beats raw LRU here because recency lies on bursty traffic.

Tier 1 — what makes the cold path tolerable:

  1. Snapshot-based fast restore for cold loads. Cuts cold start from ~45 min (naive S3 pull) to 2−5 min (snapshot+clone) or 30−60 s (shared-FS warm-restore on a warm node). See the full breakdown in item 3.3. This is the same machinery used for new-pod cold start; the multi-model case just exercises it more often.
  2. Pre-warm on signal. If you know tenant A tends to use models X, Y, Z together, warm Y and Z when the first request for X arrives. Costs HBM and load-time, saves the second and third request's cold start. Worth it when models co-cluster by tenant or by workflow.
  3. Separate the cold-start SLO from the hot SLO. Don't mix them in one p95 — that's dishonest math that misleads your own tuning. Report “hot TTFT p95: 200 ms; cold-load p95: 45 s; cold-fraction: 0.3%” separately. This also lets you set legitimately different latency budgets for first-request vs. subsequent-request traffic.
  4. Surface the cold state to the client. On a cold load, return an event like {"status": "warming", "eta_seconds": 30} instead of silently holding the connection. Lets the client UI show a meaningful spinner and prevents retries from triggering more cold loads on more replicas.

Tier 2 — advanced or workload-specific:

  1. LoRA / adapter swapping. If your “different models” are actually one base model + many LoRA adapters (~100 MB each instead of 200 GB), you don't have a hot/cold problem — you have an adapter-swap problem, orders of magnitude cheaper. The base stays hot; adapters page in milliseconds. vLLM, SGLang, and TGI all support multi-LoRA serving. If you have flexibility on the modeling side, push for LoRA over full fine-tunes.
  2. Tiered pools by model class. Run a dedicated pool for “hot base models” (pinned, low-latency) and a separate “cold pool” (multiplexed, best-effort SLO, cheaper hardware). Tenants opt into one or the other at the API level. Mirrors the multi-tenancy tiers in item 1.7.
  3. Peer-to-peer weight distribution. When the autoscaler needs to load model X onto 20 cold pods simultaneously, pulling all 20 copies from S3 saturates the egress link. BitTorrent-style fanout between pods drops aggregate load time substantially. Worth it only at large fleet sizes — below 10 pods the implementation cost isn't justified.
What to say in an interview

A senior answer frames it as a caching problem, not a loading problem: “HBM is the cache, models are the entries, and you can't fit them all. So I'd pin the top-N by RPS, route requests to replicas where the model is already hot, evict on long idle to avoid thrash, and keep snapshot-based fast restore as the path for the long tail. The cold-load SLO is reported separately from the hot SLO — mixing them hides the real distribution. If the 'different models' are LoRA adapters, the problem mostly evaporates — that's the question I'd ask first.”

The trap to avoid: treating this as “just load weights faster.” The interviewer is watching for whether you recognize that the cold path is unavoidable on the long tail, and the engineering is about minimizing how often you traverse it and how much it hurts when you do.

4.3 · Operations

Fast weight rollouts at scale (model updates across a cluster)

Importance: HIGH
What it is

How do you push a new version of a model — same architecture, new weights — across a cluster of hundreds of replicas, without downtime, quality regressions, or breaking in-flight conversations? This is not the same as deploying new code. A code deploy moves a 100 MB container image and Kubernetes handles the rest in minutes. A weight deploy moves a 100−500 GB checkpoint to every replica, and the naive math is ugly:

$$\text{naive bytes to move} \;=\; N_{\text{replicas}} \times W_{\text{weights}}$$

200 replicas × 200 GB = 40 TB. At S3's effective per-prefix throughput (~50−100 GB/s shared across the fleet), pulling that in parallel saturates the prefix and your VPC NAT — throughput collapses and the rollout takes hours instead of minutes.

That's just the distribution problem. Once weights are on each node, you still have to swap them on a running engine (the engine holds the old weights in HBM with in-flight requests using them), and you still have to shift traffic safely (a quality regression in v2 can't be allowed to hurt all users at once). The interview test is whether you recognize that this is three distinct problems stacked, each with its own scaling bottleneck.

The three independent layers

A clean mental model: each layer has its own design choices, and they compose independently. Total rollout time is roughly:

$$T_{\text{rollout}} \;\approx\; T_{\text{distribute per node}} \;+\; T_{\text{swap per replica}} \times \frac{N_{\text{replicas}}}{B_{\text{parallel}}}$$

where $B_{\text{parallel}}$ is the deploy strategy's batch size (1 for blue/green, small for canary, larger for aggressive rolling).

Layer The question it answers Options (cheap → sophisticated)
1. Weight distribution How do bytes get to each node? On-demand S3 pull · pre-staged PV/NVMe · shared parallel FS · P2P (BitTorrent-style)
2. Per-replica swap How does one replica switch v1 → v2? Cold swap (drain + restart) · Hot swap (in-place, needs 2× HBM)
3. Traffic orchestration How does the fleet transition without breaking users? Rolling · Canary · Blue/green · Shadow

The layers can't paper over each other's weaknesses. If distribution is slow (S3 pull on demand), no orchestration strategy helps — every replica still waits 30 minutes for bytes. If the per-replica swap is slow, even blue/green hurts.

Production strategy — layer by layer

Layer 1 · Weight distribution (priority order):

  1. Pre-stage to local NVMe via DaemonSet. A per-node daemon pulls new weights from S3 to host-path local storage asynchronously during off-peak hours, decoupled from the rollout window. When the rollout fires, the engine reads from local NVMe (~3 GB/s sustained) instead of S3 (~80 MB/s on a single stream). This is Tier 0 — without it, no other layer's fast strategy matters.
  2. Shared parallel filesystem (WekaFS, Lustre, FSx for Lustre) for fresh rollouts that can't be pre-staged. 50−100 GB/s aggregate read bandwidth, handles fanout naturally. Tradeoff: filesystem cost, single point of contention.
  3. Peer-to-peer distribution for cluster-wide simultaneous deploys. N replicas help each other fetch from a small set of seeds, so aggregate bandwidth scales with $N$ instead of being capped by S3 egress. Tools: NVIDIA's nimble, BitTorrent-derived libraries, custom implementations at hyperscalers. Implementation cost is non-trivial — only worth it at fleet sizes >50.
  4. Per-region CDN / object cache for cross-region rollouts. Reduces cross-region S3 egress (expensive) and improves first-replica latency in each region.

Most large platforms layer these: DaemonSet pre-stages popular models nightly, shared FS handles fresh in-region deploys, P2P kicks in for cluster-wide moves, S3 is the source of truth. See item 3.3 for the full cold-start machinery that backs this.

Layer 2 · Per-replica swap:

  1. Cold swap (the common default).
    1. Router drains traffic from the replica (stop sending new requests; let in-flight finish).
    2. Engine stops, weight files replaced from local NVMe (milliseconds if pre-staged), engine restarts.
    3. Engine warms up: KV pool allocation, CUDA-graph recapture, synthetic forward pass.
    4. Health check passes; router puts replica back in rotation.
    Total: ~30 s – 2 min per replica with weights pre-staged. Replica is out of rotation the entire time, so your fleet temporarily loses one unit of capacity per concurrent swap.
  2. Hot swap (in-place weight replacement).
    1. Engine stays running; traffic continues hitting the replica.
    2. New weights stream into a second HBM allocation alongside the old ones.
    3. At a quiescent scheduler boundary, atomically flip the weight pointer.
    4. Free old weights; invalidate KV and prefix caches (now stale w.r.t. new weights).
    Total: ~5–15 s per replica, no rotation removal. But requires 2× weight HBM headroom during the swap — most production deployments don't have that, so this is rare in practice. Used mainly for small models or when KV pool can be temporarily shrunk to make room.

Layer 3 · Traffic orchestration:

Strategy Capacity needed Wall-clock When to use
Rolling 1× Slow (hours) Low-risk minor updates; cost-constrained envs
Canary + rolling 1.05–1.2× Medium (~70 min for a 200-replica fleet) Default for most rollouts — ramp 1% → 5% → 25% → 100% with quality gates between
Blue/green 2× Fast (limited by distribution + swap) Major/risky updates; instant rollback required
Shadow 2× (compute), 1× (serving) Variable Before promoting to canary — mirror traffic to v2, discard responses, compare quality offline

Canary + rolling is the sweet spot for most rollouts. 1.05−1.2× cost overhead, quality regressions caught at 1% blast radius, instant rollback via router config (~10 s) if anything goes wrong. Blue/green is reserved for major updates where 2× cost during the window is justified by the safety guarantee.

Cold-swap walkthrough — which component does what (shared-FS example)

Concrete example for Layer 2 cold swap on a node where weights are served from a shared parallel filesystem (WekaFS / Lustre / FSx for Lustre). Every step is owned by a specific component — the value of the table is seeing where the boundaries are, because each boundary is a separately tunable layer.

# Component (owner) What happens
1 kube-scheduler
K8s control plane
Binds the pod to a new node based on GPU resource requests, affinity rules, taints.
2 kubelet + CSI driver
Node-local K8s + storage vendor
kubelet asks the CSI driver to attach + mount the shared FS PVC. Sub-second on warm CSI; 5–10 s on cold attach. Result: /mnt/weights/llama-70b/v2/ is visible inside the pod.
3 containerd (or cri-o)
Container runtime
Pulls the engine image (cached locally if recent), starts the container, runs the engine binary as the entrypoint with args like --model /mnt/weights/llama-70b/v2/.
4 Engine entrypoint
vLLM LLMEngine.__init__ / SGLang ModelRunner / TRT-LLM Executor
Reads engine config, parses CLI args, initializes torch + CUDA context, sets up tensor-parallel process group across the 8 GPUs (NCCL init).
5 ModelLoader
Engine module (e.g. vllm.model_executor.model_loader.DefaultModelLoader)
Walks /mnt/weights/llama-70b/v2/, parses config.json, identifies the safetensors shards (model-00001-of-00030.safetensors …), picks the loader strategy.
6 safetensors library
Third-party Python lib, calls into libc
safe_open() invokes mmap(2) on each shard file — see the callout below. No bytes read yet; just virtual-memory mappings created.
7 Per-rank Worker + parallel_state
One engine subprocess per GPU; uses TP/PP sharding rules from vllm.distributed.parallel_state
For each named parameter (e.g. model.layers.0.self_attn.q_proj.weight): compute this rank's slice (1/TP of the rows for a column-parallel layer); touch the corresponding mmap pages.
8 Linux kernel (VFS + page cache + FS driver)
Kernel space, transparent to the engine
Touched mmap page → page fault → kernel reads bytes from the shared FS through the page cache into host RAM, maps the page into the engine's virtual address space. Engine just sees bytes appear in its pointer.
9 CUDA driver (cudaMemcpyAsync)
NVIDIA driver + runtime, invoked via PyTorch's .to(device)
Copies the slice from host RAM → HBM over PCIe (~25 GB/s per link, 8 links in parallel). Once copied, host RAM page can be evicted — data lives in HBM now.
10 Worker.initialize_cache
Engine module
After all weights are placed, measures remaining HBM, computes block count for KV pool, allocates the PagedAttention block pool, builds block-table data structures.
11 CUDAGraphRunner
Engine compilation module + CUDA runtime
Runs a synthetic forward pass at each pre-decided batch-size bucket (1, 2, 4, 8, 16, 32, …), captures the kernel sequence as a CUDA Graph for replay. 30–90 s — the bulk of "engine warm-up."
12 kubelet readiness probe + Router
K8s + service mesh / inference router
kubelet polls the engine's /health endpoint; once it returns 200, kubelet marks the pod READY; the router (Dynamo / Envoy / custom) starts shifting traffic to this replica.

Where mmap fits — the userspace/kernel bridge. mmap is the mechanism that makes step 8 efficient. It deserves a closer look because understanding it is the difference between "the engine reads weights" and "the engine and kernel cooperate to stream weights into HBM."

USERSPACE (engine process) KERNEL HARDWARE ───────────────────────────────────────────────────────────────────────────────────────── safe_open("model-001.safetensors") │ └─► mmap(fd, len, PROT_READ, ──► sets up page-table entries MAP_PRIVATE) (NO disk I/O yet) returns virtual address V engine touches V + offset ──────► page fault! │ └─► VFS asks FS driver to fetch ──► shared FS read (WekaFS / Lustre client over RDMA or TCP) │ populate page cache ◄────────────── bytes arrive │ map page into engine's address space │ ◄─────────────┘ engine now sees bytes at V+offset engine: torch.Tensor.to('cuda') │ └─► cudaMemcpyAsync(hbm_dst, virt_src=V+offset, ──► CUDA driver pins host page, len) issues DMA descriptor ──► PCIe DMA host RAM → HBM

Three properties of mmap that make this work at scale:

  1. Lazy. Mapping a 200 GB file costs ~microseconds (just page-table setup). Actual bytes are read only when touched. Engines that load weights in layer order never have more than one layer's bytes resident in host RAM at a time, even though the whole 200 GB is "mapped."
  2. Page-cache shared. If two ranks on the same node touch the same shard (e.g., embeddings or shared layers in PP), the kernel reads the file once and serves both from the page cache. Multi-GPU loads on one node naturally dedupe.
  3. Zero-copy compatible. The pointer the engine hands to cudaMemcpyAsync points into the page cache; CUDA can DMA directly from there with no intermediate copy. (Pinning the page first speeds it up further; tools like NVIDIA GDS can DMA straight from FS into HBM, skipping host RAM entirely.)

Byte-flow summary for the shared-FS cold swap:

shared FS ──► FS client (kernel) ──► page cache (host RAM) [WekaFS/Lustre over RDMA] [populated lazily by mmap page faults] │ ▼ cudaMemcpyAsync (CUDA driver, PCIe DMA) │ ▼ HBM (GPU)

Each arrow is a separately tunable layer. FS bandwidth (Lustre stripe count, WekaFS client threads) tunes arrow 1. mmap behavior + page-cache size tunes arrow 2. cudaMemcpyAsync + pinned memory + stream parallelism tunes arrow 3. The slowest arrow dominates the cold-swap wall-clock — usually arrow 1 unless your FS is beefy, in which case it shifts to arrow 3 (PCIe-bound).

Cross-cutting concerns (not a layer but you'll be asked)
  • Cache invalidation. Weights changing means KV and prefix-cache entries computed under v1 are mathematically invalid for v2 — different weights produce different K/V tensors. On every per-replica swap, both caches must be flushed. Prefix-cache hit rate drops to 0% immediately after swap and recovers over minutes as the new cache warms. This is a real compute cost during the ramp window — budget for it.
  • Conversation continuity. Multi-turn chat sessions started on v1 shouldn't mid-conversation switch to v2 — the assistant's "voice" changes, which users notice. Solution: sticky routing by session_id during transition. The router tracks session → version mapping; new sessions prefer v2, existing sessions stay on v1 until they end. Drain-then-switch is the cruder alternative.
  • Per-region staging. Roll one region (or one AZ) at a time. If region A breaks, regions B/C are untouched. Standard blast-radius limitation, free safety win.
  • S3 thundering herd. If 200 replicas all pull at once you get throttled by S3's per-prefix rate limits, plus your VPC NAT chokes. Stagger initial fetches even if pre-staged. The DaemonSet should also stagger its background pulls.
A concrete production playbook

200-replica deploy of Llama-70B v2 in one region, weights pre-staged the night before via DaemonSet:

T-12h: DaemonSet pulls v2 weights from S3 to local NVMe on all 200 nodes (off-peak, async, no traffic impact) T+0: Bring up 10 v2 replicas (5% of fleet), warm from local NVMe (~60s each) T+5m: Route 1% of traffic to v2 canary, monitor for 15 min — TTFT, ITL, error rate, output-quality comparison vs v1 T+20m: Quality + latency look fine → scale v2 to 50, route 25% T+35m: Continue ramp → 100 v2 replicas, route 50% T+50m: v2 at full capacity (200 replicas), route 100%, drain v1 sessions T+65m: v1 replicas at 0 in-flight → decommission T+70m: Rollout complete; ~70 min wall-clock, ~5% peak over-provisioning ROLLBACK PATH (any time): Router config flip → 100% traffic back to v1 (~10 s) Keep v1 fleet alive until v2 is debugged
What to say in an interview

A senior answer leads with the three-layer model and the scale problem: “Three problems stacked: distributing 100s of GB to 100s of replicas without saturating S3 or NAT, swapping weights on each replica without long downtime per node, and orchestrating the traffic shift without breaking quality or in-flight conversations. For distribution, pre-stage to local NVMe via DaemonSet during quiet hours; use a shared FS or P2P for fresh rollouts. Per-replica cold swap from pre-staged files is ~60 s — hot swap is faster but needs 2× HBM headroom most deployments don't have. For orchestration, canary at 1−5% with quality gates is the sweet spot — 1× cost overhead, ~70 minutes wall-clock, instant rollback via router config. The non-obvious traps are cache invalidation (prefix-cache hit rate drops to 0% on swap), conversation continuity (sticky routing by session_id during transition), and the S3 thundering-herd problem during initial fetch.”

The trap to avoid: treating it as a Kubernetes rolling-deploy problem. The interviewer is watching for whether you recognize that distribution dominates wall-clock at LLM scale, and that the per-replica swap and traffic orchestration questions are downstream of solving the distribution problem first. Most candidates jump to canary-vs-blue-green; the senior signal is starting with "how do bytes get to the node."

4.4 · Operations

Fault tolerance & monitoring

Importance: HIGH
What it is

Inference serving fails differently from stateless web services. A failed GPU node doesn't just drop a request — it loses 32−128 in-flight streaming sessions, holds 60 GB of KV cache that needs to be recomputed elsewhere, and takes 30−60 s of weight-load time to replace. A failed router can blackhole an entire pool because there's no DNS-style fallback for "send to a replica with this prefix cached." So fault tolerance for inference is less about preventing failures (you can't — GPUs die, ECC errors spike, kernels panic) and more about containing the blast radius and recovering fast enough that users don't notice. Monitoring is the prerequisite: you can't recover from what you can't see, and the metrics that matter for LLM serving aren't the ones a standard web stack ships with.

Two questions to answer concretely: (a) what happens when a router or GPU node dies, and (b) what do you watch on the dashboard so you know before users do.

Failure modes — what dies and what happens
Failure What's affected Detection Recovery
Single router pod Requests in flight on that pod (likely zero — router is stateless per request); none that haven't hit it yet K8s liveness probe (5−10 s) Other router replicas absorb traffic immediately; replacement pod up in ~10 s
Entire router fleet All new traffic to the pool blackholes — in-flight streams continue to client until they end External health check / synthetic monitoring Multi-AZ replicas + LB front; if entire fleet is down, traffic must fail over to another region or be 503'd
GPU node (hardware or kernel crash) All in-flight requests on that replica are lost (typically 32−128 concurrent sessions); KV cache evaporates; weights need to be reloaded on the replacement kubelet readiness probe fail (~30 s), router heartbeat timeout (~5−10 s) Router evicts replica from routing table (~1 s after detection); traffic re-routes to remaining replicas; autoscaler spawns replacement (30−60 s if snapshot-restore is set up, item 3.3; 30−60 min naive)
Engine crash (OOM, kernel bug) Same as GPU node failure from the request perspective — in-flight lost Process exit, kubelet restart, ~5−15 s Same engine restarts on the same node; weights re-loaded; KV reinitialized; ~30−60 s back in rotation
NVLink / TP-group fault A TP-sharded replica becomes useless (the all-reduce fails) NCCL timeout in the engine, then crash Whole replica goes down (not just the bad GPU); needs full node replacement
Object store (S3) failure Can't cold-load new replicas — existing replicas serve fine because weights are already in HBM S3 5xx rate spike; cold-start pods stuck Wait it out; pre-staged local weights save you; multi-region buckets reduce blast radius
Auth service down Can't issue new JWTs; existing JWTs work until expiry Auth /token endpoint error rate HA auth service with read replicas; cache JWKS aggressively (1 hr TTL) so verification keeps working
Redis (rate limit / control plane) Rate-limit checks fail; routing-state lookups stale Redis client errors Fail-open for rate limits (better than blocking all traffic); router falls back to last-known good state
Network partition (cross-AZ) Replicas in one AZ unreachable from others; cross-AZ KV routing fails Network probe failures; KV-transfer timeouts in disaggregated serving Pool-internal traffic stays AZ-local; degrade to higher latency by routing across surviving AZs

The brutal asymmetry: a router pod dying loses ~zero requests; a GPU node dying loses 32−128. That's because **router is stateless per request, GPU replica is stateful** (KV cache + in-flight decode). Architecture choices flow from this asymmetry — routers can be casual about replication (3 pods suffices), GPU nodes need careful blast-radius design (cross-AZ, overprovisioning, fast replacement).

Production strategy for resilience — high-priority items first

Tier 0 — you must have these:

  1. N+1 (or N+2) capacity per pool. Always run enough replicas that losing one doesn't push the rest over saturation. For a pool sized to handle peak with 10 replicas at 80% KV occupancy, run 11−12. The extra cost is small; the alternative is "one failure cascades into latency spike for everyone."
  2. Cross-AZ replica placement. K8s pod anti-affinity rules spread replicas across availability zones. A single AZ outage shouldn't take down the whole pool. The router must be aware of zones too — don't route 100% of a pool's traffic through one AZ if other AZs have capacity.
  3. Fast health checks + auto-eviction. Router heartbeats to engines every 1−2 s; mark replica unhealthy after 3 consecutive misses (~5 s detection). Detection time is the dominant component of recovery time — tune this aggressively. Slow detection means the router keeps sending traffic into a black hole.
  4. Idempotency on retries. When a request fails and the client retries, you must not double-bill or double-execute. Idempotency keys (client-generated UUID in a header, deduped at the gateway for ~10 minutes) are standard. Without this, every transient failure during a long generation can result in two charges and two billings.

Tier 1 — what makes recovery fast:

  1. Pre-warmed standby pool. Keep 1−2 fully-loaded extra replicas per critical pool, not serving traffic. When a primary dies, promote a standby instantly (router config update, ~10 s). This is the difference between "30 s recovery" and "60 s recovery" — standby skips the cold-load phase entirely.
  2. Snapshot-based fast restore. Item 3.3 — cuts new-pod cold start from 30+ min to 30−60 s. Gates how fast you can scale into a failure event. Without this, a single GPU failure during a traffic spike is unrecoverable in real time.
  3. Client-side retry with jittered backoff + circuit breaker. SDKs should default to retry on 5xx and connection-reset, with exponential backoff + full jitter. Pair with a client-side circuit breaker: if 10 consecutive requests fail to a region, stop hitting it for 30 s and try a peer region.
  4. Bulkhead by pool. A failure in the LoRA pool shouldn't take down the foundation pool. Separate K8s namespaces, separate routers, separate autoscalers. The control plane is shared; the data plane is bulkheaded.

Tier 2 — advanced patterns:

  1. Request hedging for VIP traffic — send to two replicas, take whichever responds first, cancel the loser. Doubles cost; used selectively for latency-critical requests on critical models.
  2. Streaming reconnect tokens. When a streaming response is cut off mid-flight (replica died), the client can reconnect with a token that resumes from the last delivered token instead of re-prompting from scratch. Requires the engine to checkpoint output to a fast store (Redis); not all platforms support this.
  3. Multi-region active-active. Run the same fleet in 2+ regions, with global load balancing. Worst-case single-region outage just shifts traffic with ~50−200 ms latency added. Expensive; reserved for very high-availability tiers.
Key metrics to monitor — tiered by what they tell you

The metrics fall into five tiers. The order matters: user-visible SLO metrics tell you something already broke; engine-internal metrics tell you it's about to break; hardware metrics catch what the engine can't see. All three layers are needed.

Tier 1 — SLO metrics (user-visible health, page on miss):

MetricWhy it matters
ttft_seconds (p50, p95, p99)Time to first token — the most-watched latency metric. Page if p95 > SLO for 2 min.
itl_seconds / tpot_seconds (p50, p95, p99)Inter-token latency during streaming. Determines how smooth the response feels.
e2e_request_latency_seconds (p50, p95, p99)End-to-end wall clock. Catches issues TTFT alone misses (e.g., decode slowdown).
error_rate (4xx, 5xx separately)5xx = your fault; 4xx = client's fault. Track separately. Page on 5xx spike.
goodput_tokens_per_secondTokens/sec delivered within SLO. Not raw throughput — the honest serving capacity.
streaming_disconnect_rateSessions that ended mid-stream due to error. Catches failures the latency metrics miss.

Tier 2 — engine internals (leading indicators, drive autoscaling):

Metric (vLLM names)Why it matters
vllm:gpu_cache_usage_percKV-cache occupancy — the leading indicator for saturation. Scale on this, not CPU.
vllm:num_requests_waitingAdmission queue depth — if this is non-zero you're already saturated.
vllm:num_requests_runningActive batch size — capped at max_num_seqs; sanity check.
vllm:num_requests_swappedRequests preempted to host RAM — non-zero means you're in pain (KV exhaustion).
vllm:prefix_cache_hits_total / queries_totalPrefix-cache hit rate. Trend matters — sudden drop = workload shift or eviction churn.
vllm:time_to_first_token_seconds (engine view)TTFT as measured by the engine (excludes network); compare to gateway TTFT to attribute latency.

Tier 3 — hardware (DCGM exporter, catches what engine misses):

MetricWhy it matters
DCGM_FI_DEV_FB_USED / FB_TOTALHBM used vs capacity (hardware view). Cross-check with engine's KV occupancy.
DCGM_FI_DEV_MEM_COPY_UTILHBM bandwidth utilization — the binding resource for decode. ~100% on healthy decode-bound serving.
DCGM_FI_DEV_GPU_UTILSM utilization. Less useful for inference (decode is memory-bound) but tracked for sanity.
DCGM_FI_DEV_POWER_USAGEWatts — proxy for actual work being done. Sudden drop on a "healthy" replica = something silently wrong.
DCGM_FI_DEV_ECC_*ECC error counts. Rising rate predicts GPU hardware failure — drain proactively.
DCGM_FI_DEV_NVLINK_BANDWIDTH_*NVLink throughput per link. Asymmetry across links signals a degraded link.

Tier 4 — workload characterization (drives capacity planning):

MetricWhy it matters
request_prompt_tokens distribution (p50/p95)Prompt-length distribution. Shifts here change your KV math overnight.
request_generation_tokens distributionOutput-length distribution. Determines decode cost per request.
requests_per_second (per model, per tenant)Traffic rate. Per-model breakdown drives autoscaling and hot/cold decisions.
cache_hit_rate per modelLower than expected = workload shift, eviction problem, or routing misconfiguration.

Tier 5 — infra / operational (catches platform-level issues):

MetricWhy it matters
pod_restart_countEngine crashes / OOMs. Should be ~0; any spike = investigate.
cold_start_duration_secondsTime from pod scheduled to first request served. Gates your scale-up speed.
weight_load_duration_secondsSub-component of cold start — isolates FS/distribution problems.
autoscaler_decisions_totalScale-up and scale-down events. Frequent oscillation = thrashing, fix your hysteresis.
cost_per_million_tokensDerived metric (GPU-hours / tokens served). The business-level health check.
Alerting tiers — what wakes someone up
PAGE (wake on-call): - p95 TTFT > SLO for 5 min - 5xx error rate > 1% for 2 min - Pool capacity unavailable (autoscaler can't bring up new pods) - Region-wide outage signals (router blackholing, DNS errors) - ECC error storm on multiple GPUs (potential hardware fleet issue) SLACK / WARN (notify, not page): - KV cache > 85% for 2 min on any pool - Queue depth growing without scale-up - Cold-start latency above baseline (FS / distribution degrading) - Prefix-cache hit rate dropped > 10 pp from baseline - Preemption rate non-zero on any replica TICKET (investigate during business hours): - Per-replica throughput regressed > 5% vs fleet median - ECC errors incrementing slowly on a single GPU - Cost per token trending up week-over-week - Pod restart count above zero on any deployment

The principle: page on user-visible damage, warn on saturation predictors, ticket on slow degradation. Paging on warning-level metrics burns on-call; failing to page on damage hurts users. The split is load-bearing.

Cross-cutting principles
  • Detection time dominates recovery time. 10-second detection + 30-second replacement = 40-second user impact. 60-second detection + 30-second replacement = 90 seconds. Tune health-check frequency aggressively — it's cheap and it's the highest-leverage knob.
  • State asymmetry drives blast radius. Stateless components (router, gateway) tolerate failure cheaply — just kill and replace. Stateful components (engine replicas with KV cache, in-flight streams) are expensive to lose. Spend engineering on the stateful side.
  • Multi-AZ > single-AZ over-provisioning. 12 replicas in one AZ is more failure-vulnerable than 6+6 across two AZs, despite same capacity. An AZ going down is a routine cloud event.
  • Three views of the same number catch different bugs. Engine's KV occupancy + DCGM's HBM-used + workload's avg context length should all correlate. When they diverge, something is wrong (memory leak, accounting bug, another tenant on the GPU).
  • The dashboard order matters. Top row is autoscale signals (KV, queue, preemption). Second row is SLO (TTFT, ITL, error rate). Third is workload (prompt/output length, RPS, cache hit). Fourth is hardware (DCGM). If the SLO row lights up before the autoscale row, your thresholds are wrong.
What to say in an interview

A senior answer separates the two questions cleanly: “For fault tolerance: the asymmetry that matters is router-vs-GPU-node. Router failure loses ~zero requests because routers are stateless per request, so 3-replica HA + LB handles it. GPU-node failure loses 32−128 in-flight sessions and 60 GB of KV cache — far more expensive. The defense is N+1 capacity per pool, cross-AZ placement, fast health checks (5-10 s detection), snapshot-based fast restore (30-60 s replacement), and idempotency keys for retries. Pre-warmed standby pools cut recovery further for critical models. For monitoring: five tiers — SLO metrics (TTFT/ITL/error rate, page on miss); engine internals (KV occupancy, queue depth, preemption rate — the leading indicators that drive autoscaling); hardware via DCGM (HBM bandwidth, ECC errors); workload characterization (prompt/output length distribution drives capacity planning); and infra/operational (cold-start time, pod restarts, cost per token). Page on user-visible damage, warn on saturation predictors, ticket on slow degradation.”

The trap to avoid: treating LLM serving like a stateless web service. CPU-based autoscaling, simple round-robin load balancing, and "just restart the pod" don't work because the per-replica state (KV cache, in-flight streams, model weights in HBM) makes both failure and recovery much more expensive than people expect. The interviewer wants to hear that you understand why LLM fault tolerance is different, not just a list of standard SRE practices.

4.5 · Architecture

Model registry & routing data model

Importance: HIGH
What it is

The control-plane data model that connects three layers: what models exist (the registry), what's deployed (the pools), and which pool serves a given request (the routing rules). A naive design assumes one model = one pool, but real platforms have to support many-to-many cardinality on both axes:

  • Single-model fan-out: one foundation model in many pools (regions, tiers, quantizations) — e.g., llama-70b-v2 in us-east-1, us-west-2, eu-west-1, plus a canary v3 pool, plus dedicated pools per tenant.
  • LoRA fan-in: many adapter model_ids share one pool (a base model + many adapters all loaded in HBM together).
  • Colocated multi-model fan-in: many small full models packed into one pool, sharing HBM (embedders, small LLMs, classifiers).

Without a schema that handles all three, you either constrain the platform's deployment patterns or you build special cases everywhere. The right data model treats model_id as a logical identifier and pool as physical capacity, joined by a routing table — same shape as DNS hostname → IP, or service mesh service name → endpoint.

The three pool types — visual taxonomy

1. Single-model pool — one (model, version, quantization) per pool. Fans out across regions/tiers/quants:

SINGLE-MODEL POOLS (one model per pool, model fans out across regions/tiers) llama-70b-v2 ──┬──► [llama-70b-v2-shared-us-east-1] 8×H100 fp8 ├──► [llama-70b-v2-shared-us-west-2] 8×H100 fp8 ├──► [llama-70b-v2-shared-eu-west-1] 8×H100 fp8 ├──► [llama-70b-v2-dedicated-exampleco] 8×H100 fp8 (tenant=exampleco) ├──► [llama-70b-v2-batch-pool] 8×H100 fp8 (batch tier) └──► [llama-70b-fp16-quality-pool] 8×H100 fp16 llama-70b-v3-canary ──► [llama-70b-v3-canary-us-east-1] 1×H100 (small canary)

2. LoRA pool — one base model + many adapters share HBM. Many adapter model_ids route to the same pool:

LORA POOL (1 base + N adapters in HBM, many model_ids → 1 pool) ┌─────────────────────────────────────────────────────────┐ │ llama-70b-lora-pool-us-east-1 (8×H100, fp8) │ │ │ │ HBM layout per node: │ │ ┌─────────────────────────────────────────────────┐ │ │ │ Base model: llama-70b-v2 (~70 GB sharded) │ │ │ ├─────────────────────────────────────────────────┤ │ │ │ Loaded adapters (~100 MB each, ~50 fit easily): │ │ │ │ • exampleco-customer-support-v3 │ │ │ │ • exampleco-product-faq-v1 │ │ │ │ • corp-legal-qa-v2 │ │ │ │ • startup-coding-assistant │ │ │ │ • ... ~46 more │ │ │ ├─────────────────────────────────────────────────┤ │ │ │ KV cache pool (~400 GB total across node) │ │ │ └─────────────────────────────────────────────────┘ │ └─────────────────────────────────────────────────────────┘ ▲ ▲ ▲ ▲ ▲ │ │ │ │ │ exampleco-cs-v3 exampleco-faq corp-legal startup llama-70b-v2 (base, no adapter) │ │ │ │ │ └──────────┴──────────┴──────────┴──────────┘ All these model_ids route to this one pool

3. Colocated multi-model pool — several distinct full models packed into HBM, served by independent processes:

COLOCATED MULTI-MODEL POOL (N distinct full models in HBM, separate processes per model) ┌─────────────────────────────────────────────────────────┐ │ small-models-pool-us-east-1 (1×A100-80G, fp16) │ │ │ │ ┌──────────────────┐ ┌──────────────────┐ │ │ │ Engine process A │ │ Engine process B │ │ │ │ bge-large-en │ │ all-MiniLM-L6 │ │ │ │ (1.3 GB) │ │ (90 MB) │ │ │ │ own KV pool │ │ own KV pool │ │ │ └──────────────────┘ └──────────────────┘ │ │ ┌──────────────────┐ ┌──────────────────┐ │ │ │ Engine process C │ │ Engine process D │ │ │ │ gemma-2b │ │ qwen-1.5b │ │ │ │ (2 GB) │ │ (1.5 GB) │ │ │ │ own KV pool │ │ own KV pool │ │ │ └──────────────────┘ └──────────────────┘ │ │ │ │ (Each model batches independently — no cross-model │ │ batching like LoRA, because they're different │ │ architectures) │ └─────────────────────────────────────────────────────────┘ ▲ ▲ ▲ ▲ │ │ │ │ bge-large all-MiniLM gemma-2b qwen-1.5b │ │ │ │ └──────────────┴──────────────┴──────────────┘ All these model_ids route to this one pool (router dispatches to the right engine process)

The cardinality picture: a model_id can be in many pools (single-model fan-out); a pool can host many model_ids (LoRA + colocated fan-in). The relationship is genuinely many-to-many.

The full schema

Layer 1 — Model catalog (what exists):

-- Models: the catalog. Foundation models + fine-tunes + LoRA adapters are -- all "models" with different metadata. CREATE TABLE models ( model_id VARCHAR(64) PRIMARY KEY, -- "llama-3-70b", "exampleco-customer-support" model_class VARCHAR(16) NOT NULL, -- "foundation" | "fine_tune" | "adapter" architecture VARCHAR(64) NOT NULL, -- "llama", "qwen", "gemma" parameter_count BIGINT NOT NULL, base_model_id VARCHAR(64) NULL -- NULL for foundation; REFERENCES models(model_id), -- set for fine-tunes & adapters license VARCHAR(32), owner_id VARCHAR(64) NOT NULL, visibility VARCHAR(16) NOT NULL, -- "public" | "private" | "org:X" created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), deprecated_at TIMESTAMPTZ NULL ); -- Versions: each model has versions over time CREATE TABLE model_versions ( model_id VARCHAR(64) REFERENCES models, version VARCHAR(32), -- "v2", "v3-canary" weights_uri TEXT NOT NULL, -- "s3://models/llama-3-70b/v2/" config_json JSONB NOT NULL, quantizations TEXT[] NOT NULL, -- {"fp16","fp8","int4-awq"} max_context_len INT NOT NULL, status VARCHAR(16) NOT NULL, -- "available"|"beta"|"deprecated" sha256 CHAR(64) NOT NULL, size_bytes BIGINT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), PRIMARY KEY (model_id, version) ); -- Adapter-specific metadata (LoRA rank, alpha) — only populated when -- model_class='adapter' in models table CREATE TABLE adapter_specs ( model_id VARCHAR(64) PRIMARY KEY REFERENCES models, rank INT NOT NULL, -- LoRA rank (8-64 typical) alpha INT NOT NULL, target_modules TEXT[] NOT NULL, -- which layers it adapts base_model_id VARCHAR(64) NOT NULL REFERENCES models, base_version VARCHAR(32) NOT NULL, FOREIGN KEY (base_model_id, base_version) REFERENCES model_versions(model_id, version) );

Layer 2 — Pool definitions (what's deployed):

-- Pools: physical serving capacity. Type-discriminated. CREATE TABLE pools ( pool_id VARCHAR(64) PRIMARY KEY, pool_type VARCHAR(16) NOT NULL, -- "single"|"lora"|"colocated" -- Single-model: which exact (model, version, quant) this pool serves single_model_id VARCHAR(64) NULL, single_version VARCHAR(32) NULL, -- LoRA: which base model this pool is built on lora_base_id VARCHAR(64) NULL, lora_base_version VARCHAR(32) NULL, -- Common spec quantization VARCHAR(16) NOT NULL, region VARCHAR(32) NOT NULL, tier VARCHAR(16) NOT NULL, -- "shared"|"dedicated"|"batch" gpu_type VARCHAR(32) NOT NULL, -- "H100", "A100-80G", "L40S" tenant_id VARCHAR(64) NULL, -- NULL = shared across tenants -- Deployment params min_replicas INT NOT NULL, max_replicas INT NOT NULL, autoscale_metric VARCHAR(32) NOT NULL, -- "gpu_cache_usage_perc" autoscale_target REAL NOT NULL, -- 0.80 desired_state VARCHAR(16) NOT NULL, -- "active"|"draining"|"deleted" created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), -- Type consistency CHECK ( (pool_type = 'single' AND single_model_id IS NOT NULL AND lora_base_id IS NULL) OR (pool_type = 'lora' AND lora_base_id IS NOT NULL AND single_model_id IS NULL) OR (pool_type = 'colocated' AND single_model_id IS NULL AND lora_base_id IS NULL) ) ); -- Junction: which model_ids each pool actually hosts. -- For single-model: one row pointing to the single model. -- For LoRA: one row for base + one row per loaded adapter. -- For colocated: one row per hosted model. CREATE TABLE pool_hosted_models ( pool_id VARCHAR(64) REFERENCES pools, model_id VARCHAR(64) NOT NULL REFERENCES models, is_adapter BOOLEAN NOT NULL, pinned BOOLEAN NOT NULL DEFAULT FALSE, -- never evict from HBM loaded_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), last_used_at TIMESTAMPTZ, PRIMARY KEY (pool_id, model_id) ); CREATE INDEX idx_hosted_by_model ON pool_hosted_models (model_id); CREATE INDEX idx_hosted_by_pool ON pool_hosted_models (pool_id);

Layer 3 — Routing rules (resolve request to pool):

-- Routing rules: how to map a request (model_id + tenant + region + tier) -- to a ranked list of candidate pools. CREATE TABLE routing_rules ( rule_id VARCHAR(64) PRIMARY KEY, -- Match conditions (NULL = wildcard, matches anything) match_model_id VARCHAR(64) NOT NULL, -- required: which model match_version_pref VARCHAR(32), -- NULL=latest, "v2"=pinned match_tenant_id VARCHAR(64), -- NULL = any tenant match_tenant_tier VARCHAR(16), -- "enterprise"|"paid"|"free" match_region VARCHAR(32), -- NULL = any region match_quantization VARCHAR(16), -- NULL = any quantization -- Result target_pool_id VARCHAR(64) NOT NULL REFERENCES pools, priority INT NOT NULL, -- lower = preferred (fallback order) weight INT NOT NULL DEFAULT 100, -- for weighted routing (canary) active BOOLEAN NOT NULL DEFAULT TRUE, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ); CREATE INDEX idx_routing_lookup ON routing_rules (match_model_id, match_tenant_tier, match_region) WHERE active = TRUE;

Layer 4 — Runtime state (NOT stored; derived):

-- NOT persisted in PostgreSQL. Fetched live from: -- K8s API: pool_id → [pod IPs, ready status, age] -- Engine /metrics: pod_ip → KV occupancy, queue depth, batch size -- Engine /adapters: pod_ip → [currently loaded adapters with timestamps] -- -- The router caches this in memory with ~1-2 second refresh.
Routing resolution — the lookup query

When a request arrives at the router, the resolution is a single JOIN through the junction table:

-- Find ranked candidate pools for a request: -- model_id = "exampleco-customer-support" (a LoRA adapter) -- tenant = "org_exampleco" (enterprise tier) -- region = "us-east-1" SELECT r.target_pool_id, r.priority, r.weight, p.pool_type, p.region, p.tenant_id FROM routing_rules r JOIN pools p ON p.pool_id = r.target_pool_id AND p.desired_state = 'active' JOIN pool_hosted_models h ON h.pool_id = p.pool_id AND h.model_id = 'exampleco-customer-support' -- the JOIN that -- enforces "pool hosts model" WHERE r.active = TRUE AND r.match_model_id = 'exampleco-customer-support' AND (r.match_tenant_id IS NULL OR r.match_tenant_id = 'org_exampleco') AND (r.match_tenant_tier IS NULL OR r.match_tenant_tier = 'enterprise') AND (r.match_region IS NULL OR r.match_region = 'us-east-1') ORDER BY r.priority ASC, r.weight DESC;

Returns a ranked list of candidate pools. Router picks the highest-priority healthy one, then runs Tier-2 (KV-aware) routing inside that pool to pick a specific replica.

Where pool state lives — config vs infrastructure vs runtime

Common confusion: the schema above stores pool definitions, but where does runtime state live — which pods are alive right now, what's their KV occupancy, which adapters are currently loaded? The answer: three state layers, each with a different source of truth. SQL only owns the slow-changing config; runtime state never goes there.

Layer What it contains Source of truth Update rate How router consumes
Config (definitions) Pool definitions, routing rules, model registry PostgreSQL (or K8s CRDs) Minutes to hours Pulled with cache, refreshed every 30 s
Infrastructure (existence + health) Which pods exist, their IPs, ready / not-ready K8s API (etcd backend) Sub-second K8s watch stream (push-based)
Runtime (load + capacity) KV occupancy, queue depth, loaded adapters, batch size Each engine's /metrics endpoint Real-time (per iteration) HTTP polling every 1−5 s

Pool definition lives in SQL: “pool_X should exist with these specs.” Pool state lives in K8s + engines: “pool_X currently has 3 pods at these IPs, KV occupancy is 72%, 65%, 80%.” These are completely separate stores. The general principle: persist what doesn't change often; query directly what does.

Why runtime state is NOT in SQL
  1. Update rate. KV occupancy changes every decode iteration (~6 ms). Writing 30 GPUs × 200 metrics every 6 ms to SQL would crater the DB.
  2. No persistence value. KV occupancy from 5 seconds ago is useless for routing. Nothing to archive.
  3. Source of truth is elsewhere. K8s API already knows which pods exist. Engines already know their KV state. SQL would be a stale mirror at best.
  4. Read latency. Per-request SQL lookups would add 1+ ms to every request. Direct in-memory metrics is sub-microsecond.
How the router assembles its view

The router runs three concurrent state-management loops, combining all layers into a single in-memory view used for routing decisions:

Router pod (in-memory state): ┌─────────────────────────────────────────────────────────────────┐ │ │ │ CONFIG (refreshed every 30s) │ │ ┌─────────────────────────────────────────────────────────┐ │ │ │ pool_definitions ← Postgres / K8s CRDs │ │ │ │ routing_rules ← Postgres / K8s CRDs │ │ │ │ pool_hosted_models ← Postgres / K8s CRDs │ │ │ └─────────────────────────────────────────────────────────┘ │ │ │ │ INFRASTRUCTURE (K8s watch stream, real-time push) │ │ ┌─────────────────────────────────────────────────────────┐ │ │ │ pool_X.pods = [ │ │ │ │ { ip: 10.0.0.5, ready: true, age: 5m }, │ │ │ │ { ip: 10.0.0.6, ready: true, age: 2h }, │ │ │ │ { ip: 10.0.0.7, ready: false, age: 30s } │ │ │ │ ] │ │ │ └─────────────────────────────────────────────────────────┘ │ │ │ │ RUNTIME (polled every 1-5s from each pod's /metrics) │ │ ┌─────────────────────────────────────────────────────────┐ │ │ │ pod_metrics = { │ │ │ │ "10.0.0.5": { kv: 72%, queue: 0, adapters: [...] } │ │ │ │ "10.0.0.6": { kv: 88%, queue: 3, adapters: [...] } │ │ │ │ "10.0.0.7": (not ready, no metrics yet) │ │ │ │ } │ │ │ └─────────────────────────────────────────────────────────┘ │ │ │ │ PER-REQUEST: combine all three for routing decision (~100 ns) │ └─────────────────────────────────────────────────────────────────┘

Each layer has its own update mechanism:

  • Config loop: polls Postgres every 30 s, OR subscribes to Redis pub/sub for invalidations on writes.
  • K8s watch loop: opens long-lived watch stream to K8s API server; gets push notifications on every pod change (created, ready, deleted, IP changed).
  • Metrics poll loop: for each known pod, scrapes its /metrics endpoint every 1−5 s.
The K8s watch gives sub-second pod awareness

Kubernetes API supports Watch — a long-lived HTTP/2 stream that pushes events when matching resources change. This is the standard way controllers (and routers) stay in sync with K8s state without polling.

# Pseudocode for router's K8s watch def watch_pool_pods(pool_label): for event in k8s_client.watch( kind="Pod", label_selector=f"pool={pool_label}" ): if event.type == "ADDED" and event.object.status == "Running": self.add_pod(event.object.pod_ip) self.start_metrics_polling(event.object.pod_ip) elif event.type == "DELETED" or event.object.status != "Running": self.remove_pod(event.object.pod_ip) self.stop_metrics_polling(event.object.pod_ip)

When a new pod comes up, the router knows within ~1 s (the time for kubelet to mark it Ready and K8s API to fan out the event). No SQL involvement.

Metrics polling gives load awareness

For each known pod, the router pulls /metrics (Prometheus-format text):

GET http://10.0.0.5:8000/metrics # HELP vllm:gpu_cache_usage_perc GPU KV cache usage # TYPE vllm:gpu_cache_usage_perc gauge vllm:gpu_cache_usage_perc 0.72 vllm:num_requests_running 12 vllm:num_requests_waiting 0 vllm:num_requests_swapped 0 ...

Router parses, updates in-memory state, uses for per-request routing. 1−5 s polling is the sweet spot: fast enough to react to load, slow enough to not hammer engines with HTTP traffic.

Flow examples — new pod joins, pod becomes unhealthy

Example A: new pod joins the pool

T+0: Autoscaler updates K8s Deployment replica count from 3 → 4 T+0.5: K8s scheduler binds new pod to a node K8s API: pod_4 created, status=Pending Router watch event: pod_4 added, ignore (not ready yet) T+5: Pod container starts; engine begins loading weights T+45: Weights loaded; engine /health returns 200 kubelet marks pod_4 READY K8s API: pod_4 status=Running, ready=true Router watch event: pod_4 is ready → add to routing T+45.1: Router opens /metrics polling to pod_4 First poll: KV=0%, queue=0 T+45.1: Router includes pod_4 in routing candidates T+46: Router polls again, KV still ~0%, routes some traffic to pod_4 T+50: pod_4 starts seeing requests, KV rising SQL was never queried in this flow. Pool definition in SQL said "pool_X should have 3-10 replicas"; everything else came from K8s + engine.

Example B: pod becomes unhealthy

T+0: Pod_2's engine crashes (OOM) T+0.5: Process exits; container restarts Engine /health starts returning 503 (warming up again) T+1: kubelet probe fails: pod_2 not ready K8s API: pod_2 status=Running, ready=false Router watch event: pod_2 NOT ready → remove from routing T+1.1: Router stops sending traffic to pod_2 T+45: Engine reloads weights, /health returns 200 kubelet probe succeeds: pod_2 ready Router watch event: pod_2 ready → re-add to routing Recovery bounded by detection (kubelet probe interval) + replacement (weight reload time). SQL is irrelevant.
The common mistake: mirroring runtime state to SQL

People sometimes try to put runtime state in SQL for "observability" or "easier queries." This is almost always wrong:

-- DON'T DO THIS CREATE TABLE pod_runtime_state ( pod_id VARCHAR(64), kv_occupancy REAL, queue_depth INT, updated_at TIMESTAMPTZ ); -- Update every 5 seconds from every pod...
  • Write amplification: 100 pods × 12 updates/min × many metrics = thousands of writes/sec to one table. SQL throughput chokes.
  • Stale-by-design: the row says KV was 72% at the last update, but it's 5 s old. Router would route on stale data.
  • Wrong source of truth: K8s + engine endpoints already have this; mirroring adds lag and bugs.
  • No persistence value: nobody cares about KV occupancy from yesterday at this exact second.

The right place for historical metrics is Prometheus (or VictoriaMetrics, Mimir), which is purpose-built for time series. Prometheus scrapes the same /metrics endpoints every 15−30 s and stores them in a TSDB for dashboards. The router doesn't query Prometheus — too slow, too historical. Router goes direct to the engine.

┌──────────────────────────┐ │ Each engine's /metrics │ └──────────┬───────────────┘ │ ┌──────────────┼──────────────┐ │ │ │ ▼ ▼ ▼ ┌────────┐ ┌────────┐ ┌──────────┐ │ Router │ │ Router │ │Prometheus│ │ pod A │ │ pod B │ │ (TSDB) │ │ 1-5s │ │ 1-5s │ │ 15-30s │ └────────┘ └────────┘ └────┬─────┘ │ ▼ ┌────────┐ │Grafana │ │(humans)│ └────────┘ Same source endpoint, three consumers, very different refresh rates.
What lives where — the unambiguous reference
QuestionAnswered by
Does pool_X exist as a deployment intent?Postgres (or K8s CRDs)
Which model does pool_X serve?Postgres / K8s CRDs
Which pods belong to pool_X right now?K8s API (live)
Is pod_5 ready to take traffic?K8s API (live, derived from kubelet probes)
What's pod_5's IP?K8s API (live)
Is pod_5 healthy from the engine's perspective?engine /health
What's pod_5's current KV occupancy?engine /metrics
How many requests in pod_5's queue?engine /metrics
Which LoRA adapters are loaded on pod_5?engine /adapters (or /metrics)
What was pod_5's KV occupancy at 3pm yesterday?Prometheus

If you find yourself wanting to put rows 3−9 in Postgres, stop. That's the design mistake.

Adapter lifecycle in LoRA pools (the dynamic case)

LoRA pools have a runtime dimension single-model and colocated pools don't: which adapters are loaded in HBM right now. The pool spec lists *registered* adapters; the engine decides which subset is actually resident.

$$\text{adapter capacity} \;=\; \frac{\text{HBM} - \text{base weights} - \text{KV pool}}{\text{adapter size}}$$

For Llama-70B-FP8 base (~70 GB) on 8×H100 (640 GB total), KV pool ~400 GB → ~150 GB headroom → ~1000 LoRA adapters at 150 MB each. So you don't usually need to evict; the constraint is "registered adapters" not HBM capacity.

Two-tier loading policy:

  • Pinned adapters (pool_hosted_models.pinned=TRUE): always loaded in HBM, never evicted. Top-20 by traffic.
  • On-demand adapters: loaded on first request (~1-2 s from S3), evicted LRU after long idle (~10 min).

This is the hot/cold problem from item 4.2 but two orders of magnitude cheaper because adapters are ~100 MB instead of ~200 GB. Cold-load is seconds instead of minutes; eviction is sub-millisecond.

The end-to-end flow — publish a model, deploy, serve
[REGISTER] 1. ML pipeline trains model → uploads weights to s3://models/qwen-72b/v1/ 2. Pipeline calls registry API: POST /api/models { model_id: "qwen-72b", model_class: "foundation", architecture: "qwen", parameter_count: 72_000_000_000 } 3. Pipeline calls version API: POST /api/models/qwen-72b/versions { version: "v1", weights_uri: "s3://...", quantizations: ["fp16","fp8"], status: "beta" } 4. Postgres: INSERT into models, model_versions. [REGISTER LORA] 1. Customer trains LoRA on Llama-70b-v2 2. SDK uploads adapter → s3://adapters/exampleco/customer-support/v3/ 3. SDK calls: POST /api/models { model_id: "exampleco-customer-support", model_class: "adapter", base_model_id: "llama-3-70b", owner_id: "org_exampleco", visibility: "org:org_exampleco" } POST /api/adapter_specs { model_id: "exampleco-customer-support", rank: 32, alpha: 64, base_model_id: "llama-3-70b", base_version: "v2" } [DEPLOY POOL] 1. Platform engineer authors a pool CRD or calls deployment API: { pool_id: "llama-70b-lora-pool-us-east-1", pool_type: "lora", lora_base_id: "llama-3-70b", lora_base_version: "v2", hosted_adapters: ["exampleco-cs-v3", "exampleco-faq", "corp-legal", ...], region: "us-east-1", min_replicas: 3, max_replicas: 10 } 2. Controller: - INSERT pools row - INSERT pool_hosted_models rows (1 for base + N for adapters) - Create K8s Deployment with vLLM pods: vllm --model s3://models/llama-3-70b/v2/ --enable-lora --lora-modules exampleco-cs-v3=s3://... exampleco-faq=s3://... ... 3. Pods cold-load (~30-60 s) → engine reports adapters loaded → READY 4. Router refreshes from Postgres → pool is now serving 5. Controller (or admin) creates routing_rules entry pointing to this pool [SERVE] 1. Client request: model_id="exampleco-customer-support" from org_exampleco in us-east-1 2. Gateway: validate JWT → extract org_id="org_exampleco", tier="enterprise" 3. Router resolution: - Tier-1: routing_rules JOIN pool_hosted_models → llama-70b-lora-pool-us-east-1 - Tier-2: KV-aware pick within pool → replica_5 4. Engine on replica_5: dispatches "exampleco-customer-support" adapter 5. Response streams back via router → gateway → client
Cardinality summary — the full many-to-many picture
model_id (logical) pool (physical capacity) ───────────────────────────────────────────────────────────────────────── llama-3-70b ────────────────┬─────────► llama-70b-v2-shared-us-east-1 ├─────────► llama-70b-v2-shared-us-west-2 ├─────────► llama-70b-v2-shared-eu-west-1 ├─────────► llama-70b-v2-dedicated-exampleco ├─────────► llama-70b-v2-batch-pool └─────────► llama-70b-lora-pool-us-east-1 (as the base model) exampleco-customer-support ──────┬─────────► llama-70b-lora-pool-us-east-1 (LoRA adapter) └─────────► llama-70b-v2-dedicated-exampleco (exampleco's adapters loaded here too) corp-legal-qa ──────────────────────► llama-70b-lora-pool-us-east-1 (LoRA adapter) (only in shared LoRA pool) bge-large-en-v1.5 ──────────────────► small-embedders-pool-us-east-1 (small embedder) (sharing pool with 3 others) gemma-2b ───────────────────────────► small-llms-pool-us-east-1 (small LLM) (sharing pool with 3 others) ┌───────────────────────────────────────────────────────────────────┐ │ Cardinality: │ │ model_id → pool: many (foundation models fan out across pools) │ │ pool → model_id: many (LoRA pools + colocated pools fan in) │ │ Relationship is many-to-many on both axes. │ │ Junction table `pool_hosted_models` is what makes it work. │ └───────────────────────────────────────────────────────────────────┘
What to say in an interview

A senior answer makes the abstraction explicit: “Treat model_id as a logical identifier and pool as physical capacity, joined by a routing rules table. The relationship is many-to-many on both axes — a foundation model fans out across regions/tiers/quants (many pools per model_id); LoRA pools fan in (many adapter model_ids per pool); colocated pools fan in (many small model_ids per pool). The schema needs three layers: a model catalog in Postgres (`models`, `model_versions`, `adapter_specs`), pool definitions with a type discriminator (`pools.pool_type` = single|lora|colocated), and a junction table `pool_hosted_models` that tells the router which model_ids each pool can serve. Routing rules add the tenant/region/tier match conditions on top, with priority + weight for canary rollouts. The router's lookup is one JOIN through `pool_hosted_models` to filter candidates that actually host the requested model. The conceptual analogy is DNS hostname → IP or service mesh service → endpoint — logical-to-physical resolution with rich policy.”

The trap to avoid: assuming "one pool = one model" because it's simpler. That assumption breaks the moment you add LoRA serving or colocated small models, both of which are standard production patterns. The schema needs to support all three pool types from day one — retrofitting the junction table later is painful.

5
SOTA KV technology (Jun. 2026)
Section 5 · frontier architectures for long-context inference
⌃

Chapters 1 and 2 explain how serving systems move, allocate, quantize, and reuse KV cache. This chapter is the modeling-side complement: how frontier architectures reduce the amount of KV that needs to exist in the first place. The unifying idea is simple: uniform full-resolution attention memory is waste. A production model should not pay the same HBM price for every old token, every layer, and every head.

As of June 2026, the SOTA pattern is no longer "make full attention faster." It is "make attention non-uniform" along five axes: representation, distance, token selection, memory resolution, and time. These axes overlap in real models, but separating them is the cleanest way to reason about the serving economics.

5.1 · Foundation

The long-context memory wall

Importance: HIGH
Problem
Full attention makes decode cost grow with context length.
Key phrase
Uniform full-resolution memory — every cached token remains equally expensive to store and scan.
What it does & why it matters

In standard decoder-only transformers, every generated token runs a full forward pass through the model. During decode, that forward pass has two memory bills:

  1. Weights: the model weights are resident in HBM, but small-batch decode still streams them from HBM into the compute units every token.
  2. KV cache: attention reads the historical K/V vectors for the session. This bill grows linearly with context length.

The toy arithmetic is useful because it explains the product behavior. For Llama 3 70B at FP16, the raw KV per token is:

$$2 \times 80 \times 8 \times 128 \times 2\text{ bytes} = 327{,}680\text{ bytes} \approx 320\text{ KiB}$$

That is K and V, across 80 layers, 8 KV heads, 128 dimensions per head, at 2 bytes per FP16 value. At 1M tokens, the raw KV cache is about 330 GB before allocator overhead. Add roughly 141 GB of FP16 weights, and a single decode step may require on the order of 470 GB of HBM traffic.

$$\frac{8{,}000\text{ GB/s}}{470\text{ GB/token}} \approx 17\text{ tokens/s}$$

This is a structural ceiling, not a benchmark promise. It assumes perfect bandwidth utilization, no communication overhead, and instantaneous math. Real systems do worse unless the architecture reduces what has to be read.

The architectural lesson

Full attention treats a recent instruction, a useful retrieved fact, boilerplate from 700K tokens ago, and an irrelevant old sentence as the same class of memory: full-resolution KV, always eligible to be read. That is the waste. A scalable long context model needs a memory hierarchy:

  • Hot memory: recent and control-critical tokens, fast and high fidelity.
  • Warm memory: older useful context, compressed or selectively retrieved.
  • Cold memory: old context represented as summaries, indexes, or recurrent state.
  • Discarded memory: tokens that stop being useful enough to justify HBM.

The rest of this chapter is a taxonomy of how frontier models build that hierarchy.

References
5.2 · Representation sparsity

Multi-Query, Grouped-Query & Multi-head Latent Attention

Importance: HIGH
Question
What can we stop storing per token?
Examples
Multi-Query Attention (MQA) Grouped-Query Attention (GQA) Multi-head Latent Attention (MLA) DeepSeek-V2/V3/V4 Kimi K2.6
What it does & why it matters

Multi-head attention stores K and V separately for each attention head. But those heads are correlated views of the same token. Representation sparsity reduces the number or dimensionality of the cached views. The three names to know are Multi-Query Attention (MQA), Grouped-Query Attention (GQA), and Multi-head Latent Attention (MLA).

  • Multi-Query Attention (MQA): all query heads share one K/V head. Very small cache, sometimes quality cost.
  • Grouped-Query Attention (GQA): groups of query heads share a K/V head. This is the practical compromise used by Llama, Gemma, Qwen, and many serving-oriented dense models.
  • Multi-head Latent Attention (MLA): cache a low-dimensional latent representation, then reconstruct the K/V information needed by each head on demand.

GQA changes the KV formula by reducing kv_heads. Llama 3 70B has 64 query heads but only 8 KV heads, so its KV cache is 8× smaller than full MHA with 64 KV heads. That is why the earlier formula uses 8, not 64.

MLA: latent object, not every camera angle

MLA takes the same idea further. Instead of storing per-head K/V tensors, DeepSeek stores a compact latent vector per token and reconstructs the useful attention representation inside the layer. This trades extra compute for much lower HBM traffic, which is exactly the right trade in memory-bound decode.

DeepSeek-V2 reported 93.3% KV-cache reduction versus DeepSeek 67B while maintaining or improving capability. The important point is not just the headline number; it is that MLA changes the roofline position of decode. Less HBM per token means the same GPU can sustain more sessions or longer contexts before KV becomes the admission bottleneck.

By 2026, MLA is no longer a one-off DeepSeek trick. Kimi K2.6 lists MLA as its attention mechanism, with 1T total parameters, 32B activated parameters, and 256K context. That is a clear signal that frontier open-weight models are adopting latent KV as a standard long-context primitive.

Serving implication

GQA and MLA are architecture-level versions of KV quantization: they reduce the number of bytes stored and read per cached token. The difference is that quantization compresses the same tensors after training, while GQA/MLA change what the model learns to cache. For long-context serving, this affects three systems directly:

  • Admission control: lower KV/token means more concurrent sessions under the same HBM budget.
  • P/D disaggregation: less KV to transfer from prefill pool to decode pool.
  • Hierarchical KV: less data to spill, load, hash, and index across tiers.
References
5.3 · Distance sparsity

Local/global attention: most layers look nearby

Importance: MEDIUM-HIGH
Question
What can we stop attending to at every layer?
Examples
Sliding-window attention Gemma 3 gpt-oss Mistral-style windows
What it does & why it matters

Language is not uniformly long-range. Syntax, local references, code indentation, and short idioms are mostly nearby. Long-range dependencies matter, but not every layer needs to serve them.

Distance sparsity uses many local layers plus occasional global layers. A local layer only attends to a recent sliding window, so its KV cache can be bounded by window size instead of total context length. Global layers preserve whole-document reach, but they are fewer.

Layer pattern: local local local local local global 1K 1K 1K 1K 1K full Only the global layer pays the full-context KV bill.
Concrete designs

Gemma 3 is the cleanest published example. Google states that Gemma 3 interleaves 5 local layers for every 1 global layer, with local layers using a 1024-token span, specifically to reduce KV-cache growth at long context. OpenAI's gpt-oss model card similarly describes alternating dense and locally banded sparse attention, with grouped multi-query attention for memory efficiency.

This is not a hack around attention; it is a statement about where language processing happens. Early and middle layers can resolve most local structure from a short window. Fewer layers need expensive global memory.

Serving implication

Local/global attention reduces the slope of KV growth. If only a fraction of layers are global, then only that fraction needs unbounded KV cache. Local-layer KV can be evicted as the window slides. The result is a model that still supports long context but has a better HBM curve than full attention.

The caveat: local/global attention is predictable but blunt. It does not know which old tokens are important; it simply gives full reach to selected layers. That is why it composes naturally with token sparsity and compressed memory.

References
5.4 · Token sparsity

Learned sparse retrieval over history

Importance: HIGH
Question
What old tokens can we stop reading for this query?
Examples
Native Sparse Attention (NSA) DeepSeek Sparse Attention (DSA) learned indexer block selection
What it does & why it matters

If someone asks a question about a book, you do not reread the book from page one. You use an index, find the relevant pages, and read those carefully. Token sparsity gives attention the same shape.

The hard version is not "drop random tokens." It is: for each query, cheaply find the relevant historical blocks, then run exact attention on those blocks plus a local window. That requires the selector to be learned, hardware-aligned, and trained with the model rather than bolted on afterward.

Native Sparse Attention: compressed, selected, local

DeepSeek's Native Sparse Attention (NSA) is the canonical design. It has three paths:

  1. Compressed coarse history: blocks of old tokens are summarized so the query can scan the whole context cheaply.
  2. Selected fine-grained blocks: the model chooses important blocks and attends to their exact K/V.
  3. Sliding local window: recent tokens remain high fidelity because local continuity is always important.

The important systems detail is blockwise selection. Selecting scattered individual tokens destroys memory coalescing and Tensor Core utilization. Selecting contiguous blocks keeps the kernel closer to FlashAttention-style access patterns. NSA reports up to 11.6× expected decoding speedup at 64K context in its setup because memory access volume falls sharply.

DeepSeek Sparse Attention: sparse attention in a frontier model

DeepSeek-V3.2-Exp took the idea into a production-scale model as DeepSeek Sparse Attention (DSA). The model card frames V3.2-Exp as an intermediate step toward a next-generation architecture, preserving output quality while improving long-context training and inference efficiency.

The serving implication is direct: sparse attention changes decode from "read the whole prefix" to "read the useful selected blocks." That is the difference between nominal 1M context and economically usable 1M context.

References
5.5 · Resolution sparsity

Compressed deep memory: old context gets cheaper

Importance: HIGH
Question
What old memory can stay useful at lower resolution?
Examples
Compressed Sparse Attention (CSA) Heavily Compressed Attention (HCA) DeepSeek-V4 million-token context
What it does & why it matters

Human memory is multi-resolution. This morning is high fidelity; last month is episodic; last year is mostly compressed themes with a few sharp details. Long-context attention is moving the same way.

Resolution sparsity keeps old context available, but not necessarily as full-resolution K/V for every token. The model maintains a compressed global view and selectively recovers detail when needed. This is different from a simple sliding window: old context is not gone, it is represented more cheaply.

DeepSeek-V4: Compressed Sparse + Heavily Compressed Attention

DeepSeek-V4 is the clearest 2026 example. The public report describes a hybrid attention architecture combining Compressed Sparse Attention (CSA) and Heavily Compressed Attention (HCA). In the 1M-token setting, DeepSeek reports that V4-Pro requires only 27% of the single-token inference FLOPs and 10% of the KV cache compared with DeepSeek-V3.2.

That headline is the product significance: 1M context moves from a special-case surcharge tier toward a default capability. It is not because HBM became infinite. It is because the model stopped storing and scanning the deep past at uniform resolution.

Serving implication

Resolution sparsity changes the shape of the decode curve. In a full-attention model, each new token gets slower as context grows because the KV scan grows. In a compressed-memory model, the far-history component grows much more slowly, so decode throughput can remain flatter across 4K, 100K, and 1M contexts.

The tradeoff is complexity. A serving engine now needs kernels and cache managers that understand multiple memory forms: local full-resolution KV, selected sparse KV, compressed global memory, and possibly indexer state. PagedAttention-style block tables are still useful, but the "block" is no longer always the same object.

References
5.6 · Temporal sparsity

State instead of cache

Importance: MEDIUM-HIGH
Question
What if some layers stop storing token history at all?
Examples
Mamba / State Space Models (SSMs) Gated DeltaNet Kimi Delta Attention (KDA) Qwen3-Next Kimi Linear
What it does & why it matters

The boldest KV reduction is to avoid per-token KV for some layers entirely. Linear attention, state-space models, and gated recurrence maintain a running state:

$$\text{state}_{t+1} = f(\text{state}_t, x_t)$$

If this worked perfectly for all tasks, memory would be flat with context length. There would be no 1M-token KV cache to scan. The reason full attention survived so long is that fixed-state memory is hard: compressing arbitrary history into a fixed vector can lose details that exact attention would recover.

The 2026 compromise: hybrid layers

The frontier pattern is not pure recurrence. It is a hybrid: many recurrent or linear-attention layers for cheap temporal memory, plus occasional attention layers for exact retrieval and alignment.

Qwen3-Next-80B-A3B uses a repeated layout of 3 Gated DeltaNet layers followed by 1 Gated Attention layer. The model card reports 80B total parameters, 3B activated parameters, native 262K context, and validation up to about 1M tokens with YaRN. It also reports 10× inference throughput over Qwen3-32B-Base for contexts over 32K.

Qwen3.5-397B-A17B continues the same broad direction at larger scale: Gated Delta Networks plus sparse MoE, with 397B total parameters, 17B activated, and a 3:1 DeltaNet-to-attention layout.

Kimi Linear is the other key signal. It combines Kimi Delta Attention (KDA) with MLA and reports up to 75% KV-cache reduction and up to 6× decoding throughput at 1M context, while claiming to outperform full-attention baselines under fair comparisons.

Serving implication

Temporal sparsity changes the state object the engine must manage. Instead of storing a growing KV tensor for every layer, hybrid models store:

  • Attention-layer KV: still paged, transferred, quantized, and evicted like normal KV.
  • Recurrent state: fixed-size per sequence, cheap to carry across decode steps.
  • MoE routing state: expert dispatch remains the dominant complexity for sparse FFN layers.

The interview-level summary: full attention is exact but grows with context; recurrence is cheap but compressive; frontier models increasingly use both.

References

Interview summary

A senior answer should not list acronyms. It should name the memory bill: decode is HBM-bound, and full attention makes the KV bill grow with context. Then map each architecture to the bill it reduces: GQA/MLA reduce bytes per token; local/global attention reduces which layers pay full context; NSA/DSA read selected blocks; CSA/HCA lower the resolution of deep history; DeltaNet/KDA replace part of the cache with fixed-size state.

The clean thesis: SOTA long-context models are becoming memory-hierarchy systems. Serving stacks still need PagedAttention, quantization, prefix routing, and disaggregation, but the biggest gains now come from giving those stacks less uniform KV to move.