AI System Design and Architecture

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.

Where the 139 ms p99 actually goesPostprocessout — 7 msModelforward — 38 msQueue wait — 61 msTokenise — 6 msLB, auth,validate — 15 msTLSconnection — 12 mstopbottomQuantising the forward pass from 24 ms to 11 ms moves the p99 from 139 ms to 120 ms — about 14 percent for three weeks.
Three weeks went into a stage that was a quarter of the tail; the queue was costing more than the model ever did.

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:

L=λWL = \lambda W

where LL is the average number of requests in the system, λ\lambda is arrival rate, and WW 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=24L = 600 \times 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=2008/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.8600 \times 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:

TechniqueEffect on latencyEffect on throughput
Larger batchesWorse (wait for fill)Much better
More replicasBetter (less queueing)Better
QuantisationBetterBetter
CachingMuch better on hitsBetter (offloads compute)
Running at higher utilisationMuch worseMarginally better
Adding a retry layerBetter 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:

Stagep50p99Share of p99
TLS / connection (pooled)0.4 ms12 ms8.6%
Load balancer + gateway1.2 ms9 ms6.5%
Auth lookup0.3 ms4 ms2.9%
Validation0.5 ms2 ms1.4%
Tokenisation / preprocessing2.1 ms6 ms4.3%
Queue wait8.0 ms61 ms43.9%
Model forward pass24.0 ms38 ms27.3%
Postprocessing1.5 ms4 ms2.9%
Response serialisation0.8 ms3 ms2.2%
Total38.8 ms139 ms100%

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−ρ)\rho/(1-\rho) 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.90440.99^{10} = 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

TechniqueTypical speedupMemory changeAccuracy costEffort
fp32 → fp16 / bf161.5–2×÷2NegligibleOne flag
Post-training int8 quantisation2–4×÷40.5–2% on most classification tasksCalibration set, a few hours
Quantisation-aware training2–4×÷4<0.5%A retraining run
ONNX Runtime / graph fusion1.3–2.5×UnchangedNoneExport + validate
Structured pruning1.2–2×Reduced1–3%Fine-tuning required
Distillation to a smaller model3–10×Much smaller2–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.3750.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:

Python
import torch, onnxruntime as ortbatch, seq = torch.export.Dim("batch"), torch.export.Dim("seq")torch.onnx.export(                  # dynamo-based exporter, the default in PyTorch 2.9+    model, (sample_ids, sample_mask), "model.onnx",    input_names=["input_ids", "attention_mask"], output_names=["logits"],    dynamic_shapes=({0: batch, 1: seq}, {0: batch, 1: seq}),)sess = ort.InferenceSession(    "model.onnx",    providers=[("CUDAExecutionProvider", {"cudnn_conv_algo_search": "EXHAUSTIVE"}),               "CPUExecutionProvider"],)logits = 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,924210 + 300/28 \times 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.2bT(b) = 6 + 1.2b ms — a fixed per-batch overhead plus a marginal per-item cost. Throughput b/T(b)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).

Text
  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 stage

Connection pooling matters for the same Little's Law reason as thread pools. Size the pool at a small multiple of λW\lambda 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 kk of NN items is approximately Hk/HNH_k/H_N, where Hn≈ln⁡n+0.5772H_n \approx \ln n + 0.5772.

For N=1,000,000N = 1{,}000{,}000 distinct inputs, caching the top 1,000:

H1000H1,000,000=ln⁡1000+0.5772ln⁡106+0.5772=7.48514.393=0.52\frac{H_{1000}}{H_{1{,}000{,}000}} = \frac{\ln 1000 + 0.5772}{\ln 10^6 + 0.5772} = \frac{7.485}{14.393} = 0.52

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%(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.60.52 \times 2 + 0.48 \times 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

LayerKeyWhere it livesTTLMain risk
Full query resulthash(normalised input + model version)RedisHoursStale after model change if version omitted
Embeddingshash(doc id + encoder version)Redis or vector storeUntil re-embeddedStorage growth
Partial / intermediatehash(preprocessed tensor)Local memoryMinutesMemory pressure per replica
Prompt prefix (KV cache)hash(shared prefix tokens)GPU memorySeconds–minutesEvicts under concurrency
Auth / configapi key hashLocal, 60 s TTL60 sRevocation delay

Local versus distributed, and why you want both

Text
   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.50.30 \times 0.01 + 0.25 \times 0.9 + 0.45 \times 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×1/0.48 = 2.08\times the load it was provisioned for. If you provisioned at 60% utilisation, a cold start puts you at 0.60×2.08=1.250.60 \times 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

MetricWhat it tells youAlert when
p50 / p95 / p99 per stageWhere the time goesAny stage's p99 share moves by 10 points
Throughput and batch-size histogramWhether batching is workingMedian batch size drops (queueing is starving)
Utilisation ρHow close to the 1/(1−ρ) cliffSustained above 0.75
Cache hit rate by tierWhether cost savings are realFalls more than 10 points from baseline
Queue depthUnmet demand right nowDepth per replica above target for 2 minutes
Timeout and 5xx rateWhether the tail is turning into errorsAny 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.