The average of three p95s is nobody's p95
In one line
There is more than one way to compute a percentile, and percentiles cannot be combined. The average of per-shard p95s is neither any shard's p95 nor the overall p95.
Why this matters
A table like this often shows up in performance meetings — shard A's p95 is 92 milliseconds, shard B's is 171 milliseconds, and shard C's is 818 milliseconds. Someone averages the three and writes "overall p95 is 360 milliseconds". That number is wrong. If you pool the raw data and measure again, it is 253 milliseconds. Then someone else takes a weighted average by traffic share — 189 milliseconds. That is also wrong, and this time it is wrong in the opposite direction.
There is one more place where things diverge. Ask the same "p95" of the same samples and different tools give different values. A tool that uses nearest-rank returns a value that actually exists in the samples, while a tool that uses linear interpolation makes up a value between two samples. With only 20 samples, the difference spreads as far as 140 milliseconds versus 178 milliseconds. If the tool you used yesterday differs from the one you use today, p95 looks like it got worse even though nothing changed.
These two divergences happen quietly on dashboards. Nobody sees an error, the numbers look plausible, and the SLO verdict is issued on those numbers.
How it works
A percentile is an order statistic. It is the job of pulling out the value at a particular position from sorted samples, so the methods diverge on "how to decide the position". Nearest-rank uses the k = ceil(q x n)th value as it is. Linear interpolation makes a real-valued position h = (n - 1) x q and divides proportionally between the samples before and after it. When there are many samples, the difference between the two disappears, and when samples are few or there is a lone value in the tail, it spreads widely.
The reason they cannot be combined is more fundamental. An average is made of a sum and a count, so adding the sums of the parts gives the whole, but a percentile is a position in sorted order, so adding the positions of the parts does not give the position of the whole. The Prometheus documentation says the same thing — precomputed quantiles cannot be combined with each other, and what can be combined is buckets.
So there are only two right ways. One is to pour all the raw observations into one pot and sort again. It is accurate, but you must carry all the observations around. The other is to add histogram buckets. If for each shard you count "how many at or below 50 milliseconds, how many at or below 100 milliseconds", those counts simply add up. You then estimate the percentile again from the summed table. In exchange, because you answer with buckets rather than values, the boundaries decide the error. If the boundaries are sparse, the estimate gets smeared somewhere inside a bucket, and if they are dense, it gets accurate but the time series multiply. This trade-off is why data structures like HdrHistogram exist.
One more thing: a combined percentile has a fence that always holds. The overall p95 must fall between the minimum and the maximum of the per-shard p95s. That is because if 95% of every shard is at or below some value, then even after mixing, 95% is at or below the maximum. So a report that says "after combining, it came out larger than every shard" should make you suspect the calculation, not the data.
What it looks like in the field
The most common accident is a dashboard that averages per-shard p95s and draws them as one line. A slow shard that has only 10% of traffic is pulling the overall p95 up by a factor of two, and the average line hides that fact. Conversely, if you use a weighted average, the slow shard is erased and the picture is more optimistic than reality. Both betray the intuition that "using an average errs on the safe side".
The same thing happens in load testing. There is a habit of splitting a test that is too heavy to run once into three runs and writing the average of each run's p95 in the report. If the target or the load differs per run, that average measures nothing. The right way is to keep the raw response times of every run and combine and recompute them, which is why you need the habit of not throwing away the tool's raw output (hey -o csv, k6's result file).
The same trap exists on the time axis. It is common for a dashboard to draw a one-hour graph from p95s computed every minute and lay an average line over them, but that line is not the p95 of the hour. The average of 1-minute p95s counts minutes with little traffic and minutes with much traffic with the same weight, so the more the incident happened during busy traffic, the lower than reality it comes out. If you need the p95 of one hour, you must add that hour's buckets and compute again.
There is also something to check when you change tools. Some tools write their calculation method in the documentation and some do not. If you run two tools once on the same data and look first at whether the values differ, you can later tell whether a report that "p95 got worse than last week" is the tool's fault or the service's. With tens of thousands of samples, the difference between the two methods disappears, so the check is meaningful on the side with few samples (short windows, low-traffic periods).
What you will do in the next lab
Using latency samples of three shards made with a fixed seed, you implement two percentile methods by hand and confirm that the values diverge. Then you compare the simple average and the weighted average of the per-shard p95s against the overall p95 to see them be wrong in two directions, and confirm the fence that the overall p95 lies between the minimum and the maximum of the per-shard p95s. Next you implement the method of combining by adding histogram buckets and measure the error by changing the boundaries between sparse and dense ones. Finally you run hey twice, measure the difference between the value from combining the raw data and the value from averaging p95s, and build and submit a tool that combines multiple runs correctly.