All writing
Distributed Systems6 min read

At fan-out, your p99 becomes your median

A service where one request in a hundred is slow sounds healthy. Fan a single user request out to a hundred such services and 63% of your users wait. The arithmetic is unforgiving, and it is why per-service latency targets quietly fail.

Here is a claim that sounds reassuring in a review: "our p99 is one second, but the typical response is 10ms." One request in a hundred is slow. Ninety-nine percent of the time everything is fast.

Now let that service be one of a hundred your request has to consult before it can answer. The user waits for the slowest one. Suddenly 63% of your users are experiencing that one-second tail.

This is the central observation of The Tail at Scale by Jeffrey Dean and Luiz André Barroso, published in Communications of the ACM in 2013. It is seven pages long, it is not paywalled, and it invalidates most of how teams reason about latency in a microservice architecture.

The arithmetic

If each backend independently exceeds some latency threshold with probability p, and a request must wait on n of them in parallel, then the chance the user escapes the tail entirely is (1−p) raised to n. So:

P(user request is slow) = 1 - (1 - p)^n

  p = per-backend probability of a slow response
  n = number of backends on the critical path

That exponent is doing violent work. Holding a 1-in-100 slow rate constant and varying only the fan-out:

Backends on the pathp = 1% (p99)p = 0.01% (p99.99)
11.0%0.01%
109.6%0.1%
5039.5%0.5%
10063.4%1.0%
20086.6%2.0%
50099.3%4.9%
2,000~100%18.1%
Probability that at least one backend lands in the tail, by fan-out.

The paper states the consequence directly:

If a user request must collect responses from 100 such servers in parallel, then 63% of user requests will take more than one second. Even for services with only one in 10,000 requests experiencing more than one-second latencies at the single-server level, a service with 2,000 such servers will see almost one in five user requests taking more than one second.
Dean & Barroso, The Tail at Scale, CACM 56(2), pp. 74–80

Read that second sentence again. Four nines of per-server latency compliance still leaves a fifth of your users waiting, once the fan-out is wide enough. There is no per-service latency target that saves you; the target has to be set against the fan-out.

What it looks like in production

The paper includes measurements from a real Google service where a root server distributes a request through intermediate servers to a large number of leaves — the same shape as any fan-out query.

Measurementp99 latency
A single random request finishing10 ms
95% of the requests finishing70 ms
All requests finishing140 ms
Measured at the root of a real Google service with large fan-out.

Waiting for the slowest 5% of leaf responses accounts for half of the total 99th-percentile latency. Half your tail is produced by one twentieth of your requests — which is also the good news, because it means work targeted at those outliers has outsized leverage.

Where the variability comes from

Tail latency is rarely a property of the request. It is usually interference, which is why the same request is fast on the next attempt. The paper enumerates the sources: contention for shared resources between applications and within them, background daemons doing periodic maintenance, queueing at multiple layers, garbage collection, energy management transitions, and maintenance activities like log compaction and index rebuilds.

That last observation is what makes the mitigations work. If slowness were inherent to the request, retrying elsewhere would not help. Because it is mostly interference, a second copy of the request sent to a different replica will usually be fast.

Hedged requests

Send the request to one replica. If it has not answered within the 95th-percentile expected latency for that request class, send a second copy to another replica and take whichever returns first. Because only 5% of requests reach the threshold, the additional load is capped at roughly 5%.

The measured result is striking. In a Google benchmark reading 1,000 keys from a BigTable spread across 100 servers, issuing a hedging request after a 10ms delay reduced the 99.9th-percentile latency for retrieving all 1,000 values from 1,800ms to 74ms, while sending just 2% more requests.

A 24× reduction in tail latency for 2% more traffic is the best trade in this entire article, and it requires no changes to the backends.

Tied requests

Hedging still has a window in which two servers do the same work. Tied requests close it: enqueue the request on two servers simultaneously, each tagged with the identity of the other. Whichever server starts executing first sends a cancellation to its counterpart, which drops the copy if it is still sitting in the queue.

This attacks queueing delay specifically, which is a large share of variability, since completion time is far more predictable once a request actually starts running. In Google's distributed file system, tying to a second replica after 1ms reduced median latency by 16% and cut nearly 40% off the 99.9th percentile — with less than 1% overhead in disk utilisation, because cancellation genuinely prevents the redundant read.

The most interesting result is the second scenario: with a large concurrent sorting job competing for the same disks, the latency profile with tied requests nearly matched an idle cluster without them. That is what lets you consolidate workloads instead of over-provisioning dedicated capacity for latency-sensitive services.

Structural fixes

  • Micro-partitions. Generate far more partitions than machines — BigTable runs 20 to 1,000 tablets per machine. With roughly 20 partitions per machine you can shed load in 5% increments and in a twentieth of the time a one-to-one mapping would take. Recovery improves too, since many machines each pick up a small piece of a failed one.
  • Selective replication. Predict the items likely to cause imbalance and make extra copies of them, so hot partitions spread across machines without being moved. Google's web search does this for popular documents.
  • Latency-induced probation. Detect a persistently slow machine and remove it from the serving set, while continuing to send it shadow requests to know when it recovers. Counter-intuitively, removing capacity under high load can improve latency.

What to change on Monday

  1. 01Instrument the fan-out. Record how many backend calls each user request actually makes; most teams are surprised, and you cannot reason about the tail without n.
  2. 02Set leaf targets from the fan-out, not from habit. A service consulted 100 times per user request needs its p99.9 controlled, because user-visible latency samples deep into its tail.
  3. 03Stop reporting averages. An average latency graph cannot show you the thing that determines user experience, and it will look fine throughout the incident.
  4. 04Measure the tail where the user is, at the root, not per-service. Every backend can hit its target while the composed request misses badly.
  5. 05Try hedging on your widest read path first. It is the cheapest intervention with the largest measured effect, and it does not require touching the backends.

The framing that has stayed with me from the paper is its comparison to fault tolerance. We long ago accepted that a reliable system has to be built from unreliable parts, and we built the discipline to do it. Dean and Barroso argue the same is now true of latency: a predictably responsive whole has to be assembled from unpredictable parts. Tail tolerance is engineering work, not a tuning exercise you finish.

Sources

  1. 01Jeffrey Dean and Luiz André Barroso. The Tail at Scale. Communications of the ACM 56(2), pp. 74–80, 2013.