Tail latency, percentiles and queueing
Why p99 matters more than the average, how fan-out amplifies tail latency, Little's law and queueing near full utilisation, hedged requests, timeouts and deadlines, and how to set and measure latency budgets.
Reading is half of it. See this used in a real interview: walk through Design a News Feed →
"Average latency is 40 ms" sounds healthy until you learn that one request in a hundred takes two seconds, and that your page makes fifty requests. Interviewers who ask about latency want percentiles, an understanding of why the slow tail grows with scale, and concrete techniques to tame it.
Percentiles, not averages
- p50 (median): what a typical request sees.
- p99: 1 in 100 requests is slower than this.
- p99.9: 1 in 1,000.
Averages hide the tail: a service with 99 requests at 10 ms and one at 2 s averages 30 ms. Set targets as percentiles ("p99 under 200 ms") and measure them with histograms, never by averaging percentiles across servers (percentiles do not average; merge histograms instead).
Fan-out makes the tail everyone’s problem
If one backend call is slow 1 % of the time, a request that waits for 100 such calls in parallel is slow when any of them is:
P(at least one slow) = 1 − 0.99^100 ≈ 63 %
So a backend with a fine-looking p99 makes most fan-out requests slow. This is why search, feeds and any scatter-gather design care about the p99 and p99.9 of their leaves, a point Google’s "The Tail at Scale" made famous.
Where tail latency comes from
- Queueing: requests wait behind others, especially at high utilisation.
- Garbage collection pauses, background compaction, log rotation.
- Noisy neighbours on shared hardware.
- Cache misses: a 99 % hit rate means 1 % of requests take the slow path.
- Retries and timeouts: a request that timed out once and succeeded on retry took timeout + retry time.
- Cold starts after deploys or scaling.
- Network: packet loss and TCP retransmissions add hundreds of milliseconds.
Queueing: why utilisation matters
Little’s law: items in a system = arrival rate × time each spends in it. At 1,000 requests per second and 50 ms each, about 50 requests are in flight; size pools and concurrency limits from it.
Latency explodes near full utilisation. In a simple queueing model, waiting time grows roughly with utilisation / (1 − utilisation): at 50 % utilisation a request waits about as long as its service time; at 90 % about nine times as long; at 99 %, ninety-nine times. That is why services run at 50 to 70 % of capacity, and why a small traffic increase near saturation causes a large latency jump. See autoscaling and capacity planning.
Techniques that cut the tail
- Hedged requests: send the request to one replica; if no answer within (say) the p95 time, send a second copy to another replica and take whichever answers first. Costs a few percent extra load, cuts the tail sharply. Only for idempotent reads.
- Tied requests: send to two replicas at once, and cancel the second as soon as one starts processing.
- Timeouts and deadlines: give each request a deadline that is passed down to every call, so work for an already-failed request is abandoned instead of queueing behind live ones.
- Load shedding: reject early (fast 503) when queues are long, instead of accepting work you will finish too late. See rate limiting and resilience.
- Fewer, smaller hops: every sequential call adds its tail; parallelise independent calls and collapse chatty ones.
- Cache hot paths, and protect against stampedes so misses do not pile up. See caching.
- Isolate slow work: separate pools for slow and fast requests so one does not block the other.
- Partial results: in scatter-gather, return what arrived by the deadline and mark the response as partial (search engines do this).
Latency budgets
Split the end-to-end target across hops. For a 200 ms p99 page:
| Stage | Budget |
|---|---|
| Network and TLS (edge-terminated) | 40 ms |
| Gateway and auth | 10 ms |
| Fan-out to services (parallel) | 100 ms |
| Rendering and serialisation | 20 ms |
| Slack | 30 ms |
Remember that the p99s of sequential stages do not simply add (the slow cases rarely align), but the budget forces each team to know its share. Compare with the latency numbers in the estimation guide.
Measuring properly
- Record latency as histograms at the client, not just at the server; client-side includes queueing and network.
- Beware coordinated omission: load generators that wait for each response before sending the next under-report the tail. Use open-model load testing at a fixed arrival rate.
- Break down by endpoint, region and dependency; a global p99 hides the region that is on fire.
In the interview
State latency targets as percentiles, point out where fan-out amplifies the tail, keep utilisation moderate, and name two or three tail-cutting techniques that fit the design (hedged reads for replicas, deadlines, load shedding, caching).
Checklist
- Targets as p50 and p99 (p99.9 for fan-out leaves).
- Fan-out width and its effect on the tail.
- Utilisation headroom and Little’s law for pool sizes.
- Deadlines propagated to every call.
- Hedged requests or partial results where they fit.
- A latency budget per stage, measured with histograms.