Why does fan-out make fast services slow?
Short answer: a request that waits for many calls in parallel is as slow as the slowest of them. If one call in a hundred is slow, a request that makes 100 calls meets a slow one 63 % of the time. The rare case of the single call is the normal case of the request.
In the live model, one call to a service with 10 ms of work has a median of about 9 ms and a p99 of about 50 ms. A request that makes 100 such calls in parallel has a median of about 55 ms: its typical case is slower than the slowest 1 % of single calls. What helps is fewer calls per request and a narrower spread of each call, long before more servers.
A search page asks a hundred index shards. A product page asks for the price, the stock, the reviews and twenty recommendations. A feed asks for the posts of fifty accounts. Every one of these services looks healthy on its own dashboard: median 7 ms, p99 50 ms. The page is slow anyway. Jeffrey Dean and Luiz André Barroso described why in “The Tail at Scale” (2013), with numbers from Google's services. This page walks through the arithmetic and then shows the same effect in a design you can open.
The arithmetic: one slow call is enough
Call a response “slow” if it is in the slowest share p of a service's responses. A request that needs N answers and waits for all of them is slow if at least one answer is:
slow share = 1 − (1 − p)^N
| Parallel calls | 1 call in 100 is slow | 1 in 1,000 | 1 in 10,000 |
|---|---|---|---|
| 1 | 1 % | 0.1 % | 0.01 % |
| 10 | 9.6 % | 1 % | 0.1 % |
| 100 | 63 % | 9.5 % | 1 % |
| 1,000 | 100 % | 63 % | 9.5 % |
| 2,000 | 100 % | 86 % | 18 % |
Read the other way round, the table says which part of a service's latency curve its callers live in. The median of a request with 100 parallel calls is the 99.3rd percentile of the single call, and its p99 is the single call's 99.99th percentile. A team that watches the p99 of its service is watching a number that a wide caller left behind long ago: for that caller, the service's p99.99 is what counts.
What Google measured
The paper reports one real service that sends a request through a tree of servers to a very large number of leaves. Measured at the root, a single leaf answers within 10 ms at the 99th percentile. Waiting for all leaves, the 99th percentile is 140 ms. Waiting for 95 % of them, it is 70 ms: half of the whole tail is spent on the slowest twentieth of the answers.
| Measured at the root | Median | 95th percentile | 99th percentile |
|---|---|---|---|
| One random leaf finishes | 1 ms | 5 ms | 10 ms |
| 95 % of all leaves finish | 12 ms | 32 ms | 70 ms |
| All leaves finish | 40 ms | 87 ms | 140 ms |
The same effect in the live model
A design you can open in Stackrig: users send 50 requests a second to a root service, and the root asks a leaf service N times in parallel and waits for all answers. A leaf call needs 10 ms of work on average, with the spread that request handlers usually have: most calls are quicker than the mean and a few take several times as long. The leaf has far more servers than it needs (its CPUs stay below 20 %), so nothing in the table is queueing: it is only the spread of the work itself.
| Parallel calls per request | Median | 95th percentile | 99th percentile |
|---|---|---|---|
| 1 | 9 ms | 30 ms | 50 ms |
| 10 | 27 ms | 62 ms | 100 ms |
| 50 | 45 ms | 95 ms | 135 ms |
| 100 | 55 ms | 110 ms | 160 ms |
Three things to read from it:
- The median moves most. From one call to a hundred, the median grows about six times and the p99 about three times. Fan-out does not only add a bad tail: it makes the ordinary request slow.
- The first ten calls cost the most. Going from 1 to 10 calls triples the median; going from 10 to 100 doubles it. A page that makes “only” ten parallel calls is already most of the way there.
- The leaf's own numbers never changed. In all four rows the leaf service reports the same median and the same p99. Its dashboard cannot show the problem.
The spread of the work is the lever
The same design with 100 parallel calls, and one setting changed on the leaf: its “Service time variation” is 0.3 instead of 1, so that 98 of 100 calls take between about 5 and 20 ms instead of between 1 and 50. The mean stays 10 ms.
| Leaf's service time variation | Median of the request | 99th percentile of the request |
|---|---|---|
| 1 (the default) | 55 ms | 160 ms |
| 0.3 | 22 ms | 30 ms |
No server was added and no call was removed. Making every call take about the same time cut the request's p99 to about a fifth. In a real service that means: the same query for every call instead of a cheap one and an expensive one behind the same endpoint, no stop-the-world pauses, no background job on the machines that answer.
For comparison, the same 100 calls one after the other take about a second (100 × 10 ms): calling in parallel is what makes a wide request possible at all, and the tail is its price.
Have an invite? Open the designs:
In each, select the connection from Root to Leaf to change the calls per request, and the Leaf to change its service time variation.
What the paper recommends
Dean and Barroso's point is that a large service cannot remove every source of slowness (shared machines, background daemons, garbage collection, queueing, maintenance work), so it has to be built to tolerate slow parts, the way it tolerates broken ones. Their techniques, in our words:
- Hedged requests. Send the same call to a second replica if the first has not answered after a short delay, and take whichever answer comes first. Waiting until the 95th percentile of the expected latency keeps the extra load near 5 %. Their example: a benchmark that reads 1,000 keys from a table spread over 100 servers; a second request after 10 ms brought the 99.9th percentile from 1,800 ms to 74 ms with 2 % more requests.
- Tied requests. Queue the call on two servers at once and let each cancel the other's copy when it starts work, so that the queueing delay of one server no longer decides.
- Good-enough answers. Return when most of the leaves have answered, if the result is still useful without the rest. Their own table shows what that buys: 70 instead of 140 ms.
- Take slow servers out for a while, and keep sending them shadow requests to see when they have recovered.
- Many small partitions per machine, so that load can be moved in small steps and hot items can get extra replicas.
- Canary requests. Try a new kind of request on one or two leaves before sending it to thousands, so that one bad request cannot stall all of them.
A hedged request is a retry sent early. It needs calls that are safe to repeat, and a limit on how many are sent, or it turns a slow minute into an overload: Retry storms: why retries take your service down.
What to do in a smaller system
- Count the calls per request. For every page or API call, how many backend calls does it wait for? Ten is already a fan-out.
- Make one call instead of many. One query for a hundred rows, one batch call for a hundred keys: one draw from the latency curve instead of a hundred.
- Watch the percentile your callers feel. A service called 100 times per request needs its p99.9 and p99.99 on the dashboard, not its p99.
- Narrow the spread before adding servers: separate cheap and expensive work, remove pauses. Why queueing widens the spread as machines fill up is in Why does p99 latency explode long before 100 % CPU?
- Decide what the page does without the slowest answer: a time limit per call and a page that renders without the recommendations.
What this page does not tell you
- The arithmetic assumes independent calls. Real slowness is often shared: a pause on one machine slows every call it serves, a slow network slows all of them. Then fewer requests are hit than the table says, and those that are hit are hit harder.
- The shape of the spread is an assumption. The model gives work times a skewed spread described by one number, the service time variation; with the default of 1 the far tail is longer than in the simplest textbook case. Your service has its own curve: measure its p99.9 and p99.99 before trusting any table.
- The model has no hedged or tied requests. It shows the problem and what the number of calls and the spread do to it, not the paper's remedies.
- Google's numbers are Google's. One service, measured before 2013, with far more leaves than most systems have.
- The live model is not a measurement. How it compares with a real machine is answered for one system in the $12 server. What the live model is checked against, an exact simulation, is on How accurate is Stackrig?
Sources
- Jeffrey Dean and Luiz André Barroso, The Tail at Scale, Communications of the ACM 56(2), February 2013, pages 74–80: the 63 % example and the 2,000-server example (page 76), Table 1 (page 77), hedged and tied requests (pages 77–78), the other techniques (pages 78–79).
Open “Web app with database” in the playground
The playground is invite-only during the private preview: join the waitlist to get an invite.