Course Content
AI System Design and Architecture
3 sections · 7 lessons
Managing Latency, Throughput, and Caching
A team is told their inference endpoint is too slow. The model takes 24 ms, so they spend three weeks quantising it down to 11 ms — a genuine 2.2× speedup on the forward pass. They deploy. The p99 that users complain about moves from 139 ms to 120 ms. Thirteen per cent, for three weeks of work.
The forward pass was never the problem. When they finally measured every stage, 61 of those 139 milliseconds were spent waiting in a queue because the fleet ran at 91% utilisation. Two extra replicas — twenty minutes of work — would have taken p99 to about 92 ms, a bigger gain than the whole quantisation project.
This is the central discipline of performance work: measure the whole path, in percentiles, before you optimise anything. Almost every wasted optimisation comes from improving a stage that was not the bottleneck, or from optimising the average when users experience the tail.
Two different numbers
Latency is how long one request takes, measured end to end from the client's point of view. Throughput is how many requests the system completes per second. They are related but not interchangeable, and improving one often costs the other.
The link between them is Little's Law, which holds for any stable system regardless of its internals:
where L is the average number of requests in the system, λ is arrival rate, and W is average time in system. It is one of the most useful equations in systems engineering because you can rearrange it three ways:
- How much concurrency do I need? To serve 600 req/s at 40 ms each: L=600×0.040=24 requests in flight.
- What is my throughput ceiling? A server with 8 worker threads and 40 ms per request tops out at 8/0.040=200 req/s — no matter how fast the GPU is. This is where people get it wrong: they buy a bigger GPU when the limit was a thread pool size in a config file.
- How big should the connection pool be? If each request makes a 3 ms database call at 600 req/s, average concurrent DB calls is 600×0.003=1.8. A pool of 10 is generous; the pool of 100 people habitually configure just moves the queue into the database.
Why they trade against each other
Batching is the clearest case. Waiting to accumulate 32 requests before running the GPU raises throughput substantially and raises the latency of the first request in the batch by the whole collection window. Every technique below sits somewhere on this axis:
| Technique | Effect on latency | Effect on throughput |
|---|---|---|
| Larger batches | Worse (wait for fill) | Much better |
| More replicas | Better (less queueing) | Better |
| Quantisation | Better | Better |
| Caching | Much better on hits | Better (offloads compute) |
| Running at higher utilisation | Much worse | Marginally better |
| Adding a retry layer | Better at p99 (masks slow nodes) | Worse (extra work) |
Where the milliseconds actually go
Instrument every stage separately and record percentiles, not averages. Here is the endpoint from the opening scenario, measured properly:
| Stage | p50 | p99 | Share of p99 |
|---|---|---|---|
| TLS / connection (pooled) | 0.4 ms | 12 ms | 8.6% |
| Load balancer + gateway | 1.2 ms | 9 ms | 6.5% |
| Auth lookup | 0.3 ms | 4 ms | 2.9% |
| Validation | 0.5 ms | 2 ms | 1.4% |
| Tokenisation / preprocessing | 2.1 ms | 6 ms | 4.3% |
| Queue wait | 8.0 ms | 61 ms | 43.9% |
| Model forward pass | 24.0 ms | 38 ms | 27.3% |
| Postprocessing | 1.5 ms | 4 ms | 2.9% |
| Response serialisation | 0.8 ms | 3 ms | 2.2% |
| Total | 38.8 ms | 139 ms | 100% |
Now redo the three-week quantisation with these numbers. Halving the forward pass takes p50 from 38.8 to 26.8 ms (31% better) and p99 from 139 to 120 ms (14% better). Dropping utilisation from 91% to 70% by adding replicas changes queue wait from a ρ/(1−ρ) factor of 10.1 to 2.3 — roughly 61 ms down to 14 ms — taking p99 to about 92 ms, and combined with quantisation to about 73 ms. The cheap fix was worth more than the hard one, and only the measurement could tell you.
Optimise the stage with the largest share of p99, not the stage you find most interesting. A 2× speedup of a stage that is 27% of your tail buys you 14%.
Why the tail is the number that matters
Suppose a page makes 10 independent calls to your API and renders only when all have returned. The probability that none of them lands in the p99 tail is 0.9910=0.9044. So 9.6% of page loads contain at least one p99 request, and the page's own p90 is your API's p99. Fan-out converts a rare tail into a common experience. This is why tail latency, not average latency, is the SLO that customers feel.
Fan-out converts a rare tail into a common experience: with ten parallel calls, your p99 becomes the page's p90.
Cutting latency
Model-level optimisations
| Technique | Typical speedup | Memory change | Accuracy cost | Effort |
|---|---|---|---|---|
| fp32 → fp16 / bf16 | 1.5–2× | ÷2 | Negligible | One flag |
| Post-training int8 quantisation | 2–4× | ÷4 | 0.5–2% on most classification tasks | Calibration set, a few hours |
| Quantisation-aware training | 2–4× | ÷4 | <0.5% | A retraining run |
| ONNX Runtime / graph fusion | 1.3–2.5× | Unchanged | None | Export + validate |
| Structured pruning | 1.2–2× | Reduced | 1–3% | Fine-tuning required |
| Distillation to a smaller model | 3–10× | Much smaller | 2–6% | Full training pipeline |
Why quantisation works is worth understanding, because it explains when it will not help. At small batch sizes a transformer forward pass is memory-bandwidth-bound: the GPU spends most of its time reading weights, not multiplying them. BERT-base has 110 million parameters. In fp32 that is 440 MB; on a T4 with about 320 GB/s of memory bandwidth, simply reading the weights takes 0.440/320=1.375 ms. In int8 the weights are 110 MB and the read takes 0.34 ms. The ceiling on the speedup is therefore roughly 4×, and you approach it only while you remain bandwidth-bound. At batch 64 the same model is compute-bound and int8 buys far less on hardware without dedicated int8 tensor cores.
Graph-level export is the cheapest win on the table and is routinely skipped:
1import torch, onnxruntime as ort23batch, seq = torch.export.Dim("batch"), torch.export.Dim("seq")4torch.onnx.export( # dynamo-based exporter, the default in PyTorch 2.9+5 model, (sample_ids, sample_mask), "model.onnx",6 input_names=["input_ids", "attention_mask"], output_names=["logits"],7 dynamic_shapes=({0: batch, 1: seq}, {0: batch, 1: seq}),8)910sess = ort.InferenceSession(11 "model.onnx",12 providers=[("CUDAExecutionProvider", {"cudnn_conv_algo_search": "EXHAUSTIVE"}),13 "CPUExecutionProvider"],14)15logits = sess.run(["logits"], {"input_ids": ids, "attention_mask": mask})[0]The gain comes from operator fusion — a LayerNorm followed by GELU becomes one kernel instead of five, eliminating four round trips to GPU memory. Always re-validate accuracy after export: a mismatched opset can silently change a padding behaviour and cost you two points of F1 with no error message. Leave the opset at the exporter's default unless your runtime needs an older one; forcing an old opset makes the exporter down-convert the graph, which can fail. Use a sample batch of at least two, or the exporter may treat a batch of one as a fixed size.
System-level latency wins
- Connection reuse. A fresh TCP plus TLS handshake costs about three round trips: roughly 4.5 ms in-region, 90 ms across continents. Reusing a pooled keep-alive connection removes it entirely. This is often the single largest saving for a client that is far away.
- Speculative retry (hedging). If p99 is far above p50, send a second request to a different replica when the first exceeds p95, and take whichever returns first. At a p95 threshold this adds only 5% extra load and can pull p99 close to p95. Only safe for idempotent operations.
- Do the work concurrently. If a request needs an auth lookup (4 ms) and a feature fetch (9 ms), running them in parallel costs 9 ms instead of 13.
- Stream when the output is generated. For a generative model producing 28 tokens/s with 210 ms to first token, a 300-token answer takes 210+300/28×1000=10,924 ms in total — but the user sees text at 210 ms. Streaming changes perceived latency by a factor of fifty without changing a single millisecond of compute.
Raising throughput
Horizontal scaling is the blunt instrument, and it works: capacity is linear in replicas, provided the load balancer spreads work by outstanding requests rather than round robin, and provided nothing shared behind them (a database, a Redis instance, a feature store) becomes the new bottleneck.
Batching is the efficient instrument. Model a forward pass as T(b)=6+1.2b ms — a fixed per-batch overhead plus a marginal per-item cost. Throughput b/T(b) goes from 139 req/s at b=1 to 635 req/s at b=16 to 721 req/s at b=32. Between 16 and 32 you gain 14% of throughput for 76% more per-batch latency, which is where you should usually stop.
Pipelining overlaps stages that use different resources. While the GPU runs batch n, the CPU tokenises batch n+1 and the network sends results for batch n−1. If tokenisation is 2.1 ms and inference 24 ms, serial execution is 26.1 ms per batch and pipelined execution is 24 ms — an 8% throughput gain for free, larger when preprocessing is expensive (image decoding often is).
serial: [tok n][infer n][ser n][tok n+1][infer n+1][ser n+1] |------------ 26.1 ms ------------| pipelined: [tok n ][tok n+1][tok n+2] [infer n ][infer n+1][infer n+2] [ser n ][ser n+1 ] |-- 24 ms --| steady state = slowest stageConnection pooling matters for the same Little's Law reason as thread pools. Size the pool at a small multiple of λW for the downstream call, then stop. An oversized pool does not increase throughput; it moves congestion downstream and turns a clean queue in your service into a resource exhaustion event in your database.
Caching
Caching is the only technique on this page that improves latency, throughput and cost simultaneously. It is also the one with the most ways to be subtly wrong.
How much can you expect to hit?
Real request distributions are heavily skewed. Under a Zipf distribution with exponent 1 — a decent model for search queries, product views and prompt templates — the fraction of traffic covered by the top k of N items is approximately Hk/HN, where Hn≈lnn+0.5772.
For N=1,000,000 distinct inputs, caching the top 1,000:
Caching 0.1% of the keyspace serves 52% of traffic. Extend to the top 10,000 (1% of keys) and you get (9.788)/(14.393)=68%. This is why a small cache is usually enough and why doubling cache size past the head yields so little.
Translate a 52% hit rate into effect. On a 45 ms model with a 2 ms cache read: effective latency is 0.52×2+0.48×45=1.04+21.6=22.6 ms, half the original. GPU load falls 52%, so a 10-replica fleet becomes 5.
What to cache
| Layer | Key | Where it lives | TTL | Main risk |
|---|---|---|---|---|
| Full query result | hash(normalised input + model version) | Redis | Hours | Stale after model change if version omitted |
| Embeddings | hash(doc id + encoder version) | Redis or vector store | Until re-embedded | Storage growth |
| Partial / intermediate | hash(preprocessed tensor) | Local memory | Minutes | Memory pressure per replica |
| Prompt prefix (KV cache) | hash(shared prefix tokens) | GPU memory | Seconds–minutes | Evicts under concurrency |
| Auth / config | api key hash | Local, 60 s TTL | 60 s | Revocation delay |
Local versus distributed, and why you want both
request │ ▼ ┌─────────────────┐ hit: 0.01 ms │ L1: in-process │──────────────► response │ LRU, 10k items │ └────────┬────────┘ miss ▼ ┌─────────────────┐ hit: 0.9 ms │ L2: Redis │──────────────► response (and fill L1) │ shared, 5M items│ └────────┬────────┘ miss ▼ ┌─────────────────┐ 45 ms │ model inference │──────────────► response (and fill L2, L1) └─────────────────┘With a 30% L1 hit rate and a further 25% caught by L2, effective latency is 0.30×0.01+0.25×0.9+0.45×45=0.003+0.225+20.25=20.5 ms. A local cache alone would give 0.30 × 0.01 + 0.70 × 45 = 31.5 ms; Redis alone (55% hit) gives 0.55 × 0.9 + 0.45 × 45 = 20.75 ms. The two-tier arrangement wins on latency and, because L1 absorbs the hottest keys, it also takes 30% of lookups — the hottest ones — off Redis, which matters when Redis becomes your next bottleneck.
Invalidation, done the way that works
Do not invalidate on model deploy. Put the model version in the cache key. Deploying v8 makes every v7 key unreachable and it ages out on its own; rolling back to v7 finds its cache still warm. A flush-on-deploy strategy sends 100% of traffic to cold GPUs at the exact moment you are also rolling pods — which is how a routine release becomes an incident.
For content that genuinely changes, prefer short TTLs plus a version stamp on the underlying entity over explicit purges, because purge messages get lost and leave permanent inconsistency. And guard against the thundering herd: when a hot key expires, hundreds of concurrent requests all miss and all call the model. The fix is a per-key lock — one request recomputes, the others wait briefly for the fill — or serving the stale value while a single background refresh runs.
Warm-up, with the arithmetic
A fleet sized for a 52% hit rate is sized for 48% of raw traffic reaching the model. Start with an empty cache and 100% reaches the model — that is 1/0.48=2.08× the load it was provisioned for. If you provisioned at 60% utilisation, a cold start puts you at 0.60×2.08=1.25, or 125% of capacity: the queue grows without bound and the service falls over while the dashboards say "deployment successful".
So warm before you shift traffic. Replay the top few thousand keys from yesterday's logs against the new replicas before marking them ready, or ramp traffic in over several minutes so the cache fills as the load grows.
Measuring the right things
| Metric | What it tells you | Alert when |
|---|---|---|
| p50 / p95 / p99 per stage | Where the time goes | Any stage's p99 share moves by 10 points |
| Throughput and batch-size histogram | Whether batching is working | Median batch size drops (queueing is starving) |
| Utilisation ρ | How close to the 1/(1−ρ) cliff | Sustained above 0.75 |
| Cache hit rate by tier | Whether cost savings are real | Falls more than 10 points from baseline |
| Queue depth | Unmet demand right now | Depth per replica above target for 2 minutes |
| Timeout and 5xx rate | Whether the tail is turning into errors | Any sustained rise |
Two measurement mistakes are near-universal. Averaging percentiles across replicas is arithmetically meaningless — the mean of ten p99 values is not the fleet p99. Aggregate histograms and compute the percentile once. Measuring latency inside the server misses queue time, connection setup and the load balancer entirely; the server reports 24 ms while the client sees 139 ms, and the resulting argument wastes a week. Measure at the client, or at the edge, and break it down inward.
Working through a real budget
Suppose you are given a target: p99 under 100 ms at 600 req/s. Work backwards.
Fixed costs you cannot avoid — gateway, auth, validation, serialisation — total about 27 ms at p99 from the table above. That leaves 73 ms for queueing plus inference. Batched inference at b=13 costs 21.4 ms of execution, and its p99 is around 34 ms. So the queue must stay under roughly 39 ms at p99, which with an exponential-ish tail implies an average wait near 8 ms, which implies ρ ≈ 0.6. At 600 req/s with 600 req/s of capacity per batched replica, that means running two replicas rather than one — and the second replica is not headroom, it is the thing that buys the SLO.
Then add a cache. At a 52% hit rate, 288 req/s reach the model instead of 600, so across the two replicas ρ falls to about 0.24, queue wait at p99 drops to a few milliseconds, and p99 lands near 65 ms with room to spare. The cache did more for the tail than any model optimisation would have.
The lesson to carry into your own systems is to write the budget down as a table with a row per stage and a p99 column, before touching code. It takes an hour and it tells you exactly which of the dozen techniques above is worth your next three weeks — which, in the opening scenario, would have been two replicas and a Redis instance rather than a quantisation project.