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
Calculated: the share of requests that meet at least one slow call, if the calls are slow independently of each other. The paper gives two of these cells: 63 % for 100 servers at one slow answer in 100, and “almost one in five” for 2,000 servers at one in 10,000.
Parallel calls1 call in 100 is slow1 in 1,0001 in 10,000
11 %0.1 %0.01 %
109.6 %1 %0.1 %
10063 %9.5 %1 %
1,000100 %63 %9.5 %
2,000100 %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.

From Table 1 of Dean and Barroso, The Tail at Scale (2013): finishing times of leaf requests in one of Google's fan-out services, measured at the root.
Measured at the rootMedian95th percentile99th percentile
One random leaf finishes1 ms5 ms10 ms
95 % of all leaves finish12 ms32 ms70 ms
All leaves finish40 ms87 ms140 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.

From Stackrig's live model, not from real machines: latency of the whole request over two minutes, rounded. In the app, p99 is shown for the last second and moves around these values.
Parallel calls per requestMedian95th percentile99th percentile
19 ms30 ms50 ms
1027 ms62 ms100 ms
5045 ms95 ms135 ms
10055 ms110 ms160 ms

Three things to read from it:

  1. 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.
  2. 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.
  3. 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.

From Stackrig's live model: 100 parallel calls per request, the leaf's mean work 10 ms in both rows; rounded.
Leaf's service time variationMedian of the request99th percentile of the request
1 (the default)55 ms160 ms
0.322 ms30 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:

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

  1. 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.
  2. 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.
  3. 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.
  4. 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?
  5. 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

Sources

Open “Web app with database” in the playground

The playground is invite-only during the private preview: join the waitlist to get an invite.

More