Design an LLM inference platform
A user types a question and waits. The first word has to appear within about a second or the product feels broken, and the next four hundred have to arrive faster than they can read. Behind that is a machine that costs more per hour than the engineer who provisioned it, and which sits idle through most of the generation unless you do something unusual with it.
This question is newer than the rest of the course and it is being asked constantly. It is also the one where the standard web intuitions are actively wrong. The unit of capacity is not requests per second. The bottleneck is not CPU. And a bigger batch makes latency worse in a way that surprises almost everybody the first time.
Step 1: Understand the problem
Two of these answers are the whole design, and the hardware supply answer is the one that changes what kind of problem this is.
| You ask | They say | What it settles |
|---|---|---|
| Interactive chat, or batch jobs? | Interactive, streaming tokens as they are produced. | Two separate latency numbers, time to first token and the gap between tokens after that. One number called "latency" hides the only interesting failure in this system. |
| How long are prompts, and how long are answers? | Prompts from a hundred tokens to a hundred thousand. Answers a few hundred. | The most important number on the page. Reading the prompt and writing the answer are different kinds of work on the same chip, and the ratio between them decides how you schedule. |
| Can we add GPUs when traffic rises? | Not quickly. The fleet is fixed and supply constrained. | This is a scheduling problem under scarcity, not an autoscaling problem. Admission control and queueing become product decisions rather than safety valves. |
| One model or many? | A few sizes, plus customer fine tunes. | Routing by model, and a hard constraint: loading a hundred gigabytes of weights takes minutes, so you cannot swap models on demand. Either a replica is dedicated, or variants share a base model through adapters. |
| Do requests from one customer share anything? | Yes. Every request carries the same 2,000 token system prompt. | Prefix caching, which can remove most of the prompt reading work for most requests. It is the cheapest large win available and it is invisible unless you ask this question. |
| Under overload, do we go slow or refuse? | Refuse clearly, rather than making everybody slow. | Bounded queues and an explicit 429. An unbounded queue in front of a GPU means every user gets a bad experience instead of most users getting a good one. |
What you are building, and what you cut
- Serve streaming token responses. Held connections, cancellable, with per tenant priority.
- Batch many requests onto one replica. Continuously, not in fixed groups. This is where the throughput comes from.
- Manage KV cache memory. It is the real capacity limit, and it moves with context length.
- Queue, prioritise and shed. On hardware you cannot buy more of this week.
- First token in under 1 second at p95.
- Under 50ms between tokens at p95, which is faster than anyone reads.
- GPUs busy with useful work above 80% of the time.
- One 100,000 token prompt must not stall everybody else on that replica.
- Training and fine tuning. A different fleet, a different failure model, and nothing here is on its path.
- The model itself: quality, evaluations, and safety behaviour. Assume the weights are given.
- Retrieval and tool calling. Those are services that call this one, and treating them as part of it is how the design loses focus.
- Guardrail classifiers. Name them as a sidecar on the way in and the way out, with their own latency budget.
Back of the envelope
Drag the context length last, and watch what it does to everything else.
The 320 KB per token is for a 70 billion parameter model with grouped query attention, and it is the number worth memorising. Now drag the context length to 128,000: concurrency collapses, because each sequence is holding tens of gigabytes of cache. Long context is not a feature you switch on, it is capacity you spend, and the person who adds it to the product is spending most of your fleet.
Almost everyone sizes this in requests per second. Requests per second is an output of this system, not an input to it, because a request holds memory for its whole lifetime and that memory scales with how much text is in it. The second mistake is trusting GPU utilisation: a chip can report 100% while doing almost no useful arithmetic, because generating one token at a time is limited by how fast weights can be read out of memory rather than by how fast they can be multiplied.
Step 2: Propose the high level design
The API
{
"model": "big-70b",
"messages": [...],
"max_tokens": 600,
"stream": true
}{ "requests": [...], "completeWithin": "24h" }The data model
There is almost no durable state in the serving path, and saying that out loud is worth a mark. The interesting state is a few thousand rows in the memory of each replica, and the scheduler rewrites most of them tens of times a second.
| request_id | uuid | PK | Lives only as long as the generation. Nothing here survives a restart, and nothing needs to. |
| phase | enum | IDX | queued, prefill, decode, done. The scheduler picks a mix of prefill and decode work for every iteration, so this field is read constantly. |
| kv_pages | array of page ids | Which cache pages this sequence owns. Not a contiguous block, which is the subject of the third deep dive. | |
| generated | int | Tokens produced so far, against max_tokens. A sequence that hits its limit frees its pages this iteration. | |
| tenant, priority | varchar, enum | Fairness lives here. Without it, one customer with a script takes the whole replica and it looks like an outage to everyone else. | |
| arrived_at | timestamp | Queue age, which is the metric to alert on. Queue depth looks alarming during a normal burst and tells you nothing. |
The whole system on one whiteboard
Walking Figure 1:
- The client opens a streaming request. That connection stays open for the whole generation, which is seconds rather than milliseconds, so the gateway tier is sized by concurrent connections and not by requests a second.
- Quota and safety checks run here, on cheap hardware, before anything reaches a GPU. Rejecting a request after it has occupied a GPU slot is the expensive way to say no.
- The router picks a model pool, and then the request waits in a bounded queue. If the queue is full it is refused immediately rather than joining the back of something hopeless.
- Admission is where fairness happens: priority, per tenant limits, and a check that there is enough free cache memory to hold this sequence if it is let in.
- The scheduler builds the work for the next iteration, and it does this again for the next one, tens of times a second. The batch is not a thing you assemble once.
- The replica runs that iteration on the GPU, reading the prompt for new sequences and producing one more token for sequences already going.
- Tokens stream back out through the gateway as they appear. Nothing waits for the generation to finish, including the user.
The prefix cache on the right is the one box that is pure profit. Two thousand tokens of system prompt, identical on every request from a customer, get read once rather than ten thousand times.
Step 5 rebuilds the batch every iteration instead of filling one up and running it to completion. What specifically goes wrong with the simpler version, where you gather eight requests and run them together until they are all done?
Step 3: Design deep dive
Reading the prompt and writing the answer are two different machines
Everything confusing about serving these models comes from one fact: a request has two phases with opposite characteristics, running on the same chip.
Two consequences to say out loud, because they are what the metrics mean:
Time to first token is dominated by queueing and by prefill, so it gets worse when the system is busy or when prompts are long. The gap between subsequent tokens is dominated by how many sequences are decoding together, so it gets worse as the batch grows. They move in opposite directions under the same load, which is exactly why a single average latency number tells you nothing useful about this system.
Continuous batching, and why a bigger batch is not the fix
The obvious way to batch is to collect requests until you have a group, run them together, return the results, and start again. It is also the version that wastes most of the fleet.
Figure 3 is the same idea in time, because the thing worth seeing is a request joining a batch that is already running.
“So should I make the batch as large as possible?” No, and the reason is the second metric. Every sequence you add to the batch makes each iteration slightly slower, so the gap between tokens grows for everybody. Past a point you are trading a visibly smooth response for throughput nobody asked for. Pick the batch size from your inter token latency target rather than from your utilisation graph, and say that out loud, because it is the trade the utilisation graph hides.
The KV cache is the capacity
Every token a sequence has seen or produced leaves behind a few hundred kilobytes of cached state that must stay resident for as long as that sequence is alive. That memory, not compute, is what limits how many conversations a replica can hold.
Say the operating system comparison in the interview, because it lands immediately. This is paging. The page table is per sequence, shared pages are copy on write, and running out of memory means you either evict a sequence or stop admitting new ones. Eviction is less painful here than it sounds, because recomputing a dropped sequence’s cache is a prefill, and prefill is the fast phase.
Break it
Trade-offs
| Choice | What you gain | What you pay | Pick it when |
|---|---|---|---|
| Continuous batching | Finished sequences free their slot immediately and waiting ones join within one iteration, which is where almost all of the throughput comes from. | A scheduler that runs tens of times a second, and a much more complicated engine than a loop over a fixed batch. | Any interactive serving. Static batching is only defensible for offline jobs where latency does not exist as a concept. |
| Paged KV cache | Memory follows actual usage rather than the worst case, which typically multiplies concurrency several times over, and makes prefix sharing possible at all. | An indirection inside the attention kernel, and a page table to manage per sequence. | Anything beyond a demo. The contiguous version is simpler and quietly throws away most of your hardware. |
| Admission control with a hard 429 | Everyone admitted gets a usable experience, and the refusal arrives fast enough for a client to do something sensible. | Visible errors on your dashboard, and a conversation with whoever owns the customer that got refused. | Always, on hardware you cannot buy more of this week. An unbounded queue converts a capacity problem into a universally bad experience. |
| Separate pool for very long contexts | A hundred thousand token request cannot stall the replica serving ordinary chats. | Scarce GPUs split into pools, so each pool is less efficient than one big one would be. | When the long tail of prompt length is an order of magnitude above the median, which it usually is the moment you ship a document feature. |
Interview replay
Checkpoint
1. Why is concurrent sequences a better capacity unit than requests per second?
2. What does continuous batching fix that simply increasing the batch size does not?
3. Same platform, but now serving image generation instead of text. What changes first?
- 320KB of KV cache per token on a 70 billion parameter model. If you memorise one number, make it this one.
- 8,000 tokens is 2.5GB per sequence, so a 320GB replica holds tens of conversations, not thousands.
- Two latency numbers: first token is queue plus prefill, inter token is batch size.
- Prefill is compute bound, decode is memory bandwidth bound. Every other decision follows from that one sentence.
The thing to establish first is that capacity here is concurrent sequences, not requests per second, because every sequence holds KV cache memory for its whole life and that memory grows with context length. On a replica holding a seventy billion parameter model, an eight thousand token conversation costs about two and a half gigabytes of cache, so you hold tens of them, and long context requests cost an order of magnitude more. A request has two phases that behave oppositely: prefill reads the whole prompt in one compute bound pass and decides time to first token, and decode produces one token per pass, limited by memory bandwidth, which is why batching is nearly free and where all the throughput comes from. So the engine rebuilds its batch every iteration rather than running fixed groups, a finished sequence frees its slot immediately, and a waiting one joins tens of milliseconds later. KV cache is paged like virtual memory, which multiplies concurrency several times and lets identical system prompts share pages with a reference count. In front of it is a bounded queue with a fast 429, because on hardware you cannot buy more of, refusing quickly is a better product than answering in ninety seconds. The failure I would call out is one enormous prompt stalling every other stream on a replica, and the fixes are chunked prefill and a separate pool for long contexts.
