System Design (High-Level Design) Module 3 – Performance and reliability fundamentals
Latency, throughput and percentiles
Tell latency from throughput and bandwidth, read p50 to p99.9 instead of averages, and merge latency histograms across servers without wrong maths.
What you will learn
- Distinguish latency, throughput and bandwidth, and say which one a requirement constrains
- Read p50, p95, p99 and p99.9 and explain why averages mislead
- Aggregate latency correctly across servers using histograms
Before you start
On this page
“Is it fast?” has two different answers. Latency is how long one request takes, from the moment it is sent to the moment its answer arrives. Throughput is how much work the system finishes per second. A design can be excellent at one and poor at the other, so requirements name both, and they name latency as a percentile: “p99 under 300 ms” means that 99 requests in 100 finish within 300 ms. This lesson shows why the average is the wrong number to watch, how to read p50 to p99.9, and how to combine percentiles from many servers without producing a number that describes nobody.
Latency, throughput and bandwidth
The three words are often swapped, and each answers a different question:
- Latency is time per request, in milliseconds. It includes every wait along the way: queues, network hops, disk reads and the work itself. Always say where it was measured, because the client sees the network and the load balancer that the server’s own timer never sees.
- Throughput is completed work per unit of time: requests, messages or rows per second. Google’s SRE book lists request latency, error rate and system throughput, in requests per second, among the indicators most services track.
- Bandwidth is the capacity of a link: the most bits per second it can carry. Throughput is what you actually achieve, and it is never more than the bandwidth of the narrowest link on the path.
The estimation module tied two of them together: requests in progress equal throughput times latency. That law also explains why they can move in opposite directions. Batching writes, sending fewer and larger messages, or keeping a deeper queue of work raises throughput, because each trip does more, and raises latency, because each request now waits for its batch or its place in line. A system can finish a million requests a second while some of them take several seconds. That is why a requirement such as “10,000 requests a second” says nothing about whether users wait, and “p99 under 200 ms” says nothing about how many users can be served.
Percentiles say what users actually wait
Sort the latencies of a period from fastest to slowest. The p50, or median, is the value in the middle: half of the requests were faster. The p95 is the value that 95% of requests were at or below, the p99 the value for 99%, and the p99.9 the value for 999 requests in 1,000. The lesson’s programs use the simplest definition, the nearest rank: the smallest value with at least p% of the sample at or below it. Many monitoring tools interpolate between neighbouring values instead, which changes the result slightly on small samples and hardly at all on large ones.
Here are 10,000 latencies from a model of a service whose requests usually take about 40 ms, while roughly 2% miss a cache and wait 200 to 900 ms longer:
"""Percentiles of a latency sample, and why the mean hides the slow requests.
10,000 request latencies from a seeded generator: most requests take about 40 ms, and about 2% miss a cache
and spend 200 to 900 ms more. The shape is an assumption for the sketch; change it and run the program again.
"""
import math
import random
def percentile(sorted_ms, p):
"""Nearest-rank percentile: the smallest value with at least p% of the sample at or below it."""
# Round before the ceiling: 99.9 has no exact binary form, and 99.9 / 100 * 10,000 is 9990.000000000002.
rank = math.ceil(round(p * len(sorted_ms) / 100, 6))
return sorted_ms[max(rank, 1) - 1]
rng = random.Random(7)
sample = []
for _ in range(10_000):
ms = rng.lognormvariate(math.log(40), 0.3) # the usual path: a median of about 40 ms
if rng.random() < 0.02: # about 2% miss the cache
ms += rng.uniform(200, 900)
sample.append(ms)
sample.sort()
mean = sum(sample) / len(sample)
print(f"{len(sample):,} requests")
print(f"{'mean':8}{mean:6.1f} ms")
for p in (50, 90, 95, 99, 99.9):
print(f"{'p' + str(p):8}{percentile(sample, p):6.1f} ms")
print(f"{'max':8}{sample[-1]:6.1f} ms")
# Two services with the same mean of 100 ms: one is always 100 ms, the other is 50 ms for 98% of requests
# and 2,550 ms for the other 2% (980 x 50 + 20 x 2,550 = 100,000 ms over 1,000 requests).
steady = [100.0] * 1_000
spiky = sorted([50.0] * 980 + [2_550.0] * 20)
print("\nSame mean, different users' experience")
print(f"{'':8}{'mean':>8}{'p50':>8}{'p95':>8}{'p99':>8}")
for name, values in (("steady", steady), ("spiky", spiky)):
stats = [sum(values) / len(values)] + [percentile(values, p) for p in (50, 95, 99)]
print(f"{name:8}" + "".join(f"{s:8.0f}" for s in stats)) Output
10,000 requests
mean 52.6 ms
p50 40.3 ms
p90 60.8 ms
p95 69.8 ms
p99 555.3 ms
p99.9 898.1 ms
max 924.8 ms
Same mean, different users' experience
mean p50 p95 p99
steady 100 100 100 100
spiky 100 50 50 2550
Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 percentiles.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
Read the first block from the top. The median is 40.3 ms and the p95 is 69.8 ms: most users are fine. The mean is 52.6 ms, which looks fine too. But the p99 is 555.3 ms, more than ten times the mean, and the p99.9 is 898.1 ms. The mean blends the 2% of slow requests into the 98% of fast ones and reports a number that hardly any request actually took.
The second block makes the same point without any randomness. Both services average exactly 100 ms. The steady one takes 100 ms for everybody; the spiky one serves 98% of requests in 50 ms and makes the other 2% wait 2.55 seconds. The averages are identical and the experiences are not: on the spiky service, one user in fifty waits long enough to give up. Google’s SRE book recommends the same habit: treat latency as a distribution, and watch the median for the typical request and a high percentile for the bad ones.
Why the tail matters more than its share suggests
One slow request in a hundred sounds rare, but users rarely make one request. A page that needs 50 backend calls, each slow 1% of the time on its own, has at least one slow call on of page loads. The same arithmetic applies to a user’s session: someone who clicks 100 times will meet the p99 about once. The next lesson follows this effect through systems that fan a request out to hundreds of servers.
Say what a percentile is of
A percentile without its window and its population cannot be checked. “p99 over five minutes, for checkout requests only, measured at the load balancer” can be; “p99 is 180 ms” cannot. The window matters because a short incident vanishes in a daily figure. The population matters because cheap requests, such as health checks or cached reads, can make up most of the traffic and pull every percentile down. A useful habit is to report the median, the p99 and the count of requests together, so a reader can tell a calm hour from a quiet one.
Averaging percentiles across servers is wrong
A service runs on many hosts, and each host can compute the p99 of the requests it served. The tempting shortcut is a dashboard that shows the mean of the hosts’ p99s as the fleet’s p99. It is wrong, because a percentile of the whole depends on where the slow requests are, and the per-host numbers have thrown that away. Ten hosts serve equal shares of traffic, and one of them has long garbage-collection pauses:
"""Why the average of per-host p99s is not the fleet's p99.
Ten hosts behind a round-robin load balancer serve 2,000 requests each. Host h07 has long garbage-collection
pauses: a fifth of its requests wait 300 to 700 ms more. The numbers are assumptions for the sketch.
"""
import math
import random
def percentile(sorted_ms, p):
"""Nearest-rank percentile: the smallest value with at least p% of the sample at or below it."""
# Round before the ceiling: 99.9 has no exact binary form, and 99.9 / 100 * 10,000 is 9990.000000000002.
rank = math.ceil(round(p * len(sorted_ms) / 100, 6))
return sorted_ms[max(rank, 1) - 1]
rng = random.Random(11)
hosts = {}
for n in range(1, 11):
name = f"h{n:02}"
latencies = []
for _ in range(2_000):
ms = rng.lognormvariate(math.log(40), 0.3)
if name == "h07" and rng.random() < 0.2: # the host with long pauses
ms += rng.uniform(300, 700)
latencies.append(ms)
hosts[name] = sorted(latencies)
print(f"{'host':6}{'requests':>9}{'p50 ms':>9}{'p99 ms':>9}")
for name, latencies in hosts.items():
note = " <- long pauses" if name == "h07" else ""
print(f"{name:6}{len(latencies):>9,}{percentile(latencies, 50):9.1f}{percentile(latencies, 99):9.1f}{note}")
per_host_p99 = [percentile(latencies, 99) for latencies in hosts.values()]
everything = sorted(ms for latencies in hosts.values() for ms in latencies)
print()
print(f"{'average of the 10 per-host p99s':34}{sum(per_host_p99) / len(per_host_p99):6.1f} ms")
print(f"{'largest per-host p99':34}{max(per_host_p99):6.1f} ms")
print(f"{f'p99 of all {len(everything):,} requests':34}{percentile(everything, 99):6.1f} ms") Output
host requests p50 ms p99 ms h01 2,000 39.8 81.0 h02 2,000 39.6 79.8 h03 2,000 39.2 77.5 h04 2,000 40.3 80.0 h05 2,000 39.9 83.7 h06 2,000 39.8 81.6 h07 2,000 43.8 725.1 <- long pauses h08 2,000 39.7 78.5 h09 2,000 39.9 79.7 h10 2,000 40.0 81.5 average of the 10 per-host p99s 144.8 ms largest per-host p99 725.1 ms p99 of all 20,000 requests 532.0 ms
Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 per_host.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 true p99 of all 20,000 requests is 532.0 ms: about 400 of them were slow, far more than the 200 that make up the slowest 1%. The mean of the per-host p99s says 144.8 ms, which would tell an engineer that the fleet is healthy. The largest per-host p99, 725.1 ms, errs the other way. Neither is a percentile of any set of requests. Prometheus’s documentation states the rule bluntly: quantiles computed by each instance cannot be aggregated, and averaging them yields statistically nonsensical values.
Two ways to get a fleet-wide p99, and only one is right
Text description of the diagram
The diagram starts at the top with three hosts, A, B and C, that measure the latency of every request they serve. Two paths lead down from them.
- The wrong path: each host reports its own p99, and a dashboard takes the mean of the three. The result is not a percentile of any set of requests, because it gives a host that served a few requests the same weight as a busy one and throws away where the slow requests were.
- The right path: each host reports how many requests fell into each latency bucket. The counts are added bucket by bucket, and the p99 is read from the merged counts, which describe every request of the fleet.
The merged estimate can only be as precise as the buckets: it is right to within the width of the bucket the p99 falls in.
The fix is to keep the information a percentile needs. Each host counts its requests in a fixed set of latency buckets, such as 0 to 10 ms, 10 to 25 ms and so on. Counts add up across hosts, so the fleet’s histogram is the sum of the hosts’ histograms, bucket by bucket, and any percentile can then be read from the sum:
// Merging latency histograms from three hosts, and estimating the fleet's p99 from the merged counts.
// Each host counts its requests per latency bucket, as a metrics library does; the counts are added bucket by
// bucket, and the p99 is read from the sum with linear interpolation inside its bucket. Hosts A and B are
// healthy; host C pauses, so a fifth of its requests take 300 to 700 ms longer. The workload is an assumption.
// A small seeded random generator (mulberry32), so every run prints the same.
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(2024);
function latency(slowShare) {
let ms = 20 + 15 * (random() + random() + random()); // the usual path: 20 to 65 ms
if (random() < slowShare) ms += 300 + 400 * random(); // a pause
return ms;
}
const hosts = [
{ name: "host A", requests: 4000, slowShare: 0 },
{ name: "host B", requests: 4000, slowShare: 0 },
{ name: "host C", requests: 2000, slowShare: 0.2 },
];
for (const host of hosts) host.latencies = Array.from({ length: host.requests }, () => latency(host.slowShare));
// What each host exports: a count per bucket; `bounds` are the buckets' upper bounds in ms.
const countsFor = (latencies, bounds) => {
const counts = bounds.map(() => 0);
for (const ms of latencies) counts[bounds.findIndex((bound) => ms <= bound)] += 1;
return counts;
};
const addCounts = (all) => all[0].map((_, i) => all.reduce((sum, counts) => sum + counts[i], 0));
// The q-quantile of bucket counts: find the bucket that holds the rank, then interpolate linearly inside it.
function quantile(q, counts, bounds) {
const rank = q * counts.reduce((a, b) => a + b, 0);
let below = 0;
for (let i = 0; i < counts.length; i++) {
if (below + counts[i] >= rank) {
const lower = i === 0 ? 0 : bounds[i - 1];
if (bounds[i] === Infinity) return lower; // no upper bound to interpolate towards
return lower + ((bounds[i] - lower) * (rank - below)) / counts[i];
}
below += counts[i];
}
return NaN;
}
const BOUNDS = [10, 25, 50, 100, 250, 500, 1000, Infinity];
const perHostCounts = hosts.map((h) => countsFor(h.latencies, BOUNDS));
const merged = addCounts(perHostCounts);
const label = (i) => `${i === 0 ? 0 : BOUNDS[i - 1]}-${BOUNDS[i] === Infinity ? "more" : BOUNDS[i]}`;
console.log(`${"bucket (ms)".padEnd(12)}${hosts.map((h) => h.name.padStart(8)).join("")}${"merged".padStart(8)}`);
BOUNDS.forEach((_, i) => {
console.log(`${label(i).padEnd(12)}${perHostCounts.map((c) => String(c[i]).padStart(8)).join("")}${String(merged[i]).padStart(8)}`);
});
const every = hosts.flatMap((h) => h.latencies).sort((a, b) => a - b);
const exact = every[Math.ceil(0.99 * every.length) - 1];
const perHost = perHostCounts.map((counts) => quantile(0.99, counts, BOUNDS));
const ms = (x) => `${x.toFixed(0)} ms`;
const show = (what, value) => console.log(`${what.padEnd(32)}${value}`);
console.log("");
show("each host's own p99:", perHost.map(ms).join(", "));
show("average of those three:", ms(perHost.reduce((a, b) => a + b, 0) / perHost.length));
show("p99 of the merged counts:", ms(quantile(0.99, merged, BOUNDS)));
// Finer buckets where the tail lives shrink the interpolation error.
const FINE = [10, 25, 50, 100, 250, 400, 500, 600, 700, 800, 900, 1000, Infinity];
const fine = addCounts(hosts.map((h) => countsFor(h.latencies, FINE)));
show("... with buckets every 100 ms:", ms(quantile(0.99, fine, FINE)));
show(`p99 of all ${every.length} raw latencies:`, ms(exact)); Output
bucket (ms) host A host B host C merged 0-10 0 0 0 0 10-25 33 30 1 64 25-50 3243 3278 1354 7875 50-100 724 692 259 1675 100-250 0 0 0 0 250-500 0 0 154 154 500-1000 0 0 232 232 1000-more 0 0 0 0 each host's own p99: 97 ms, 97 ms, 957 ms average of those three: 384 ms p99 of the merged counts: 784 ms ... with buckets every 100 ms: 634 ms p99 of all 10000 raw latencies: 630 ms
Recorded with Node.js 24.21.0 on macOS 26 arm64. To run it yourself: mise exec node@24.21.0 -- node histogram_merge.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
The table is what three hosts would export, and the “merged” column is their sum. Reading the p99 from the merged
counts gives 784 ms against a true value of 630 ms. That error comes from the buckets, not from the method: the p99
falls in the 500 to 1,000 ms bucket, and the estimate assumes the requests in a bucket are spread evenly across it,
which is what Prometheus’s histogram_quantile() does too. With a boundary every 100 ms in that range, the estimate
becomes 634 ms. The average of the per-host p99s, 384 ms, is wrong in a way no choice of buckets can repair.
So the precision you get is the precision of your bucket boundaries, and the boundaries should be dense where your
objectives are. Prometheus’s guide puts the error of such an estimate at no more than the width of the bucket the
percentile falls in, and shows that a boundary exactly at an objective’s threshold, such as 300 ms, lets you count
the requests within it exactly. A query such as
histogram_quantile(0.99, sum by (le) (rate(http_request_duration_seconds_bucket[5m]))) sums the buckets of every
instance before taking the percentile, which is the merge above. HdrHistogram takes another route: it keeps the
rounding error under a fixed share of each value, set as a number of significant digits (three digits means at most
0.1%), with a memory footprint that depends only on the range and precision chosen.
A load test can hide the tail it should measure
A load tester that sends one request, waits for the reply and only then sends the next records a 10-second stall as a single slow sample, while real users, who arrive on their own schedule, would have sent hundreds of requests that all waited. HdrHistogram’s documentation calls this coordinated omission: the samples that go missing are exactly the bad ones, so the report leans towards good results. Its fix is to record each value with the expected interval between requests, so that a long stall also fills in the samples that were skipped.
Practice
Exercise · Easy · Python
Percentiles from a sample and from merged bucket counts
Write three helpers in latency_stats.py.
nearest_rank(latencies_ms, p) returns the p-th percentile of a list of latencies by the nearest-rank method: sort the values, and return the smallest one with at least p% of the values at or below it. p is greater than 0 and at most 100, and the list is not sorted when it arrives. Do not change the caller's list.
merge(histograms) adds the bucket counts of several hosts. Each histogram is a dictionary from a bucket's upper bound in milliseconds to the number of requests in that bucket, and a host may leave out buckets it has no requests in. Return one dictionary with every bucket that any host reported.
bucket_of_percentile(counts, p) returns the upper bound of the bucket that holds the p-th percentile of merged counts: walk the buckets from the fastest to the slowest and stop at the first one where the running total reaches p% of all requests.
For example, nearest_rank([30, 10, 20, 50, 40], 50) is 30, merge([{25: 3, 50: 5}, {25: 1, 100: 2}]) is {25: 4, 50: 5, 100: 2}, and bucket_of_percentile({50: 90, 100: 8, 500: 2}, 99) is 500.
Starter code · latency_stats.py
import math
def nearest_rank(latencies_ms, p):
"""The p-th percentile of a list of latencies, by the nearest-rank method."""
# Replace this line with your code.
return 0
def merge(histograms):
"""Add per-host bucket counts: each histogram maps a bucket's upper bound in ms to its count."""
# Replace this line with your code.
return {}
def bucket_of_percentile(counts, p):
"""The upper bound of the bucket that holds the p-th percentile of merged counts."""
# Replace this line with your code.
return 0 The sample tests · test_latency_stats.py
from latency_stats import bucket_of_percentile, merge, nearest_rank
def test_nearest_rank():
"""sorts the values and returns the value at the nearest rank"""
assert nearest_rank([30, 10, 20, 50, 40], 50) == 30
assert nearest_rank([30, 10, 20, 50, 40], 100) == 50
assert nearest_rank([30, 10, 20, 50, 40], 20) == 10
assert nearest_rank(list(range(1, 1001)), 99) == 990
assert nearest_rank(list(range(1, 1001)), 99.9) == 999
def test_nearest_rank_keeps_input():
"""leaves the caller's list as it was"""
values = [5, 3, 9]
nearest_rank(values, 50)
assert values == [5, 3, 9]
def test_merge():
"""adds counts bucket by bucket, keeping buckets only one host reported"""
assert merge([{25: 3, 50: 5}, {25: 1, 100: 2}]) == {25: 4, 50: 5, 100: 2}
assert merge([{10: 7}]) == {10: 7}
assert merge([]) == {}
def test_bucket_of_percentile():
"""walks the buckets from the fastest and stops where the running total reaches p%"""
counts = {100: 1675, 25: 64, 1000: 232, 50: 7875, 500: 154}
assert bucket_of_percentile(counts, 99) == 1000
assert bucket_of_percentile(counts, 95) == 100
assert bucket_of_percentile(counts, 50) == 50
assert bucket_of_percentile({50: 90, 100: 8, 500: 2}, 99) == 500
assert bucket_of_percentile({50: 90, 100: 8, 500: 2}, 98) == 100 A hint
For the nearest rank, sort a copy (sorted(latencies_ms)); the rank is p% of the number of values, rounded up, and lists count from 0, so the value is at index rank − 1. Round before you round up: 99.9 has no exact binary form, so 99.9 / 100 * 1000 is 999.0000000000001, and math.ceil of that is 1000. math.ceil(round(p * n / 100, 6)) avoids it. For the bucket, walk sorted(counts) from the fastest bound to the slowest and compare the running total with p / 100 * total.
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 (fresher to mid level): what does “p99 latency of 300 ms” mean, and why is it reported instead of the average? It means that 99% of the requests in the stated window finished within 300 ms and 1% took longer. Teams report it because the average hides the slow requests: a few very slow ones and many fast ones give an average that almost no request took, and the users stuck in the slow 1% are often the ones who leave or retry. A percentile also compares cleanly with a target, because “99% under 300 ms” is a promise about users, while an average can stay flat while the tail doubles. Always ask for the window and the population behind the number.
Key takeaways
- Latency is time per request, throughput is work per second, bandwidth is a link’s capacity; batching and queues can raise throughput and latency together.
- Report latency as percentiles, with the window and the population: the mean hides the slow requests that users remember.
- A 1% tail per call becomes common per page: with 50 calls, 39.5% of page loads meet it.
- Never average per-host percentiles. Add per-host bucket counts and read the percentile from the sum; its precision is the width of the bucket it falls in.
- A load tester that waits for each reply under-reports stalls (coordinated omission).
References
- Site Reliability Engineering, chapter 4: Service Level Objectives (Google (O'Reilly Media))
- Histograms and summaries (Prometheus documentation) (The Prometheus Authors)
- Query functions: histogram_quantile() (Prometheus documentation) (The Prometheus Authors)
- HdrHistogram, a High Dynamic Range Histogram (HdrHistogram project)
- HdrHistogram Java API: package org.HdrHistogram (coordinated omission) (HdrHistogram project)
- The Tail at Scale (Jeffrey Dean and Luiz André Barroso) (Communications of the ACM, via Google Research)
Related tools
Report a problem with this lesson
Kept only in this browser. Your Learn progress