The whole design, in writing
Learn AI system design by building an LLM inference serving system step by step. An interactive guide covering the request queue, the prefill/decode split, continuous batching, the KV cache and paged attention, tensor sharding across GPUs, autoscaling on GPU pressure, and the latency-versus-throughput trade-offs of serving a model at scale.
Every step of the build above, written out: the problem each piece solves, the option that was taken and the ones that were not, the numbers, and how it fails in production.
The big idea
Why is serving an LLM hard?
Running a model once on your laptop is easy. Serving it to thousands of users at once, fast and affordably, on hardware that costs dollars per hour, is not. A single request can hog a GPU for seconds while it generates tokens one at a time. How do you keep scarce GPUs full and users’ answers fast?
Treat the GPU fleet as a precious, fixed resource and design everything around it: a queue to absorb spikes, continuous batching to keep GPUs full, a KV cache to make each token cheap, sharding to fit big models, and autoscaling to ride demand.
What the new pieces do
- Clientclient
- Sends a prompt and reads tokens back as they stream. Cares about time-to-first-token and tokens/sec.
Step 1 · The skeleton
A client, a router, a model
A client sends a prompt and wants tokens streamed back. The model lives on GPUs that can’t be exposed directly. What sits in between?
Requests arrive for a model running on a GPU fleet. What fronts it?
No auth, no load balancing, no limits — and a client pinned to one worker that might be busy. You need a router in front.
The router authenticates, applies limits, and sends each request to a GPU replica that has room — then streams tokens back.
Generations are dynamic and rarely identical, so caching whole answers barely helps. The expensive part is fresh inference, not delivery.
A Router fronts the fleet: it admits requests, applies rate limits, and routes each to a GPU Worker with capacity, then streams the generated tokens back to the Client. The classic client → router → workers spine, tuned for streaming.
Why this piece earns its place
A router in front of a stateless service is a solved problem; in front of a stream it is not. The response stays open for as long as the answer takes, which makes the router stateful for the life of the request and breaks three defaults at once. Any buffering proxy between you and the user holds tokens until the response completes, so the server records a fast first token while the user watches a blank screen — a latency failure invisible in your own telemetry, because your own telemetry is the thing lying. Idle timeouts have to tell a long generation apart from a dead connection, which means a heartbeat, not a bigger number. And once the first token is on the wire the status code has already been sent, so a failure halfway through has to be delivered in-band, and a retry is no longer safe: the user has read part of an answer you are about to contradict. The piece to build early is cancellation. When a client disconnects, that has to reach the scheduler, or the GPU keeps decoding tokens nobody will read and reports the waste as utilization.
What the new pieces do
- Routerbackend
- The front door. Authenticates, applies limits, and routes each request to a model replica with capacity.
Step 2 · Absorb the spikes
Put a queue in front of the GPUs
Traffic is bursty — quiet, then a flood. GPUs are a fixed number. If requests hit workers directly, a burst either overwhelms them or, when quiet, leaves them idle. How do you smooth this?
GPU count is fixed; traffic is spiky. How do you avoid stampede AND starvation?
Dropping on every spike is a terrible experience and wastes capacity that frees up moments later. Buffer first, drop only as a last resort.
GPUs take minutes to provision and cost too much to hold per request. You can’t scale per-request at GPU granularity in real time.
A queue absorbs bursts so GPUs stay fully fed without being stampeded, and admission control can shed load gracefully when the queue gets too deep.
A Request Queue sits between the router and the GPUs. Bursts fill the queue instead of crushing the workers; in quiet moments the queue drains and GPUs stay busy. Admission control can shed load ("busy, retry") when the queue grows too deep — protecting latency for everyone already in flight.
Why this piece earns its place
A queue with no bound is not a shock absorber, it is a delay line: past a certain depth, every request that enters is already condemned to expire before a GPU reaches it. So two cheap mechanisms go in with the queue itself. A hard cap, so admission control has something to fire on. And a deadline stamped at enqueue and rechecked at dequeue, so work whose client has already given up is dropped before it buys a prefill rather than after. What admission cannot do is size the job. A request carries its prompt and a max-token cap, and that cap is a ceiling rather than a forecast — the true output length is settled only when generation stops — so every decision here is made against an upper bound and has to stay conservative about it. Two consequences worth naming. A refusal is only useful if it carries a Retry-After with jitter; a bare one returns as a synchronized wave and refills the queue you just drained. And the queue belongs in one place in front of the fleet, not as a deep inbox per replica: once a request is committed to a worker it cannot be moved, so it waits behind that worker’s longest prefill while another replica sits idle.
What the new pieces do
- Request Queueservice
- Holds incoming requests so the GPUs are never starved or stampeded — the buffer that absorbs spiky traffic.
Step 3 · How a token is made
Prefill, then decode
To serve a model well you have to know how it actually computes. Generating an answer isn’t one operation — it’s two very different phases with very different costs. What are they?
Generating an answer splits into two phases. What are they?
Weights load once at startup, not per request. The per-request work is the two-phase prefill/decode, which dominates everything.
Prefill runs the full prompt through the model in parallel (compute-heavy); decode then generates output tokens one at a time (memory-bandwidth-heavy). They scale differently.
That’s a different architecture’s framing. Decoder-style LLM serving is specifically prefill (prompt) then autoregressive decode (output).
The GPU Workers do two phases. Prefill processes the entire prompt in one parallel pass (compute-bound). Decode then generates output one token per step, each step depending on the last (memory-bandwidth-bound). The asymmetry — fast parallel prefill, slow sequential decode — shapes every serving decision.
Why this piece earns its place
The reason to hold the two phases apart in your head is that they compete for the same device. A prefill is one large compute burst, and while it runs, every sequence already streaming gets no token — so one long prompt arriving mid-flight lands as a visible stutter in dozens of other people’s answers. That damage hides from both metrics named above: time-to-first-token is fine, total time is fine, and the thing that moved is inter-token latency, which almost nobody graphs. Graph it and the fixes are available — split a long prefill into chunks the scheduler interleaves between decode steps, or run prefill and decode on separate pools so the burst lands where nobody is streaming. The asymmetry also explains why the next step works at all. Decode reads the entire weight matrix to produce a single token, so running many sequences together reads those weights once instead of once each — the bandwidth that dominates the step is paid for the whole batch, not per member of it. Prefill is already compute-saturated by one long prompt, so batching buys it far less. Batching is a decode optimization.
- 1 passprefill the prompt
- 1 tokenper decode step
- TTFTset by prefill
What the new pieces do
- GPU Workersservice
- The model running on GPUs. Does a one-time prefill of the prompt, then decodes output one token per step.
Step 4 · Keep the GPUs full
Continuous batching
Decode generates one token per step per sequence — a single request barely uses the GPU. But requests start and finish at different times and have different lengths. Naive batching wastes the GPU waiting for the slowest one. How do you keep it packed?
Requests have different lengths and arrive at different times. How do you batch?
A single decode step uses a fraction of the GPU; everything else waits. You’re paying for hardware that mostly idles.
The whole batch stalls on the longest sequence, and new arrivals wait for the next batch. Idle gaps everywhere.
A scheduler packs many sequences into each GPU step, dropping finished ones and slotting in new arrivals immediately. The GPU stays saturated regardless of length mix.
A Batch Scheduler does continuous (in-flight) batching: at each decode step it runs all active sequences together, retires the ones that just finished, and slots waiting requests in immediately — no waiting for a batch to drain. The GPU stays full no matter how lengths and arrivals mix. Serving metrics (TTFT, tokens/sec, utilization) drive its decisions.
Why this piece earns its place
Continuous batching removes the batch boundary, and that boundary was quietly doing a job: it was the moment everyone waiting got a turn. Without it, nothing in the loop is obliged to ever admit a particular queued request, so a steady arrival rate plus a few long-lived residents can hold the GPU indefinitely while a handful of requests age out at the back. Median latency looks excellent throughout, which is why this is usually found by a customer rather than a dashboard. The repair is explicit — aging, a priority that climbs with wait time, a slot reserved per step for the oldest waiter — but it has to be written, because run whatever is active contains no notion of fairness. The second surprise is that batch composition changes the output. Reductions run in a different order at a different batch size, so an identical prompt with an identical seed can produce a different token depending on who happened to be decoding alongside it. Nothing is broken. But bit-exact reproducibility is not a property this server has, and any eval that assumes it will flake forever.
What the new pieces do
- Batch Schedulerservice
- Packs many in-flight requests into each GPU step, adding and retiring sequences token by token.
- Metrics / Controlbus
- Streams serving metrics — time-to-first-token, throughput, GPU memory — that drive batching and autoscaling.
Step 5 · Make each token cheap
The KV cache and paged attention
At each decode step, attention needs to look back over every previous token. Recomputing that for the whole sequence on every single token would be brutally quadratic. And the cache that avoids it can blow up GPU memory. How do you make decode both fast and memory-efficient?
How do you avoid recomputing attention over the whole prompt every token?
That’s the quadratic waste the KV cache exists to kill — recomputing all prior tokens for every new one is enormously expensive.
Store the attention K/V per token so each new token reuses prior work — and page the cache in fixed blocks so many sequences pack into GPU memory without fragmentation.
Answers rarely repeat verbatim, so that barely helps. The win is caching intermediate attention state within a single generation.
The KV Cache stores each token’s attention key/value state so every new token reuses prior work instead of recomputing — turning decode from quadratic to linear. Paged attention stores that cache in fixed-size blocks (like OS virtual memory), so many sequences share GPU memory without waste — which is exactly what lets continuous batching pack so many requests in.
Why this piece earns its place
Two things past the headline. The first is that the KV budget is a remainder: weights and activations take the device memory they need, and whatever is left decides how many sequences can be held at once. So the deployable batch size is a consequence of load-time decisions — which model, which precision, which sharding degree — not a dial you can turn when the queue grows. Block size is the dial that stays, and it is two-sided: large blocks bring back the waste paging removed, since every sequence sits on a partly filled last block, while small blocks grow the block table and the indirection the attention kernel pays on every step. The second is what blocks make possible elsewhere. Because allocation is per block rather than per sequence, two requests beginning with identical text can point at the same blocks, reference-counted and copied only where they diverge. A shared system preamble, or a chat history resent in full every turn, then costs its prefill once instead of once per request — provided both requests land on the same replica, which quietly makes prefix reuse a routing decision as much as a memory one.
What the new pieces do
- KV Cachecache
- Per-sequence attention state kept in GPU memory so each new token reuses prior work instead of recomputing.
Back of the envelope
- KV cache
- reuse attention state — decode goes O(n), not O(n²)
- paged attention
- fixed-block cache → no fragmentation, dense packing
- memory caps concurrency
- more KV memory = bigger batches
- evict / preempt
- pause low-priority sequences when memory is tight
Step 6 · Fit a giant model
Shard the weights across GPUs
A frontier model’s weights don’t fit in a single GPU’s memory. You can’t shrink the model. So how do you run something bigger than any one device?
The model’s weights are larger than one GPU’s memory. Now what?
Sometimes valid — but when you need the big model, you must run it as-is. The serving system has to handle models bigger than one device.
Split each layer’s tensors across GPUs (tensor parallelism), with fast interconnect so they act as one logical worker. Pipeline parallelism splits by layer for even larger models.
Disk is orders of magnitude too slow for per-token weight access. Weights must live in GPU memory — sharded if they don’t fit on one.
The Model Weights are tensor-sharded across several GPUs: each layer’s matrices are split so the GPUs compute a forward pass together over fast interconnect, acting as one logical worker. Even larger models add pipeline parallelism (split by layer). The fleet is then many such multi-GPU workers behind the scheduler.
Why this piece earns its place
Sharding changes what a worker is, and those consequences outlive the memory problem it solved. The shard group becomes one failure unit: the replica is healthy only while every rank is, so a single sick GPU takes down a whole worker rather than a fraction of the fleet. Worse, the characteristic failure is not a crash. The ranks synchronize at every layer, so when one stops participating the others sit inside a collective that never returns — the process is alive, the port answers, the health check passes, and the worker serves nothing. Put a timeout on the collective and fail the group loudly; an undetected hung rank is the longest outage available in this design. Rollouts pay too, since every rank must load its slice before the group can serve, so a deploy removes the whole group at once and spare capacity has to cover it. And the degree is not simply the smallest that fits: more devices also means more aggregate memory left over for KV, so the throughput-optimal shard count is often wider than the one that merely makes the weights fit.
What the new pieces do
- Model Weightsstore
- The parameters, split across multiple GPUs when the model is too big for one device.
Step 7 · Ride the demand
Autoscale on GPU pressure
Demand swings through the day. Provision for the peak and you burn money on idle GPUs at 3am; provision for the average and you melt at peak. GPUs aren’t instant to add. How do you size the fleet?
Demand swings hour to hour and GPUs are slow to provision. How do you size the fleet?
You pay for peak capacity around the clock, idling expensive GPUs most of the day. Bleeds money.
The queue smooths small bursts, but a sustained peak just grows the queue and latency without end. Buffering isn’t capacity.
An autoscaler watches the real pressure signals and adds/removes GPU replicas, keeping a warm pool so scale-up isn’t cold-start slow. Match capacity to demand.
An Autoscaler watches the true pressure signals — queue depth and GPU utilization — and adds or removes replicas to hold latency targets. Because GPUs are slow to start, it keeps a warm pool and scales ahead of the curve. Under extreme load it sheds or routes overflow to a smaller, faster model.
Why this piece earns its place
Autoscaling a stateless web tier is a reflex, and two of its assumptions are false here. The first is that a replica can be removed when you decide to remove it. A worker holds live sequences whose attention state sits on its own device and cannot be moved elsewhere, so scale-down is not a kill — it is stop admitting, then wait for the resident set to finish, which can take as long as the longest generation in flight. That drain time is the floor on how fast the fleet can shrink, and it needs a third state the router has to understand: a replica that is draining takes no new work while it keeps streaming the work it already holds. Without it, scale-in is indistinguishable from cutting people off mid-sentence. The second false assumption is symmetry. The two errors do not cost the same — one replica short is visible in every waiting user’s first token within seconds, one replica too many costs money and nothing else — so the loop should decide on a short window on the way out and a long one on the way back. It also has to count replicas that are still coming up as capacity already ordered. A loop that reads pressure, orders, then reads that same pressure again on the next tick keeps ordering into a queue already being answered, and lands a fleet far larger than anyone asked for.
What the new pieces do
- Autoscalerservice
- Watches queue depth and GPU utilization and adds or removes replicas to hold latency under load.
The payoff
You built an inference server
From "run it once" to serving at scale: a router and queue, the prefill/decode split, continuous batching, a paged KV cache, tensor-sharded weights, and autoscaling — all balancing the latency/throughput trade-off.
Now flood the queue and feel the central tension: as load climbs, time-to-first-token rises, and you watch batching, the KV cache, and autoscaling fight to protect tail latency before the fleet stalls.
Everything you assembled, in order
- Router + stream — admit, route, stream tokens — optimize TTFT
- Request Queue — shock absorber between spiky demand and fixed GPUs
- Prefill / decode — parallel prompt pass, then one token per step
- Continuous batching — add/retire sequences each step — GPUs stay full
- KV cache + paging — cheap tokens, dense concurrency, memory-bound
- Tensor sharding — split big models across GPUs to fit and scale
- Autoscaler — scale on queue depth + GPU util, keep warm pools
- Latency vs throughput — the trade-off every knob above is balancing
