Percentiles Don't Add Up

Jun. 21, 2026 · 4 min read

Your p99 latency might not be a percentile. Averaging percentiles across replicas produces a meaningless number, and t-digest is the sketch algorithm that fixes it at scale. Image: Davidson-Pilon, Cameron (2015). “Percentile and Quantile Estimation of Big Data: The t-Digest.” DataOrigami.

Percentiles don’t add up

Your on-call engineer was paged on Saturday. Users reported 3 second load times, but when the engineer checked the dashboard, the p99 latency was 180ms.

But 180ms wasn’t a meaningful value. Six replicas were each computing a local p99 and shipping it to the monitoring system, which was just averaging them. Averaging the percentiles produces an average, not a percentile. One replica was saturated while the others were fine, and it disappeared into the average. The fleet’s actual p99 was 2.8 seconds.

Percentiles fail at scale

A percentile is about rank. The p99 (99th percentile) of a set of latencies is the value below which 99% of requests fall. Computing it requires sorting the full set, every single data point.

Keeping every single data point is impractical at scale. A service handling 10,000 rps accumulates 864 million latency values per day. Assuming 32-bit integers plus an additional 64-bit timestamp per observation, that’s ~3.78 TB/year for just a single latency metric (and that ignores indexes and other storage overhead).

Monitoring systems solve this with sketches. A sketch is a compact summary that estimates statistics with a configurable error bound. Instead of storing all n observations, a sketch will use sublinear or parameter-bounded memory. Decreasing/increasing the memory budget generally degrades/improves accuracy.

What t-digest does

Ted Dunning introduced t-digest in 2013 as a sketch optimized for quantile estimation, later formalized with Otmar Ertl in 2019 treatment[1]. t-digest’s key property is that its accuracy is highest at the tails and lowest near the median.

t-digest processes observations one at a time without holding the raw data. Each observation is absorbed into the nearest eligible centroid, then it’s discarded. The only structure that persists is the list of centroids. Every centroid stores only the mean and count (weight) representing the number of observations it covers.

t-digest summarizes thousands of latency measurements into a much smaller number of centroids. Cameron Davidson-Pilon fed 8 MB of Pareto-distributed data into a t-digest and the output was 5 kB, with an error rate below 0.002%[2].

Centroids near the tails are numerous, but represent relatively few observations, which preserves the resolution where extreme-percentile accuracy matters. Centroids near the median are fewer and typically summarize many more observations, allowing this dense region to be compressed with little loss in resolution[1].

t-digest’s asymmetry matches how latency metrics are used in practice. In large distributed systems, the slowest requests tend to shape users’ impressions. Jeff Dean and Luiz André Barroso made the case for the “tail at scale” in 2013, and the industry shifted toward percentile-based service-level objectives (SLOs)[3].

A compression parameter δ is configured when creating a t-digest instance. It sets a limit on how many centroids the sketch is allowed to maintain[1]. Raising δ improves accuracy at the cost of size.

Separate t-digests merge without any loss of accuracy[1]. Combining and compressing the centroid lists is equivalent to processing every observation in a single pass. Each replica computes and ships a t-digest instead of a computed p99, and the monitoring system combines the t-digests at query time. The 2.8-second tail stays visible instead of disappearing into an average.

Netflix hit this problem, and found that storing raw latency observations at their scale was too expensive. t-digest let them summarize per instance and merge at query time. As traffic grows, memory and compute stays flat, since they depend on the compression parameter, not the number of n observations[4].

Histograms vs. Summaries in Prometheus

Prometheus faces the same aggregation problem[5], solved by implementing histograms that aggregate bucket counts across replicas, the way t-digest merges centroids.

Summaries compute quantiles inside each application process, exporting only the finished p99.

Histograms record observations in predefined buckets and ship the raw counts. The monitoring system computes quantiles from that bucket data at query time, accurate to within the width of the bucket[5].

Native histograms is a newer Prometheus feature that use adaptive bucket boundaries instead of fixed ones. This allows you to get accurate quantile estimates without choosing boundaries upfront.

How teams get this wrong

Application performance monitoring (APM) tools often receive metrics as precomputed percentile values, and there’s no way to verify how they were computed after the fact.

Prometheus summaries expose a quantile label that looks like histogram data, but they’re also precomputed percentile values.

Classic histograms require upfront bucket definitions, guessing the shape of the distribution before any requests arrive. Native histograms sidestep this by choosing boundaries adaptively.

Put it into practice

Check whether a percentile metric was computed at query time or exported as a finished number, since the two carry different guarantees.

When you see an aggregate percentile across replicas, verify the aggregation was done correctly. histogram_quantile(0.99, sum(rate(...))) is correct. avg(metric{quantile="0.99"}) is not.

t-digest is a common choice for query-time quantile computation, with native implementations in ClickHouse, Redis, and Elasticsearch.


References

  1. Dunning, Ted and Otmar Ertl (2019). "Computing Extremely Accurate Quantiles Using t-Digests." arXiv:1902.04023. https://arxiv.org/abs/1902.04023
  2. Davidson-Pilon, Cameron (2015). "Percentile and Quantile Estimation of Big Data: The t-Digest." DataOrigami. https://web.archive.org/web/20230319061101/https://dataorigami.net/2015/03/19/Percentile-and-Quantile-Estimation-of-Big-Data-The-t-Digest.html
  3. Dean, Jeff and Luiz André Barroso (2013). "The Tail at Scale." Communications of the ACM, 56(2): 74-80. https://research.google/pubs/pub40801/
  4. Ortiz, Thiara (2023). "Measuring Real-Life Latency of the Internet: A Netflix Story." SREcon23 Americas. USENIX. https://www.youtube.com/watch?v=wXWFrKiJXHE
  5. Prometheus (2024). "Histograms and Summaries." Prometheus Documentation. https://prometheus.io/docs/practices/histograms/
  6. Masson, Charles, Jee E. Rim, and Homin K. Lee (2019). "DDSketch: A Fast and Fully-Mergeable Quantile Sketch with Relative-Error Guarantees." VLDB 2019. https://arxiv.org/abs/1908.10693

Outtakes

t-digest belongs to a family of probabilistic data structures (sketches) that trade precision for a fixed memory budget. HyperLogLog estimates cardinality, the count of distinct items that are passed. Count-min sketch estimates the frequency of how often each item appeared. Bloom filters is for set membership (Cormode and Muthukrishnan, 2005).

DDSketch is Datadog’s quantile sketch where error is proportional to the quantile value itself[6].

Averaging averages. If Group A (10 people, $50K average) and Group B (100 people, $100K average) merge, the combined average is $95.4K. A naive average of the two group averages incorrectly gives $75K. Averaging percentiles fail for the same reason.

OpenTelemetry’s exponential histograms use bucket boundaries that scale geometrically, concentrating precision at the tails without configuring the boundaries upfront (OpenTelemetry, 2024).


Changelog

2026-08-08 Corrected t-digest’s attribution: Ted Dunning introduced it alone in 2013; Otmar Ertl joined for the more rigorous 2019 formalization.
2026-06-21 Initial draft.