Percentiles, Properly: You Cannot Average a p99

Averaging percentiles across shards or minutes produces a number that means nothing. Here is the arithmetic, plus why p99 gets worse as you add services.


Percentiles: two shards whose p99 values are 100ms and 400ms combine into a p99 that depends on their relative volume, so averaging the two produces a number that matches neither; histogram merging computes the true value

Here’s a number: your p99 latency dashboard is probably wrong. Not slightly off — structurally, arithmetically wrong, in a way that makes it read better than reality.

It’s the most common metrics bug I run into, it survives code review because the code looks reasonable, and it comes from one mistake: treating a percentile like a quantity you can add up.

The arithmetic, with small numbers

A percentile is not a measurement. It’s a position in a sorted list. The p99 is the value at the 99th position out of 100, once you’ve lined everything up. Positions don’t add, and that’s the entire problem.

Take two shards. Ten requests each, latency in milliseconds, already sorted:

Shard A:  10  10  10  10  10  10  10  10  10  100
Shard B:  20  20  20  20  20  20  20  20  20  400

Say p99 of each is its slowest value — 100ms for A, 400ms for B. Your dashboard averages them and reports 250ms.

Now merge the raw data and sort all twenty:

10 10 10 10 10 10 10 10 10 20 20 20 20 20 20 20 20 20 100 400

Twenty values. The 99th percentile sits at the top of that list — 400ms. Your dashboard said 250. It was off by 150ms, and it erred in the comfortable direction.

Now change one thing. Make shard B carry a hundred times the traffic. Shard A’s single 100ms outlier is now a rounding error in a sea of B’s requests, and the true p99 moves close to B’s — around 400ms. The average of the two p99s is still 250ms. Same reported number, completely different reality.

That’s the thing to take away: the error isn’t a fixed bias you could correct for. It depends on the relative volume and shape of the inputs, neither of which the two percentile values carry. There is no fudge factor. The number means nothing.

This applies anywhere you’re combining:

  • p99 across shards or replicas
  • p99 across regions
  • the “average p99 over the last hour” that a time-series query gives you when it rolls up per-minute percentiles
  • a per-endpoint p99 averaged into a service-level p99

That third one catches people constantly. Your metrics backend stored a p99 per minute. You ask for an hour. It averages sixty p99s and hands you something confident-looking. It’s the same bug with a nicer UI.

So how do you actually combine them

Aggregate the distributions, not the summaries. Then compute the percentile once, at the end, from the merged distribution.

In practice that means histograms. Instead of storing “p99 was 230ms,” you store bucket counts: 400 requests in the 0–10ms bucket, 900 in 10–25ms, and so on. Histograms do merge — adding two of them is adding their bucket counts, which is ordinary addition and works correctly across shards, regions and time.

The cost is resolution. Your percentile is only as precise as your bucket boundaries, so with coarse buckets you get “p99 is somewhere between 200 and 500ms.” Which is why the common implementations use exponentially-sized buckets: fine granularity where the values are small, coarse where they’re large, bounded relative error throughout. HdrHistogram does this, Prometheus native histograms do this, and t-digest takes a different route that keeps extra accuracy at the extremes — specifically because the extremes are the part you care about.

If you take one practical step from this post: find out whether your metrics pipeline stores histograms or pre-computed percentiles. If it’s the latter, every aggregation you do downstream is wrong, and no amount of careful dashboarding will fix it.

The measurement bug underneath the measurement bug

Suppose you fix all of the above and your percentiles are now computed correctly. They can still understate the tail badly, for a reason that has nothing to do with arithmetic.

Gil Tene named this one coordinated omission, and the name describes the mechanism precisely. Your load generator is supposed to send a request every 10ms. The system stalls for a full second. During that second the generator is blocked waiting on its outstanding request — so it doesn’t send the hundred requests it was scheduled to send.

When the stall clears, you record one request that took 1000ms. But a hundred requests should have been issued during that window, and they would have taken 1000ms, 990ms, 980ms, and so on down. Your dataset is missing ninety-nine slow measurements, and they’re missing precisely because the system was slow. The load generator and the system under test coordinated their pause. Hence the name.

The effect is large and always in the same direction: the worse the stall, the more measurements it suppresses. This is why a benchmark can report a tidy 5ms p99 for a service that visibly freezes for a second under load.

The fix is to measure each request against when it was supposed to be sent, not when it actually went out. Good load generators do this — wrk2 was written specifically to correct it, and HdrHistogram ships helpers for it. If you’re writing your own harness, this is the thing to get right, and it’s a few lines: record the intended send time up front, and compute latency from that.

The same trap exists in production, incidentally. If a queue backs up and your client-side timer only starts once a connection is acquired from the pool, you’ve omitted the queueing time — which is exactly the part that got bad. Start the clock when the work arrives, not when you get around to it.

Why your p99 is worse than everyone else’s p99

Here’s a calculation that reframes how you think about ownership of latency.

A request fans out to 10 backend services. Each one exceeds its latency threshold 1% of the time, independently. What fraction of your requests touch at least one slow call?

P(no slow call) = 0.99^10 ≈ 0.904
P(at least one) ≈ 9.6%

Every backend has a perfectly respectable 1% tail. Your user-facing request has a 9.6% tail. Nobody did anything wrong, and the end-to-end number is nearly ten times worse than any individual component’s.

Push it further. Twenty services at 1% each gives about 18%. And if the fan-out is sequential rather than parallel, the latencies add, so you’re not just more likely to be slow — you’re slower.

Two consequences I’d actually act on.

First, a per-service latency SLO is necessary and not sufficient. Ten teams can each hit their target while the composed experience misses badly, and the org chart will not surface this. Somebody has to own the end-to-end number.

Second, this is the honest argument for hedged requests — firing a duplicate to a second replica once the first exceeds, say, p95, and taking whichever answers first. You’re spending a few percent extra load to cut the tail, and because the tail compounds multiplicatively across a fan-out, that trade usually pays. It’s also why reducing fan-out width is a latency optimisation in its own right, which isn’t obvious until you’ve done this arithmetic.

What to put on the dashboard

Stop showing the mean. It’s not that averages are evil — it’s that the same mean comes from “everything takes 100ms” and from “most things take 10ms and a few take 4 seconds,” and your users are having completely different days in those two worlds.

Show p50, p95, p99 together. The spread between them is the signal. p50 flat while p99 climbs means a subset of requests found a bad path — a cold cache, an unlucky shard, a GC pause — and that’s a different investigation from all three rising together, which is saturation.

Keep a max, too. A max is a single sample and statistically almost meaningless, and it’s still the first thing I look at, because it tells you what this system is capable of doing to someone.

And write SLOs as percentiles with an explicit window: “99% of requests under 300ms over 30 days.” That composes into an error budget you can actually spend, which an average never will. If you want the full treatment of that, it’s in SLI, SLO, SLA, and error budgets.

One last thing. When you set the threshold, look at where your distribution is multimodal rather than picking a round number. Most real latency distributions have humps — cache hit, cache miss, cache miss plus contention. A threshold that lands in the valley between two humps is stable and meaningful. One that lands on a slope will flap, and you’ll spend a quarter tuning the alert instead of the system.


Related: Queueing Theory for SREs: Little’s Law and the Utilisation Knee · SLI, SLO, SLA, and error budgets · Garbage Collection: The Convenience That Shows Up in Your p99

Frequently asked questions

Why can't you average percentiles?

Because a percentile is a position in a sorted distribution, not a quantity, and positions do not add. The p99 of a combined dataset depends on how the two distributions overlap, which the two p99 values alone do not tell you. Averaging them produces a number that is not the p99 of anything — it can be far too low when one source is small and slow, and it has no error bound you can reason about. The correct approach is to aggregate the underlying distributions, usually as histograms, and compute the percentile once at the end.

What is coordinated omission?

Coordinated omission is a measurement artifact, named by Gil Tene, in which a load generator stops issuing requests while the system under test is stalled, so the stall is recorded as one slow request instead of the many requests that should have been sent during it. Because the measurement tool and the system coordinate their pauses, the resulting latency distribution systematically understates the tail. The fix is to measure latency against the time a request was scheduled to be sent rather than the time it actually went out.

Why does p99 get worse when a request calls more services?

Because the slow paths compound. If a request fans out to ten independent services and each exceeds some threshold one percent of the time, the chance that none of them do is 0.99 raised to the tenth power, about 90.4 percent. So roughly 9.6 percent of requests touch at least one slow call — the overall rate is nearly ten times the per-service rate. This is why tail latency is a system-level property rather than something each service can own in isolation.

Should an SLO use an average or a percentile?

A percentile, essentially always. An average hides the shape of the distribution: the same mean can come from consistently mediocre performance or from mostly-fast responses plus a painful tail, and users experience those very differently. Percentiles also compose into a meaningful error budget, because '99 percent of requests under 300ms' states directly how much slow traffic is acceptable, which an average cannot express.

Comments