skip to content

Why can't you average the p99 latencies of ten replicas to get the fleet's p99?

level: middleimportance: must knowfreq 68%

answer

  1. A percentile is a position, not a quantity
  2. Averaging ignores how much traffic each served
  3. The pooled value sits between min and max
  4. Add the counts, then estimate once
  5. Same error over time as across replicas

basics

~20 s

A percentile is an order statistic over a set of observations, not a quantity that adds. Averaging per-replica p99s ignores how much traffic each served and describes no real request. Merge the underlying bucket counts first, then estimate the quantile once.

solid answer

~50 s

A p99 is a position in a sorted set of observations, and positions do not add or average across sets. The mean of ten replicas' p99s gives a replica that served 1,143 requests the same weight as one that served 41,208, so the result matches no request anyone made. What *is* true is that the pooled p99 lies between the smallest and the largest of the per-replica p99s — if every replica has 99% of its requests at or below its own p99, so does the union — which makes the **maximum** a legitimate, if pessimistic, upper bound and the average merely a number inside that bracket. The correct aggregation works one level down: counts are additive, so sum the per-boundary bucket counts across the replicas and estimate the quantile once from the merged counts. The same rule applies over time and across shards.

go deeper

for a junior

Recall that a percentile describes a position in one set of measurements, so a percentile computed per replica cannot simply be averaged into a number for all of them.

for a middle

Be ready to explain both defects — unequal traffic weighting and the fact that order statistics do not add — and to give the correct method of merging bucket counts before estimating.

for a senior

Show that the same error appears across time and across services, and distinguish the pooled quantile over a period from the maximum of per-period quantiles, because alerts and objectives want different ones.

for a principal

Own the rule that aggregation happens on counts and that grouping is deferred to query time, and make sure services that must be compared share a bucket grid so the merge is even valid.

## Why the arithmetic has no meaning A percentile is an **order statistic**: sort every observation and read the value at a given position. Sums and means are defined on the values themselves; a percentile is defined on the *ranking* of a particular set. Combine two sets and the ranking changes in a way that no arithmetic on the two old positions can reconstruct, because the numbers that decide the new position — how many observations each set contributed, and where they sat relative to each other — were thrown away when each set was reduced to a single quantile. The weighting problem alone is enough to kill it. Take three replicas of a district-heating billing service over one window: | Replica | Requests | Reported p99 | | --- | --- | --- | | A | 41,208 | 180 ms | | B | 39,884 | 172 ms | | C | 1,143 | 41 ms | The average of the three p99s is 131 ms. But C, which has just come back from a restart and served 1.4% of the traffic, casts a third of the vote. Pool the 82,235 requests and the slowest 822 of them come almost entirely from A and B, so the real pooled p99 sits close to 175 ms. The dashboard is roughly a quarter low, and it is *systematically* low in exactly the situation you care about: during a rollout, when fresh replicas with little traffic and easy latency are averaged in beside loaded ones. ## What is actually true about the bracket One useful fact survives. The pooled q-quantile always lies between the smallest and the largest of the per-group q-quantiles. The argument is one line: if `x` is the largest of the replica p99s, then every replica has at least 99% of its requests at or below `x`, so their union does too, so the pooled p99 cannot exceed `x`. The mirror argument gives the lower bound. So: - **max of the per-replica p99s** — a genuine upper bound on the pooled p99, and often the number you actually want, because it answers "is any replica serving a bad tail?" - **min of the per-replica p99s** — a genuine lower bound, rarely useful. - **mean of the per-replica p99s** — a number inside the bracket with no interpretation, and one whose error grows exactly as traffic becomes uneven. ## The correct aggregation Do the combining one level down, on the counts, because counts over disjoint sets of observations are additive. 1. Take the per-boundary bucket counts from each replica for the same window and the same bucket grid. 2. Add them boundary by boundary to get the pooled distribution. 3. Estimate the quantile once, from the merged counts. This is why an aggregatable histogram is worth its cost: the grouping decision is deferred to query time, so the same stored data answers the p99 per replica, per region and for the whole fleet without anyone precomputing those combinations. It also imposes a requirement people discover the hard way — the replicas must share a bucket grid. Merge counts from two different grids and the addition is meaningless. ## The same mistake in the time dimension Averaging a series of per-minute p99 values into an hourly figure is the identical error with the groups being minutes rather than replicas. A quiet minute at 03:00 weighs as much as the busy minute in which the tail actually blew out, so the smoothing hides the event. The fix is the same: merge the bucket counts across the hour and estimate once. Be careful that these are two different questions, both legitimate: - *The p99 of the hour's requests* — merge the counts over the hour, estimate once. This is what a user-facing objective usually means. - *The worst minute's p99* — compute per minute, then take the maximum. This is what an alert usually means. A dashboard that averages per-minute p99s answers neither, and it always looks better than both. ## And across services The same instinct produces end-to-end numbers by adding one service's p99 to the next service's p99. That is only right if the same request was in the tail at every hop, which is generally false — most slow requests are slow at one hop. It is not a dependable bound in either direction, so do not treat the sum as conservative. If you need an end-to-end tail, measure end to end: instrument a histogram at the entry point over the whole request, and use per-hop histograms to explain where the time went rather than to reconstruct the total. The general rule worth stating out loud in an interview: **aggregate the counts, never the quantiles.** Any time a percentile appears on the left of an arithmetic operator, something has gone wrong.

  • Is the maximum of the per-replica p99s a number worth alerting on?
    Often yes, as long as you say that is what it is. It bounds the pooled p99 from above and it answers a genuinely useful question — whether any single replica is serving a bad tail — which a pooled figure can hide when one replica out of forty is sick. It will read worse than what users experience overall, so it belongs on a per-replica health alert rather than on a user-facing objective.
  • A dashboard builds a daily p99 by averaging hourly p99 points, and it looks far better than what users report. Why?
    Averaging spreads the bad hours across the good ones. A pooled daily p99 is dominated by the busiest hours, which are usually the worst ones, while an average of hourly values weights a quiet night hour equally with a peak hour and pulls the result down. Merge the hourly bucket counts across the day and estimate the quantile once, or, if the question is about the worst period, take the maximum of the hourly values instead.

Averaging ten replicas' p99s is like averaging ten schools' fastest sprint times and calling it the district record: the district's fastest runner is a person who actually ran, not an average of other people's bests.

saying these in an interview costs you the question

  • Averages per-replica percentiles and calls it the fleet figure
  • Adds per-service p99s to get an end-to-end tail
  • Thinks weighting the average by request count fixes it
  • Believes percentiles behave like means under aggregation
  • Merges bucket counts from replicas using different grids