Your country

Tools that support it use your country for local currency, number formats, units and paper size. Your choice is saved only in this browser.

Type a name or a two-letter code. Use the up and down arrow keys to move through the countries, Enter to choose one and Escape to close.

System Design (High-Level Design)  Module 3 – Performance and reliability fundamentals

Tail latency at scale: fan-out and hedging

See how fan-out turns rare slow replies into slow requests, cut the tail with hedged and tied requests safely, and know when to remove variance instead.

  • Advanced
  • 25 minutes
  • Examples run with Python 3.14.8, Pyodide 314.0.7, Node.js 24.21.0 and quickjs 0.32.0
  • By MySmartCoPilot

What you will learn

  • Calculate how fan-out turns rare slow responses into common slow requests
  • Apply hedged and tied requests safely, with a delay and a budget
  • Choose where to spend effort, removing variance or tolerating it

Before you start

On this page

Tail latency is the latency of the slowest requests, the p99 and beyond. It matters most in systems that split one request across many servers, because a request that waits for all of them is as slow as the slowest. If each server is slow on 1 call in 100, a request that waits for 100 servers is slow 63% of the time. Two kinds of fix exist: remove the causes of slow calls, or tolerate them by sending a second copy of a slow call to another replica (hedged requests) or by queueing it on two servers at once (tied requests). This lesson works out the numbers and shows how to use both fixes without overloading the servers they are meant to help.

Fan-out turns a rare slow server into a common slow request

Search, feeds and analytics queries often fan out: a root server sends one request to many leaf servers, each holding part of the data, and merges their answers. The request finishes when its last leaf answers. If a leaf is slow with probability p, independently of the others, a request that waits for N leaves is slow with probability 1−(1−p)N1 - (1 - p)^N:

Fan-out against how often each server is slow Python · fanout.py
"""How fan-out turns a rare slow server into a common slow request.

A request that waits for N servers is slow whenever at least one of them is slow. If each server is slow with
probability p, independently of the others, the request is slow with probability 1 - (1 - p) ** N.
"""

import math

FANOUTS = [1, 10, 50, 100, 500, 2_000]
SLOW_SHARES = [0.01, 0.001, 0.0001]  # one call in 100, in 1,000 and in 10,000 is slow

print("Share of requests that wait for at least one slow server")
print(f"{'servers':>8}" + "".join(f"{f'1 in {round(1 / p):,}':>13}" for p in SLOW_SHARES))
for n in FANOUTS:
    print(f"{n:>8,}" + "".join(f"{1 - (1 - p) ** n:>13.2%}" for p in SLOW_SHARES))

# To keep the request's p99 at some latency T, every server must be faster than T with probability
# 0.99 ** (1 / N): the request is fast only when all N servers are.
print("\nWhat each server must achieve for the request's p99 to stay at T")
print(f"{'servers':>8}  each server faster than T on this share of calls")
for n in FANOUTS:
    need = 0.99 ** (1 / n)
    one_in = 1 / (1 - need)
    one_in = round(one_in, 2 - math.floor(math.log10(one_in)))  # three significant figures
    print(f"{n:>8,}  {need:.4%}  (at most about 1 slow call in {one_in:,.0f})")

Output

Share of requests that wait for at least one slow server
 servers     1 in 100   1 in 1,000  1 in 10,000
       1        1.00%        0.10%        0.01%
      10        9.56%        1.00%        0.10%
      50       39.50%        4.88%        0.50%
     100       63.40%        9.52%        1.00%
     500       99.34%       39.36%        4.88%
   2,000      100.00%       86.48%       18.13%

What each server must achieve for the request's p99 to stay at T
 servers  each server faster than T on this share of calls
       1  99.0000%  (at most about 1 slow call in 100)
      10  99.8995%  (at most about 1 slow call in 995)
      50  99.9799%  (at most about 1 slow call in 4,980)
     100  99.9900%  (at most about 1 slow call in 9,950)
     500  99.9980%  (at most about 1 slow call in 49,800)
   2,000  99.9995%  (at most about 1 slow call in 199,000)

Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 fanout.py

The first table is the argument of Jeffrey Dean and Luiz André Barroso’s article “The Tail at Scale”, and it reproduces their two figures: with 100 servers that are each slow 1% of the time, 63.4% of requests wait for a slow one, and even if only 1 call in 10,000 is slow, a request that waits for 2,000 servers meets one 18.1% of the time. The article also gives a measurement from a real fan-out service at Google: the p99 for one random leaf to answer was 10 ms, the p99 for 95% of the leaves was 70 ms, and the p99 for all of them was 140 ms. The slowest 5% of leaves cost half of the request’s p99.

The second table turns the argument into a requirement. For the request’s p99 to stay at some latency T with 100 leaves, each leaf must beat T on 99.99% of its calls, which is a p99.99 target for every leaf. Tail targets get tighter as fan-out grows, and at some size no amount of tuning meets them. That is when tolerating slow calls becomes cheaper than preventing them.

Where slow calls come from

The article lists the usual causes, and none of them is the request’s own work: other programs on the same machine competing for its cores, caches and memory bandwidth; background daemons that wake up for a few milliseconds; programs on other machines competing for the same switches or file systems; periodic chores such as log compaction and garbage collection; queues at every layer; processors that slow down when they run hot; flash drives that reorganise themselves internally; and devices that take time to wake from power saving.

Removing variance at its source is worth doing first, because it helps every request. The article’s examples are to give interactive requests their own priority queues instead of one deep queue, to cut long requests into pieces so that a few expensive queries cannot hold up many cheap ones, and to run background work such as compactions in short bursts at the same moment on every machine of a large fan-out, so that it slows the few requests during the burst instead of a few machines always being slow. It also notes that caches do not directly fix the tail unless the whole working set fits in them. But some variance always remains in shared machines, so large systems also learn to live with it.

Hedged requests: a second copy for slow calls

A hedged request sends the call to one replica and, if no answer has come after a short delay, sends the same call to a second replica, uses whichever reply arrives first and cancels the other.

A hedged request: the call goes to replica A; only if it has not answered by the delay is a copy sent to replica B.Send the callto replica AAnswered withinthe hedge delay?Use A's reply(about 95% of calls)Send a copyto replica BUse whicheverreply comes firstCancel theother copyyesno

A hedged request, with the delay set near the p95

Text description of the diagram

The diagram is a flowchart from top to bottom.

  1. The client sends the call to replica A.
  2. It asks: has A answered within the hedge delay?
  3. If yes, the client uses A's reply. With the delay set at the p95 of the call's latency, this is what happens to about 95% of calls.
  4. If no, the client sends a copy of the same call to replica B, uses whichever reply comes first, and cancels the other copy.

Only the slow calls get a copy, which is why a delay near the p95 adds only about 5% to the load.

Sending every call twice would double the load. The article’s rule is to wait until the call has been outstanding for longer than its p95 latency, which limits the extra load to about 5% while still shortening the tail a lot, because a slow call is often slow through interference on its server, not through the request itself. In a Google benchmark that read 1,000 keys spread over 100 servers, hedging after 10 ms cut the p99.9 of the whole read from 1,800 ms to 74 ms while sending only 2% more requests.

The simulation below hedges each leaf call of a fan-out of 100, with leaf calls that usually take about 4 ms and 2% of calls that meet 20 to 80 ms of interference:

Hedging the leaf calls of 2,000 fan-out requests Python · hedge_sim.py
"""Hedged requests in a fan-out of 100 leaf calls, simulated with a fixed seed.

Each leaf call usually takes about 4 ms; 2% of calls run into interference on their server (a compaction, a
garbage-collection pause, a noisy neighbour) and take 20 to 80 ms longer. A hedge sends a second copy of a leaf
call to another replica when the first has not answered after `delay`, and uses whichever reply comes first.
Every number here is an assumption for the sketch; change them and run it again.
"""

import math
import random

LEAVES, REQUESTS = 100, 2_000
rng = random.Random(5)


def base_ms():
    return rng.lognormvariate(math.log(4), 0.35)


def interference_ms():
    return rng.uniform(20, 80) if rng.random() < 0.02 else 0.0


def percentile(values, p):
    ordered = sorted(values)
    return ordered[math.ceil(round(p * len(ordered) / 100, 6)) - 1]


# Calibrate: the leaf latency percentiles a client would know from its own measurements.
calibration = [base_ms() + interference_ms() for _ in range(50_000)]
leaf = {p: percentile(calibration, p) for p in (50, 90, 95, 99)}

# The primary calls of every request, drawn once so that every strategy sees the same requests.
requests = [[(base_ms(), interference_ms()) for _ in range(LEAVES)] for _ in range(REQUESTS)]


def run(delay, correlated=False):
    """Request latencies and the share of leaf calls that sent a hedge."""
    latencies, hedges = [], 0
    for calls in requests:
        slowest = 0.0
        for base, extra in calls:
            done = base + extra
            if delay is not None and done > delay:
                hedges += 1
                # The copy goes to another replica: independent interference, unless the cause is shared.
                backup = base_ms() + (extra if correlated else interference_ms())
                done = min(done, delay + backup)
            slowest = max(slowest, done)
        latencies.append(slowest)
    return latencies, hedges / (REQUESTS * LEAVES)


print(f"leaf calls: p50 {leaf[50]:.1f} ms, p90 {leaf[90]:.1f} ms, p95 {leaf[95]:.1f} ms, p99 {leaf[99]:.1f} ms")
print(f"\n{REQUESTS:,} requests of {LEAVES} leaf calls each")
print(f"{'strategy':36}{'extra calls':>12}{'request p50':>13}{'request p99':>13}")
rows = [("no hedging", None, False)]
rows += [(f"hedge after the leaf p{p} ({leaf[p]:.1f} ms)", leaf[p], False) for p in (50, 90, 95, 99)]
rows += [("p95 hedge, slowness shared", leaf[95], True)]
for name, delay, correlated in rows:
    latencies, extra = run(delay, correlated)
    print(f"{name:36}{extra:>12.1%}{percentile(latencies, 50):>10.1f} ms{percentile(latencies, 99):>10.1f} ms")

Output

leaf calls: p50 4.0 ms, p90 6.5 ms, p95 7.7 ms, p99 56.3 ms

2,000 requests of 100 leaf calls each
strategy                             extra calls  request p50  request p99
no hedging                                  0.0%      63.4 ms      84.9 ms
hedge after the leaf p50 (4.0 ms)          50.2%       9.2 ms      53.9 ms
hedge after the leaf p90 (6.5 ms)          10.0%      11.3 ms      55.5 ms
hedge after the leaf p95 (7.7 ms)           5.0%      12.4 ms      57.8 ms
hedge after the leaf p99 (56.3 ms)          0.9%      59.2 ms      68.9 ms
p95 hedge, slowness shared                  5.0%      63.4 ms      84.9 ms

Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 hedge_sim.py

Without hedging, most requests contain at least one interfered leaf, so even the request’s median is 63.4 ms. A hedge at the leaf p95 of 7.7 ms sends 5.0% extra calls and brings the median down to 12.4 ms and the p99 from 84.9 ms to 57.8 ms. Hedging at the p50 makes half of all calls send a copy, for a small further gain; hedging at the p99 is cheap but too late to help most requests. The p99 does not drop further because some leaves are unlucky twice, on both replicas. The last row is the important caveat: when slowness has a shared cause, so that the copy meets the same interference as the original, hedging adds 5% load and changes nothing.

Tied requests: queue twice, run once

A hedge still runs the slow call twice whenever it was merely queued. A tied request goes to two servers at once, each copy tagged with the other server’s name, and the first server to start working on it tells the other to drop its copy. The client waits a moment, about twice the network delay, before sending the second copy, so that two idle servers do not both start. The article reports small BigTable reads that had to go to the underlying file system, which keeps three replicas of each chunk: on a mostly idle cluster, tying a request to a second replica after 1 ms cut the median from 19 to 16 ms and the p99.9 from 98 to 61 ms, and with a large sorting job competing for the disks the p99.9 fell from 159 to 108 ms, for less than 1% more disk use.

Hedging safely

Three rules keep hedges from turning a latency fix into an outage:

  • Hedge only what is safe to run twice. Reads are; a payment, an append or a counter increment is not, unless the server deduplicates it with an idempotency key. The article notes that its techniques suit reads, and that writes usually tolerate latency more easily: they are fewer, can often finish after the reply, and replicated writes through quorum protocols wait for only three to five replicas.
  • Cancel the loser. A copy that runs to the end after the first reply is wasted work on a server that is already slow.
  • Give hedges a budget. A fixed delay sends more copies exactly when the service is slow everywhere, which is when extra load hurts most. This program shows the difference:
A hedged call with and without a budget JavaScript · hedge.mjs
// A hedged call with a budget, tested against a fake service in simulated time (so every run prints the same).
// The hedger sends a copy of a call to a second replica when the first has not answered after `delayMs`, uses the
// first reply, and cancels the other. Its budget is a bucket of tokens: every call adds `ratio` of a token, every
// hedge spends a whole one, so hedges stay below about `ratio` of calls however slow the service gets.

function generator(seed) {
  let a = seed >>> 0;
  return () => {
    a = (a + 0x6d2b79f5) >>> 0;
    let t = Math.imul(a ^ (a >>> 15), a | 1);
    t ^= t + Math.imul(t ^ (t >>> 7), t | 61);
    return ((t ^ (t >>> 14)) >>> 0) / 4294967296;
  };
}
const random = generator(42);

function makeHedger({ delayMs, ratio, burst }) {
  let tokens = burst;
  return function call(primaryMs, backupMs) {
    tokens = Math.min(burst, tokens + ratio);
    if (primaryMs <= delayMs) return { ms: primaryMs, hedged: false }; // answered in time: no copy
    if (tokens < 1) return { ms: primaryMs, hedged: false }; // budget spent: wait for the first replica
    tokens -= 1;
    return { ms: Math.min(primaryMs, delayMs + backupMs), hedged: true }; // first reply wins
  };
}

// A replica's latency for one call: usually 3 to 8 ms; sometimes a pause; when overloaded, often slow.
function replicaMs(overloaded) {
  let ms = 3 + 5 * random();
  if (random() < 0.02) ms += 60 + 60 * random(); // an occasional pause on one replica
  if (overloaded && random() < 0.4) ms += 80 + 70 * random(); // a saturated backend: slow everywhere
  return ms;
}

const p99 = (values) => [...values].sort((a, b) => a - b)[Math.ceil(0.99 * values.length) - 1];

function phase(name, overloaded) {
  const calls = Array.from({ length: 5000 }, () => [replicaMs(overloaded), replicaMs(overloaded)]);
  const policies = [
    ["no hedging", null],
    ["hedge after 10 ms, no budget", makeHedger({ delayMs: 10, ratio: 1, burst: 1 })],
    ["hedge after 10 ms, 5% budget", makeHedger({ delayMs: 10, ratio: 0.05, burst: 10 })],
  ];
  console.log(`${name}: ${calls.length} calls`);
  for (const [label, hedger] of policies) {
    let hedges = 0;
    const latencies = calls.map(([primary, backup]) => {
      if (!hedger) return primary;
      const result = hedger(primary, backup);
      if (result.hedged) hedges += 1;
      return result.ms;
    });
    const load = (calls.length + hedges) / calls.length;
    console.log(`  ${label.padEnd(30)} p99 ${p99(latencies).toFixed(1).padStart(6)} ms, load on replicas x${load.toFixed(2)}`);
  }
}

phase("Normal traffic", false);
phase("Backend saturated", true);

Output

Normal traffic: 5000 calls
  no hedging                     p99   91.1 ms, load on replicas x1.00
  hedge after 10 ms, no budget   p99   14.9 ms, load on replicas x1.02
  hedge after 10 ms, 5% budget   p99   14.9 ms, load on replicas x1.02
Backend saturated: 5000 calls
  no hedging                     p99  156.0 ms, load on replicas x1.00
  hedge after 10 ms, no budget   p99  144.3 ms, load on replicas x1.41
  hedge after 10 ms, 5% budget   p99  155.8 ms, load on replicas x1.05

Recorded with Node.js 24.21.0 on macOS 26 arm64. To run it yourself: mise exec node@24.21.0 -- node hedge.mjs

In normal traffic, hedging after 10 ms cuts the p99 from 91.1 ms to 14.9 ms for 2% more calls, with or without a budget. When the backend is saturated and 40% of calls are slow everywhere, the hedger without a budget adds 41% more calls for a 12 ms gain, and in a real system that extra load would make every call slower still. The budgeted hedger stops at 5%. Libraries build in such limits: gRPC’s hedging policy keeps at most 5 attempts of a call in flight, sent at a configured delay, and its retry throttling counts failed calls against a token bucket and stops sending hedges while the bucket is at or below half full.

Remove the variance, or tolerate it?

Removing variance or tolerating it
Criterion Remove variance at the sourceTolerate it: hedged or tied requests
Helps when The cause is known and shared by many calls: a queue, a background job, a noisy neighbourSlowness is random from server to server
Extra load NoneAbout 5% for hedges after the p95; under 1% for tied requests that cancel fast
Requires Engineering work on the slow componentReplicas, calls that are safe to repeat, and cancellation
Fails when The cause is outside your controlBoth replicas share the cause of the slowness
When to choose First, for every cause you can seeLarge read fan-outs over replicated data, for what remains

The article describes two more tools that work across requests rather than within one. Splitting data into many more micro-partitions than machines, about 20 per machine in its example, lets a system move load in steps of roughly 5% when one machine is slow or busy, and copying popular partitions to more machines spreads a hot spot. A server whose replies are persistently slow can be put on probation: it stops receiving live traffic but keeps getting shadow requests, so it can return once it recovers. Search systems can also answer once most leaves have replied, with “good enough” results, and send a risky new query to one or two leaves first, a canary request, so that a query that crashes servers crashes only those.

Practice

Exercise · Medium · Python

Size the tail of a fan-out and pick a hedge delay

Write three helpers in hedging.py.

slow_request_share(slow_share, fanout) returns the share of requests that wait for at least one slow call, when a request makes fanout calls in parallel and each call is slow with probability slow_share, independently of the others. For example, slow_request_share(0.01, 100) is about 0.634.

extra_load(latencies_ms, delay_ms) returns the share of calls that would send a hedge if the client waited delay_ms before sending a copy: the calls whose latency is strictly greater than the delay. For example, extra_load([1, 2, 3, 4, 10], 3) is 0.4.

pick_delay(latencies_ms, max_extra) returns the smallest delay, chosen from the measured latencies themselves, whose extra load is at most max_extra. A short delay cuts the tail most but sends the most copies, so this is the most aggressive delay the budget allows. For example, with the latencies 1 to 100 ms and a budget of 0.05, the answer is 95.

Starter code · hedging.py

def slow_request_share(slow_share, fanout):
    """Share of requests that wait for at least one slow call among `fanout` parallel calls."""
    # Replace this line with your code.
    return 0


def extra_load(latencies_ms, delay_ms):
    """Share of calls still unanswered after delay_ms (strictly slower than it): the hedges sent."""
    # Replace this line with your code.
    return 0


def pick_delay(latencies_ms, max_extra):
    """The smallest measured latency to use as the delay with extra_load at most max_extra."""
    # Replace this line with your code.
    return 0
The sample tests · test_hedging.py
import math

from hedging import extra_load, pick_delay, slow_request_share


def test_slow_request_share():
    """is one minus the chance that every call is fast"""
    assert math.isclose(slow_request_share(0.01, 100), 0.6340, rel_tol=1e-3)
    assert math.isclose(slow_request_share(0.0001, 2_000), 0.1813, rel_tol=1e-3)
    assert math.isclose(slow_request_share(0.5, 1), 0.5, rel_tol=1e-9)
    assert math.isclose(slow_request_share(0.01, 0), 0.0, abs_tol=1e-12)


def test_extra_load():
    """counts the calls strictly slower than the delay"""
    assert math.isclose(extra_load([1, 2, 3, 4, 10], 3), 0.4)
    assert math.isclose(extra_load([5, 5, 5], 5), 0.0)
    assert math.isclose(extra_load([5, 5, 5], 4), 1.0)


def test_pick_delay():
    """returns the smallest measured latency whose extra load fits the budget"""
    latencies = list(range(1, 101))
    assert pick_delay(latencies, 0.05) == 95
    assert pick_delay(latencies, 0.01) == 99
    assert pick_delay(latencies, 0) == 100
    assert pick_delay([4, 4, 4, 50], 0.25) == 4
    assert pick_delay([50, 4, 4, 4], 0.2) == 50
A hint

The chance that all fanout calls are fast is (1 - slow_share) ** fanout; the request is slow otherwise. For the extra load, count the latencies above the delay and divide by how many there are. For the delay, try the measured latencies from the smallest up (sorted(set(latencies_ms))) and return the first one whose extra load fits the budget.

The sample tests run on this device, in your browser (Pyodide): nothing is sent to mysmartcopilot.com. The first run downloads Python (about 13.5 MB), which is kept for the next runs. A check in your browser is feedback for you, not proof that the code is right for every input.

Check yourself

5 questions about this lesson. Every answer and why it is right is on the page, behind “Show the answer”. Your score stays in this browser.

  1. Question 1 of 5 A search request waits for 100 index servers, and each is independently slow on 1% of calls. On what share of searches does at least one server answer slowly? Answer in percent, to one decimal place.

    Type a number.

    Show the answer to question 1

    Answer: 63.4 % (anything from 63.2 to 63.6 counts)

    All 100 are fast with probability 0.99 to the power 100, about 0.366, so 1 − 0.366 = 63.4% of searches wait for a slow server. A tail of 1% per server becomes the common case per search.

  2. Question 2 of 5 A client hedges a call when it has not answered after the call's measured p95 latency. About how many extra calls does that send?

    Choose one answer.

    Show the answer to question 2

    Answer: About 5% more calls, because only the calls slower than the p95 get a copy

    By definition 5% of calls are still unanswered at the p95, so about 5% of calls send a copy. Dean and Barroso describe this delay as a way to keep the extra load near 5% while shortening the tail a lot.

  3. Question 3 of 5 Which call is safe to hedge without extra protection?

    Choose one answer.

    Show the answer to question 3

    Answer: Reading a product's details by its id from a replicated store

    A hedge runs the same request twice, so only a request that is safe to run twice may be hedged as it is. A read changes nothing. A charge, an append or an increment would happen twice unless the server deduplicates it, for example with an idempotency key.

  4. Question 4 of 5 Both replicas of a shard read from the same overloaded disk array, so when one is slow the other usually is too. What does hedging do for that shard's tail?

    Choose one answer.

    Show the answer to question 4

    Answer: Little, because the copy waits on the same cause; fix the shared disk or move the replicas apart instead

    Hedging works only when slowness on one replica says little about the other. When the cause is shared, the copy is slow too and only adds load; the paper notes that these techniques help only when the source of variability does not tend to hit several replicas at once.

  5. Question 5 of 5 In a tied request, the client queues the same request on two servers, each told about the other. What happens when one of them starts working on it?

    Choose one answer.

    Show the answer to question 5

    Answer: It sends a cancellation to the other server, which drops its queued copy or puts it far back

    A tied request cancels the copy as soon as one server starts on it, so the work is rarely done twice. A small gap between the two sends, about twice the network delay, keeps two idle servers from both starting.

Interview questions

Warm-up (mid to senior level): what is tail latency, and why does it matter more when a request fans out? Tail latency is the latency of the slowest requests, measured as a high percentile such as the p99 or p99.9. It matters more under fan-out because a request that waits for many servers is as slow as its slowest one: if each server is slow on 1% of calls, a request to 100 servers meets a slow one about 63% of the time, so a rare per-server problem becomes the normal experience of users. Reducing it means either removing the causes of slow calls, such as queueing and background work, or tolerating them with techniques such as hedged requests.

Key takeaways

  • A request that waits for N servers, each slow with probability p, is slow with probability 1 − (1 − p)^N: 63% for 100 servers at 1%.
  • Fan-out tightens per-server targets: with 100 leaves, a request p99 needs a p99.99 from every leaf.
  • Remove variance at its source first: priority queues, smaller units of work, coordinated background jobs.
  • Hedge after about the p95 delay for about 5% extra load, only for calls that are safe to run twice, with cancellation and a budget; tied requests cancel the second copy once one server starts.
  • Hedging cannot help when both replicas share the cause of the slowness.

References

Related tools

Report a problem with this lesson

Quick answers and tool search

Type to search tools or to get a quick answer, for example 18% of 2500. Use the up and down arrow keys to move through the results, Enter to choose, and Escape to close.