ChatGPT / LLM Serving Architecture
Sharpened prompt. Design the serving infrastructure for a large language model handling 100k concurrent conversations, streaming the first token within 500ms, sustaining high GPU utilisation across requests with wildly different prompt and generation lengths, keeping multi-turn conversations from recomputing their history, serving both free and paid tiers fairly, and doing it on hardware that costs more per hour than most of the rest of your fleet combined.
The mental shift that makes this problem tractable: the GPU is the database, VRAM is the disk, and the scheduler is the whole system. Every design decision is about what occupies expensive memory and for how long.
1. Problem framing
Functional requirements
- Streaming chat completions with multi-turn conversation history.
- Multiple models and sizes; tool calling; structured output.
- Per-tenant rate limits and quotas; free and paid tiers.
- Safety filtering on input and output.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Time to first token | P95 < 500ms | Prefill must be fast and prefix cache hits must be common |
| Inter-token latency | < 50ms (≈20 tokens/sec, faster than reading) | Batch size is bounded by latency, not just memory |
| GPU utilisation | > 70% | Continuous batching is mandatory; static batching wastes half the fleet |
| Cost per million tokens | The dominant business metric | Everything here is really a cost-optimisation problem |
| Fairness | A free user's long generation must not starve a paid request | Preemption and priority scheduling |
Back-of-the-envelope
Model: ~70B parameters at FP8/INT8 = ~70 GB of weights
-> 2× H100 (80 GB) with tensor parallelism, or 1× H200
KV cache: per token = 2 (K,V) × layers × kv_heads × head_dim × bytes
≈ 2 × 80 × 8 × 128 × 2 B ≈ 320 KB/token with GQA at FP16
A 4,000-token conversation = ~1.3 GB of KV cache. PER USER.
VRAM budget: 160 GB total − 70 GB weights − ~10 GB activations = ~80 GB for KV
80 GB / 1.3 GB = ~60 concurrent 4k-token conversations per node.
THIS is the capacity limit, not FLOPs.
Throughput: ~2,500 output tokens/sec/node at batch 60
100k concurrent conversations, ~5% actively generating = 5,000 active
-> 5,000 / 60 ≈ 85 nodes = ~170 H100s just for decode
Prefill: a 4,000-token prompt is ~4,000× the compute of one decode step.
Prefill and decode have completely different performance profiles (4.6).
The KV-cache arithmetic is the single most important calculation on this page. Concurrency is bounded by memory, not by compute — say that, and the rest of the design follows naturally.
2. High-level architecture
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Inference engine | vLLM, TensorRT-LLM, or SGLang | PagedAttention, continuous batching, and prefix caching are table stakes; do not build them |
| KV memory manager | PagedAttention block allocator | Turns fragmentation from a hard limit into a non-issue (4.1) |
| Scheduler | Iteration-level continuous batching with priority | The single largest throughput lever (4.2) |
| Router | Consistent hash on conversation prefix, load-aware | Turns a cache miss into a cache hit (4.3) |
| Transport | SSE for simple streaming; WebSocket when bidirectional | Tokens must reach the user as they are produced |
| Quantisation | FP8 weights, FP8/INT8 KV cache | Roughly doubles the concurrency the same VRAM supports |
| Safety | Small distilled classifiers, streaming on output | A large safety model would cost as much as the generation |
4. The toughest parts
4.1 KV cache memory: why naive allocation wastes most of your VRAM
Why it's hard. Each sequence needs a contiguous KV cache sized for its maximum possible length, because you cannot know in advance how long the generation will be. Reserve 4,096 tokens for a request that generates 200 and you have wasted 95% of that allocation. Across a batch, classical serving systems wasted 60–80% of KV memory to internal and external fragmentation — and since memory is the concurrency limit, that waste translates directly into needing several times more GPUs.
Solution — PagedAttention: treat KV cache like operating-system virtual memory.
Classical (contiguous reservation):
Seq A: [████░░░░░░░░░░░░] reserved 2048, used 512 -> 75% wasted
Seq B: [██████░░░░░░░░░░] reserved 2048, used 768 -> 62% wasted
A new sequence needs a contiguous 2048-token region. None available -> rejected,
despite plenty of free memory scattered across the allocation.
PagedAttention (fixed-size blocks, non-contiguous):
Physical blocks (16 tokens each): [B0][B1][B2][B3][B4][B5][B6][B7]...
Seq A logical->physical: [B0, B3, B7] grows one block at a time
Seq B logical->physical: [B1, B2, B5, B6]
Waste is bounded by (block_size - 1) tokens per sequence: under 4%.
class BlockAllocator:
"""The GPU-memory equivalent of a page table. Blocks are fixed-size and
need not be contiguous, so external fragmentation is eliminated."""
def __init__(self, total_blocks, block_size=16):
self.free = list(range(total_blocks))
self.tables = {} # seq_id -> [physical block ids]
self.refcount = defaultdict(int) # for copy-on-write sharing
def append_token(self, seq_id, position):
table = self.tables.setdefault(seq_id, [])
if position % self.block_size == 0: # need a new block
if not self.free:
raise OutOfKVMemory # triggers preemption (4.2)
blk = self.free.pop()
table.append(blk)
self.refcount[blk] += 1
def fork(self, parent_seq, child_seq):
"""Copy-on-write: parallel samples or a shared system prompt reference
the SAME physical blocks until one of them diverges."""
self.tables[child_seq] = list(self.tables[parent_seq])
for blk in self.tables[child_seq]:
self.refcount[blk] += 1 # shared, not copied
Three consequences worth stating. Waste drops from 60–80% to under 4%, which roughly triples or quadruples the number of concurrent sequences the same GPU supports — the largest single win available in LLM serving. Copy-on-write sharing means N parallel samples from one prompt, or a thousand users sharing one system prompt, store that prefix's KV once. And blocks can be evicted and swapped to host RAM, which is what makes preemption possible in 4.2.
Reduce the per-token cost too: grouped-query attention (GQA) shares K/V heads across query heads, cutting KV size several-fold with negligible quality loss, and FP8 KV quantisation halves it again. Both directly multiply concurrency.
4.2 Continuous batching: never wait for the slowest sequence
Why it's hard. GPUs need large batches to be efficient. But requests have wildly different generation lengths — one produces 20 tokens, another 2,000. With static batching, the whole batch is locked until the longest sequence finishes, so 31 of 32 slots sit idle producing nothing while one sequence grinds on. Utilisation collapses, and new requests queue behind a batch that is 97% empty.
Solution — iteration-level scheduling: re-form the batch after every single token.
def serve_loop(engine, scheduler):
running = [] # sequences currently decoding
while True:
# 1. Evict finished sequences and free their KV blocks IMMEDIATELY.
for seq in [s for s in running if s.finished]:
allocator.free(seq.id)
running.remove(seq)
seq.stream.close()
# 2. Admit new work into the freed slots, budget permitting.
while scheduler.has_waiting() and can_admit(running):
cand = scheduler.peek_highest_priority()
if allocator.can_allocate(cand.prompt_len + RESERVE):
running.append(scheduler.pop()) # joins the very next iteration
else:
break
# 3. If memory is exhausted, PREEMPT — lowest priority / newest first.
while allocator.pressure() > 0.95 and running:
victim = min(running, key=lambda s: (s.priority, -s.arrival))
allocator.swap_out(victim.id) # KV blocks -> host RAM
scheduler.requeue(victim) # resumes later, not restarted
running.remove(victim)
# 4. One decode step for the whole batch: a single fused kernel launch.
engine.step(running) # each seq gets exactly one token
for seq in running:
seq.stream.send(seq.last_token) # stream immediately, per token
Static batching typically achieves 30–40% GPU utilisation on realistic traffic; continuous batching reaches 70–90%. That is a 2–3× throughput gain on the same hardware, which for a fleet of thousands of GPUs is the difference between viable and not.
Two subtleties worth raising. Batch size is bounded by latency, not just memory: a larger batch improves throughput but increases per-token latency for everyone in it, so cap the batch where inter-token latency reaches your budget (~50ms). And preemption must swap, not restart — swapping KV blocks to host RAM preserves the work already done, whereas recomputing a 4,000-token prefill from scratch wastes seconds of GPU time. That is only possible because PagedAttention made KV relocatable.
4.3 Prefix caching and routing: don't recompute the conversation
Why it's hard. Turn 10 of a conversation sends the entire history — system prompt, plus nine previous exchanges, perhaps 6,000 tokens — of which only the last ~50 are new. Recomputing attention over all 6,000 costs roughly 6,000× a decode step and blows the TTFT budget. Worse, if the request lands on a different GPU than turn 9, whatever was cached is on the wrong machine.
Solution — a radix tree of KV blocks for automatic prefix reuse, plus routing that sends a conversation back to the node that already holds it.
# Node-local: automatic prefix sharing via a radix tree over token blocks.
class PrefixCache:
"""Keyed by the HASH of the token sequence, so sharing is automatic and
requires no application cooperation. Two requests with the same system
prompt share its KV blocks without knowing about each other."""
def match(self, token_ids):
node, matched = self.root, 0
for blk_start in range(0, len(token_ids), BLOCK):
h = hash_block(token_ids[blk_start:blk_start + BLOCK])
if h not in node.children:
break
node = node.children[h]
node.last_used = now()
matched += BLOCK
return matched, node.blocks # skip prefill for `matched` tokens
# Cluster-level: route to where the KV already lives.
def route(request) -> Node:
# Hash a stable prefix of the conversation, not the whole request — the tail
# changes every turn, and hashing it would send every turn somewhere new.
key = hash(request.system_prompt) ^ hash(request.conversation_id)
preferred = consistent_ring.get(key)
# Load-aware: honour affinity only if that node has capacity, otherwise the
# cache hit is paid for with queueing delay.
if preferred.kv_pressure < 0.85 and preferred.queue_depth < MAX_Q:
return preferred
# Fall back to the least-loaded node among the next few on the ring, so we
# keep SOME locality rather than scattering randomly.
return min(consistent_ring.successors(key, n=3), key=lambda n: n.load)
Two effects compound here. Prefix caching turns a 6,000-token prefill into a 50-token prefill — often a 10–50× TTFT improvement on multi-turn conversations, and it is entirely automatic once the radix tree exists. Affinity routing is what makes the cache hit likely; without it, a cluster of 85 nodes gives you a ~1% hit rate by chance.
The trade to name explicitly: strict affinity creates hotspots. A viral system prompt would pin enormous traffic to one node. The load-aware fallback trades some hit rate for balance, and tuning that threshold is a real operational knob. Also share the system prompt globally — precompute its KV once per node and pin it, since it is the single most reused prefix in the system.
4.4 Streaming, and the latency the user actually feels
Why it's hard. A 500-token response takes ~25 seconds to generate at 20 tokens/sec. Waiting for completion is unusable. But streaming introduces its own problems: safety filtering must happen on a partial output, a dropped connection mid-stream wastes all the GPU time spent so far, and the client must render tokens smoothly rather than in stuttering bursts.
Solution — token-by-token SSE, with safety on a sliding window and resumable streams.
async def stream_completion(request, response):
await response.prepare(content_type="text/event-stream")
buf, emitted = [], 0
async for token in engine.generate(request):
buf.append(token)
# Safety runs on a sliding window, not per token: a single token is
# meaningless, and per-token classification would cost more than generation.
if len(buf) - emitted >= SAFETY_WINDOW:
verdict = await safety.check_partial("".join(buf))
if verdict.blocked:
await response.write(sse({"type": "error",
"reason": verdict.category}))
await engine.abort(request.id) # free the KV blocks NOW
return
await response.write(sse({"delta": "".join(buf[emitted:])}))
emitted = len(buf)
# Backpressure: if the client is not consuming, stop generating.
# Continuing to burn GPU for a disconnected client is pure loss.
if response.transport.is_closing():
await engine.abort(request.id)
return
await response.write(sse({"delta": "".join(buf[emitted:]), "done": True}))
Three points. Safety on a window, not per token — classifying every token individually would cost more compute than generating them, and a single token carries no meaning. Abort on disconnect is worth real money: users close tabs constantly, and continuing to generate for an absent client wastes the most expensive resource you have; detecting it and freeing the KV blocks immediately is a direct cost saving.
And the streaming decision has a product dimension: emitting tokens the instant they exist means a blocked completion has already shown the user part of a bad response. Buffering a few tokens before emitting trades a little perceived latency for the ability to retract — a genuine trade-off worth naming rather than assuming one side.
For resumability, checkpoint the generated text against the request ID so a client that reconnects can resume from where it dropped rather than regenerating — important on mobile networks and cheap to implement.
4.5 Fairness: one free user must not starve a paying one
Why it's hard. A single request generating 8,000 tokens occupies a batch slot and its KV blocks for minutes. If admission is FIFO, a burst of long free-tier generations fills the fleet and a paying customer's short request waits behind them. But naive strict priority starves free users entirely, and preempting mid-generation wastes the work already done.
Solution — priority queues with weighted fair sharing, admission control, and preemption that swaps rather than discards.
class FairScheduler:
# Each tier gets a guaranteed share of decode slots. Unused share is lent out,
# so the fleet stays busy, but is reclaimed the moment the owner has work.
WEIGHTS = {"enterprise": 0.50, "pro": 0.35, "free": 0.15}
def select(self, running, capacity):
deficits = {}
for tier, w in self.WEIGHTS.items():
entitled = w * capacity
actual = sum(1 for s in running if s.tier == tier)
deficits[tier] = entitled - actual # most-starved first
for tier in sorted(deficits, key=deficits.get, reverse=True):
if deficits[tier] > 0 and self.queues[tier]:
return self.queues[tier].popleft()
return None
def admit(self, request) -> Decision:
# Admission control BEFORE the request occupies memory. Rejecting fast
# is far better than accepting and timing out after 60 seconds of queueing.
est_wait = self.estimated_wait(request.tier, request.prompt_len)
if est_wait > request.tier_sla:
return Decision.reject(429, retry_after=est_wait)
# Cap the KV a single request may hold, so one 100k-token context cannot
# consume a whole node's memory.
if request.max_kv_blocks() > MAX_PER_REQUEST_BLOCKS:
return Decision.reject(400, "context_too_long_for_tier")
return Decision.admit(priority=self.priority_of(request))
Three mechanisms, each necessary. Weighted fair sharing with lending keeps utilisation high while guaranteeing each tier its floor. Admission control turns an unbounded wait into a fast, honest 429 with a Retry-After — the same discipline as the rate limiter. And per-request KV caps prevent a single enormous context from monopolising a node.
Note the asymmetry with preemption: preempt the newest, lowest-priority sequence rather than the longest-running one. The long-running sequence has the most sunk cost, so evicting it wastes the most work — a small scheduling detail with a large efficiency consequence.
4.6 Prefill and decode are two different workloads
Why it's hard. Prefill (processing the prompt) is compute-bound: it runs attention over thousands of tokens in parallel and saturates the GPU's math units. Decode (generating one token at a time) is memory-bandwidth-bound: it does very little arithmetic but must read the entire model's weights from HBM for every single token. Running both on the same GPU means a long prefill blocks every decode in flight, causing visible inter-token stutter for every other user on that node.
Solution — chunked prefill at minimum; disaggregated prefill and decode pools when scale justifies it.
# Option A — chunked prefill: slice a long prompt so decodes interleave.
def schedule_iteration(running, waiting):
budget = TOKEN_BUDGET_PER_ITERATION # e.g. 2048 tokens of work
batch = []
# Decodes first: they are latency-critical and cheap (1 token each).
for seq in running:
batch.append(DecodeStep(seq)); budget -= 1
# Fill the remaining budget with a CHUNK of prefill, not the whole prompt.
for seq in waiting:
chunk = min(budget, seq.remaining_prefill_tokens)
if chunk <= 0: break
batch.append(PrefillChunk(seq, chunk)); budget -= chunk
return batch
# A 8,000-token prompt becomes 4 chunks across 4 iterations. No decode stalls
# for seconds; TTFT rises slightly; inter-token latency stays smooth for everyone.
Option B — disaggregation (at scale):
Prefill pool Decode pool
compute-bound, big batches bandwidth-bound, many sequences
fewer, high-FLOPs GPUs more GPUs, high HBM bandwidth
| ^
| KV cache transferred over |
+----- NVLink / InfiniBand -------------+
Each pool is sized and scaled independently against its own bottleneck,
and a long prefill can never stall an active decode.
The insight to articulate: these are not two phases of one workload, they are two workloads with opposite bottlenecks. Chunked prefill is the pragmatic answer for a single-pool deployment; disaggregation is what large deployments move to, because it lets you buy different hardware for each and scale them independently. Note the cost — transferring KV cache between pools needs very fast interconnect, which is why this is a large-deployment technique.
4.7 Cost, capacity, and the model rollout
Why it's hard. GPUs are the dominant cost and cannot be autoscaled in seconds — instances are scarce, boot takes minutes, and loading 70 GB of weights takes more. Traffic is diurnal with a 5× peak-to-trough ratio. Meanwhile a new model version must be rolled out without a capacity gap, and it may have different memory characteristics that invalidate your batch tuning.
Solution — tiered capacity, aggressive quantisation, request-level cost controls, and a staged rollout with real quality gates.
# Capacity: three tiers with different commitment and cost.
CAPACITY = {
"reserved": {"share": 0.60, "cost": 1.0}, # baseline; always on
"committed": {"share": 0.25, "cost": 1.4}, # scheduled for the diurnal peak
"spot": {"share": 0.15, "cost": 0.4}, # opportunistic, may vanish
}
# Spot GPUs serve ONLY preemptible free-tier traffic, and preemption is graceful
# because PagedAttention lets us swap KV out and resume elsewhere (4.2).
def route_by_cost(request):
# Not every request needs the largest model. Cascade: try small first, and
# escalate only when the small model's confidence or a verifier says so.
if request.tier == "free" and classifier.is_simple(request.prompt):
return SMALL_MODEL_POOL # ~10× cheaper per token
if request.needs_long_context:
return LONG_CONTEXT_POOL # different KV budget, different tuning
return FLAGSHIP_POOL
Three cost levers worth quantifying. Quantisation — FP8 weights plus FP8 KV roughly doubles concurrency per GPU and therefore halves cost per token, with a quality loss that must be measured rather than assumed. Model cascading — routing simple requests to a smaller model can serve a large fraction of traffic at a fraction of the cost, and the routing classifier is itself cheap. Speculative decoding — a small draft model proposes several tokens that the large model verifies in one forward pass, giving a 2–3× decode speedup when the draft is accurate, at the cost of extra memory and wasted work on mispredictions.
For rollout, weights are large and the differences matter:
1. Load the new model onto a small canary pool (loading 70 GB takes minutes —
plan for it; this is not a container restart).
2. Shadow traffic: run 1% of requests through both and compare outputs offline.
Quality regressions do not show up in latency or error metrics.
3. Ramp 1% -> 5% -> 25% -> 100%, with automatic rollback on guardrail metrics
(refusal rate, safety-classifier trigger rate, user regeneration rate).
4. Keep the previous version warm for one hour. Rollback is a routing change.
5. Re-tune batch size and KV block budget — a new model's memory profile
invalidates the old scheduler configuration.
Step 3's guardrails are the part that distinguishes LLM rollout from ordinary deploys: the failure mode is quality, not errors, and quality is invisible to conventional monitoring. Measuring user regeneration rate and refusal rate as automatic rollback triggers is the practical answer.
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| Traffic spike | KV memory exhausted before compute | Admission control with honest 429s; preempt-and-swap; per-request KV caps |
| A viral system prompt | Prefix-affinity hotspot on one node | Load-aware routing fallback; replicate hot prefixes across nodes |
| Very long context request | One sequence consumes a node's KV budget | Per-request block caps; a separate long-context pool with different tuning |
| Long prefill on a shared node | Every other user sees token stutter | Chunked prefill; disaggregation at scale |
| Spot GPU reclamation | Sequences on those nodes interrupted | Spot serves only preemptible traffic; swap KV and resume elsewhere |
| Model rollout regression | Quality drops with no error signal | Shadow comparison + regeneration-rate and refusal-rate guardrails |
| Client disconnects at scale | GPU time burned for nobody | Detect closed transports and abort immediately |
6. Cheat sheet
- The framing: GPU is the database, VRAM is the disk, the scheduler is the system. Concurrency is bounded by KV memory, not FLOPs.
- The number: ~320 KB of KV per token with GQA at FP16 → a 4k-token conversation is ~1.3 GB → ~60 concurrent conversations on 80 GB of free VRAM.
- PagedAttention: fixed-size non-contiguous KV blocks with a page table. Waste drops from 60–80% to under 4%; copy-on-write shares prefixes; blocks become swappable, which enables preemption.
- Continuous batching: re-form the batch every iteration. 30–40% → 70–90% utilisation. Cap batch size by inter-token latency, not memory.
- Prefix cache + affinity routing: radix tree over KV blocks for automatic reuse, consistent-hash routing with a load-aware fallback. 10–50× TTFT win on multi-turn.
- Prefill vs decode: compute-bound vs bandwidth-bound. Chunked prefill at minimum; disaggregated pools at scale.
- Fairness: weighted fair share with lending, admission control with honest 429s, per-request KV caps, preempt the newest not the longest.
- Cost: FP8 quantisation, model cascading, speculative decoding, spot for preemptible traffic only.
- The one-liner: "Every hard problem here is a memory-allocation problem in disguise — PagedAttention makes KV cache relocatable and shareable, and once it is, continuous batching, prefix reuse, and preemption all become possible."