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.
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 :
"""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
Runs on this device, in your browser. The first run downloads Python (about 13.5 MB), which is kept for the next runs.
Your run, in this browser
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, with the delay set near the p95
Text description of the diagram
The diagram is a flowchart from top to bottom.
- The client sends the call to replica A.
- It asks: has A answered within the hedge delay?
- 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.
- 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:
"""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
Runs on this device, in your browser. The first run downloads Python (about 13.5 MB), which is kept for the next runs.
Your run, in this browser
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 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
Runs on this device, in your browser. The first run downloads JavaScript (about 0.6 MB), which is kept for the next runs.
Your run, in this browser
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?
| Criterion | Remove variance at the source | Tolerate it: hedged or tied requests |
|---|---|---|
| Helps when | The cause is known and shared by many calls: a queue, a background job, a noisy neighbour | Slowness is random from server to server |
| Extra load | None | About 5% for hedges after the p95; under 1% for tied requests that cancel fast |
| Requires | Engineering work on the slow component | Replicas, calls that are safe to repeat, and cancellation |
| Fails when | The cause is outside your control | Both replicas share the cause of the slowness |
| When to choose | First, for every cause you can see | Large 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.
Results of the sample tests
| Test | Result | Details |
|---|
What your code printed
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.
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
- The Tail at Scale (Jeffrey Dean and Luiz André Barroso) (Communications of the ACM, via Google Research)
- The Tail at Scale (authors' copy, PDF) (Communications of the ACM (authors' copy))
- Request hedging (gRPC documentation) (The gRPC Authors)
- gRPC Retry Design, proposal A6 (hedging policy and throttling) (The gRPC Authors)
Related tools
Report a problem with this lesson
Kept only in this browser. Your Learn progress