Quick reference for the acronyms you'll hear in every inference conversation. Read this once, then come back when something below feels unfamiliar.
TTFT + (num_output_tokens - 1) × ITL.all-reduce, all-gather, broadcast in distributed training and inference.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.NVFP4 (NVIDIA Blackwell) and MXFP4 (Microscaling open standard). Native on B200. Another ~2× over FP8 at some quality cost.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.
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.
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.
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:
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.
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.
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.
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.
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.
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:
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.
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.
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:
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.
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:
Resource cost: scheduler CPU work, sub-ms per iteration. ↳ See item 2.1 Continuous batching.
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.
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:
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.
The OSS inference world has split into two camps based on how the model graph is executed.
| 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.
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.
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.
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.
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.
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:
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.
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:
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 sharpest production systems use both, in layers:
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.
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:
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:
hash → node map.
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.
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.)
[B1, B2, B3, B4, B5] — served by replica R1[B1, B2, B3, B6, B7] — served by replica R2[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.
Each tree node holds a small struct:
Two flavors of hashing appear here and they do different jobs:
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.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."
Now suppose Request D arrives with blocks [B1, B2, B3, B6, B8, B9].
The router walks the tree:
B1, find it in root's children → descend to B1 node. Note replicas: {R1, R2, R3}.B2, find it → descend to B2. Replicas still {R1, R2, R3}.B3, find it → descend to B3. Replicas narrow to {R1, R2}.B6, find it → descend to B6. Replicas narrow to {R2}.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).
replica_set to include this new replica.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.
"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.
replica_set BitMap actually addresses machinesThe 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:
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.
The bitmap dimension is fixed at engine startup, so it does cap the number of replicas the router can address. Production choices:
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.
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:
replica_set bitmaps as it caches blocks and reports the events upstream.When a replica leaves (planned drain or crash):
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.
Prefill and decode have opposite hardware appetites:
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.
Three components make up a disaggregated cluster:
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.
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:
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.
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:
Two implementation details matter:
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.
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:
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) | ~0 | essentially free |
| L2 (host DRAM, PCIe 5.0) | 5–10 | 4−8× faster |
| L4 (RDMA, 25 GB/s) | ~15 | ~2.5× faster |
| L3 (local NVMe, 10 GB/s) | ~30 | ~1.3× faster (marginal) |
| Recompute (prefill) | ~40 | baseline |
Two non-obvious consequences fall out of this table:
Prompt is 4,096 tokens. Two candidate nodes:
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.
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:
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.
Three layers stack to make this work:
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.A typical KV transfer in a disaggregated cluster:
{p_1, p_2, ..., p_n} (PagedAttention block pointers).{d_1, ..., d_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.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).
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:
all-reduce sums partial results back to the full activation. Per-layer, per-token. Extremely bandwidth-hungry.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.
The rule that gets tattooed on every inference engineer's forehead:
| Parallelism | Bandwidth needed | Goes on |
|---|---|---|
| TP | Highest (all-reduce every layer) | Inside one NVLink domain (1 node, ≤8 GPUs) |
| EP | High (all-to-all every MoE layer) | Inside one node, sometimes 2 if linked by NVLink (NVL72) |
| PP | Modest (point-to-point per micro-batch) | Across nodes via InfiniBand |
| DP (data parallelism / replication) | None during inference | Across 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:
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.torch.distributed + CUDA streams handle this; engines like Megatron also do sequence-parallel comm fusion.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:
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.
Four pieces stack to make multi-region inference work:
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.
Inference services serve many customers from one cluster. The two ends of the spectrum:
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:
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:
(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).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.
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.
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:
Continuous batching (also called in-flight batching, iteration-level scheduling) makes scheduling decisions at every decode step, not every batch:
<eos> token leave the batch immediately.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.
The engine runs a scheduler loop that, on every decode iteration:
<eos>, hit max-tokens, or got cancelled).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.Two important real-world details:
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).
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:
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:
Three pieces work together:
append_token when the current last block is full.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.
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.
The classic softmax formula is:
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:
The same trick lets you accumulate the output $\mathbf{O}$ tile by tile.
So the FA-1 inner loop, per query tile, is:
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.
H100 (Hopper) introduced:
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.
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.
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:
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.
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.
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.
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.
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.
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.
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:
--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.
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:
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.
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.
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.
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.
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.
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:
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:
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.
Two design choices:
max_num_batched_tokens.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.
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:
Why this matters for inference:
Serving MoE well is a different engineering problem than serving dense. Open-weight MoE deployments made those differences highly visible.
Split the N experts across G GPUs. Each GPU holds N/G experts and their weights. Per MoE layer, the forward pass is:
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).
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.
Train-time techniques (auxiliary load-balance losses) help but don't fully fix runtime imbalance. Inference techniques:
DeepSeek-V3 has its own twists worth knowing:
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.
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.
Any workload-feedback loop is fully described by three properties:
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.
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 |
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.
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.
Three sentences worth being able to deliver:
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:
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 technical contribution that makes BYOC work as a product is the control-plane / data-plane separation. Each half has a clear job:
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.
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:
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.
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?
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”:
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.
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.
With a regional shared FS in place, the architecture collapses into something elegant:
/weights/llama-3-70b/v2.1/, /weights/deepseek-v3/v3.2/, ...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.
A careful reader hits two questions at this point:
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:
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.
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:
s5cmd, aria2),
saturates its NIC. Maybe 5 minutes for 200 GB. Off the hot path.
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:
CreateVolume from snapshot
+ attach + mount, ~20–40 s of K8s and cloud-API round-trips before any
data moves.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.
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):
Three things stand out from the breakdown:
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.
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.
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:
Tier 0 — you must have these or you don't have a serving system:
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.”
Tier 1 — the buffer that buys time for Tier 0:
Tier 2 — helpful but not always justified:
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.
Terminology, since this trips people up:
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.
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:
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.
Tier 0 — required for any multi-model platform:
Tier 1 — what makes the cold path tolerable:
{"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:
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.
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:
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.
A clean mental model: each layer has its own design choices, and they compose independently. Total rollout time is roughly:
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.
Layer 1 · Weight distribution (priority order):
nimble, BitTorrent-derived libraries, custom implementations at
hyperscalers. Implementation cost is non-trivial — only worth it at fleet
sizes >50.
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:
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.
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-schedulerK8s control plane |
Binds the pod to a new node based on GPU resource requests, affinity rules, taints. |
| 2 | kubelet + CSI driverNode-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 entrypointvLLM 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 | ModelLoaderEngine 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 libraryThird-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_stateOne 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_cacheEngine 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 | CUDAGraphRunnerEngine 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 + RouterK8s + 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."
Three properties of mmap that make this work at scale:
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:
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).
200-replica deploy of Llama-70B v2 in one region, weights pre-staged the night before via DaemonSet:
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."
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 | 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).
Tier 0 — you must have these:
Tier 1 — what makes recovery fast:
Tier 2 — advanced patterns:
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):
| Metric | Why 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_second | Tokens/sec delivered within SLO. Not raw throughput — the honest serving capacity. |
streaming_disconnect_rate | Sessions 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_perc | KV-cache occupancy — the leading indicator for saturation. Scale on this, not CPU. |
vllm:num_requests_waiting | Admission queue depth — if this is non-zero you're already saturated. |
vllm:num_requests_running | Active batch size — capped at max_num_seqs; sanity check. |
vllm:num_requests_swapped | Requests preempted to host RAM — non-zero means you're in pain (KV exhaustion). |
vllm:prefix_cache_hits_total / queries_total | Prefix-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):
| Metric | Why it matters |
|---|---|
DCGM_FI_DEV_FB_USED / FB_TOTAL | HBM used vs capacity (hardware view). Cross-check with engine's KV occupancy. |
DCGM_FI_DEV_MEM_COPY_UTIL | HBM bandwidth utilization — the binding resource for decode. ~100% on healthy decode-bound serving. |
DCGM_FI_DEV_GPU_UTIL | SM utilization. Less useful for inference (decode is memory-bound) but tracked for sanity. |
DCGM_FI_DEV_POWER_USAGE | Watts — 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):
| Metric | Why it matters |
|---|---|
request_prompt_tokens distribution (p50/p95) | Prompt-length distribution. Shifts here change your KV math overnight. |
request_generation_tokens distribution | Output-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 model | Lower than expected = workload shift, eviction problem, or routing misconfiguration. |
Tier 5 — infra / operational (catches platform-level issues):
| Metric | Why it matters |
|---|---|
pod_restart_count | Engine crashes / OOMs. Should be ~0; any spike = investigate. |
cold_start_duration_seconds | Time from pod scheduled to first request served. Gates your scale-up speed. |
weight_load_duration_seconds | Sub-component of cold start — isolates FS/distribution problems. |
autoscaler_decisions_total | Scale-up and scale-down events. Frequent oscillation = thrashing, fix your hysteresis. |
cost_per_million_tokens | Derived metric (GPU-hours / tokens served). The business-level health check. |
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.
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.
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:
llama-70b-v2 in us-east-1, us-west-2, eu-west-1, plus a canary v3 pool, plus dedicated pools per tenant.model_ids share one pool (a base model + many adapters all loaded in HBM together).
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.
1. Single-model pool — one (model, version, quantization) per pool. Fans out across regions/tiers/quants:
2. LoRA pool — one base model + many adapters share HBM. Many adapter model_ids route to the same pool:
3. Colocated multi-model pool — several distinct full models packed into HBM, served by independent processes:
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.
Layer 1 — Model catalog (what exists):
Layer 2 — Pool definitions (what's deployed):
Layer 3 — Routing rules (resolve request to pool):
Layer 4 — Runtime state (NOT stored; derived):
When a request arrives at the router, the resolution is a single JOIN through the junction table:
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.
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.
The router runs three concurrent state-management loops, combining all layers into a single in-memory view used for routing decisions:
Each layer has its own update mechanism:
/metrics endpoint every 1−5 s.
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.
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.
For each known pod, the router pulls /metrics (Prometheus-format text):
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.
Example A: new pod joins the pool
Example B: pod becomes unhealthy
People sometimes try to put runtime state in SQL for "observability" or "easier queries." This is almost always wrong:
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.
| Question | Answered 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.
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.
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:
pool_hosted_models.pinned=TRUE): always loaded in HBM, never evicted. Top-20 by traffic.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.
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.
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.
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:
The toy arithmetic is useful because it explains the product behavior. For Llama 3 70B at FP16, the raw KV per token is:
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.
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.
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:
The rest of this chapter is a taxonomy of how frontier models build that hierarchy.
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).
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 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.
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:
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.
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.
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.
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.
DeepSeek's Native Sparse Attention (NSA) is the canonical design. It has three paths:
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-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.
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 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.
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.
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:
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 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.
Temporal sparsity changes the state object the engine must manage. Instead of storing a growing KV tensor for every layer, hybrid models store:
The interview-level summary: full attention is exact but grows with context; recurrence is cheap but compressive; frontier models increasingly use both.
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.