skip to content

Why can Elasticsearch's percentiles aggregation report a p99 that differs from the exact value?

level: middleimportance: must knowfreq 55%

answer

  1. percentiles need ordering, sums do not
  2. memory must not grow with document count
  3. the sketch is uneven on purpose
  4. the tails matter more than the middle
  5. one parameter buys resolution with memory

basics

~20 s

The percentiles aggregation summarizes the value distribution in a bounded TDigest sketch rather than sorting every value, so results are approximate. TDigest is most accurate at extreme percentiles and least accurate near the median, and per-shard sketches are merged.

solid answer

~50 s

An exact percentile requires ordering every value, which cannot be done in bounded memory across shards. Elasticsearch's `percentiles` aggregation therefore feeds values into a **TDigest** sketch: a set of centroids that keep fine resolution at the tails and coarse resolution in the middle. That shape is deliberate — p95 and p99 latency numbers are what people act on, so accuracy is highest there and lowest around the median, the opposite of a naive sampling scheme. Accuracy is also proportional to the `compression` parameter, which defaults to 100 and trades memory for precision, and for small numbers of values the sketch is effectively exact. Each shard builds its own sketch and the coordinating node merges them, which adds a little error. `percentile_ranks` is the inverse: give it values and it returns the percentage of documents at or below each.

code

json · 12 lines
json
{
  "size": 0,
  "aggs": {
    "latency": {
      "percentiles": {
        "field": "latency_ms",
        "percents": [50, 95, 99],
        "tdigest": { "compression": 200 }
      }
    }
  }
}

go deeper

for a junior

Know that the percentiles aggregation returns estimates from a sketch rather than exact order statistics, and that you name the percentiles you want with percents.

for a middle

Explain TDigest: centroids that are fine-grained at the tails and coarse in the middle, compression trading memory for accuracy, and per-shard sketches merged by the coordinating node.

for a senior

Catch the misuse in real systems — averaging percentiles across time buckets, nesting a high-compression percentiles aggregation under a wide histogram, or promising an exact tail number to a stakeholder.

for a principal

Own how latency and distribution figures are defined and published across the platform: which percentiles are alerted on, what accuracy is guaranteed, and how sketch merging and shard layout are accounted for in SLA reporting.

## Why percentiles cannot be exact here Sum, count and average are *additive*: each shard produces a partial result and the coordinating node combines them with arithmetic. A percentile is not additive. To know the 99th percentile you must know the ordering of the values, and to compute it exactly across shards you would have to bring every value — or at least a full sorted structure — to one place. Memory would scale with the number of documents, which is precisely what an aggregation framework must avoid. ## What TDigest does The `percentiles` aggregation defaults to a TDigest sketch. TDigest groups values into centroids, each holding a mean and a count. The clever part is that centroid size is not uniform: centroids near the extremes of the distribution are kept small, so the tails are described in fine detail, while centroids in the middle are allowed to absorb many values. Querying a percentile means walking the centroids and interpolating. The consequences follow directly from that shape: - **Extremes are accurate, the middle is not.** p1, p99 and p99.9 are estimated far more precisely than p50. This is the opposite of most people's intuition and is the single best thing to say in an interview, because it also explains why the sketch was designed this way: tail latency is what teams alert on. - **Small data sets are effectively exact.** If the number of values is small enough that each lands in its own centroid, no information is lost and the reported percentile is the real one. Developers testing on a few hundred documents therefore see exact numbers and wrongly conclude the aggregation is exact. - **Accuracy scales with compression.** The `compression` parameter, default 100, controls how many centroids the sketch may keep. Higher compression means more centroids, more memory and CPU, and a tighter estimate. It is the percentile analogue of `precision_threshold` on `cardinality`. ```json { "aggs": { "latency_pct": { "percentiles": { "field": "latency_ms", "percents": [50, 95, 99, 99.9], "tdigest": { "compression": 200 } } } } } ``` With no `percents` specified the aggregation returns 1, 5, 25, 50, 75, 95 and 99. ## The HDR alternative Elasticsearch also offers an HDR Histogram implementation via the `hdr` option with `number_of_significant_value_digits`. HDR fixes memory up front for a bounded, non-negative value range and gives a guaranteed relative precision, which suits latency measurements with a known maximum. It costs more memory than TDigest for wide ranges and cannot represent values outside the configured range, so TDigest remains the sensible default unless you have measured a reason to switch. ## Merging across shards Each shard builds a sketch over its own matching documents; the coordinating node merges them into one and reads the percentiles from the merged sketch. Merging TDigests is well defined but not lossless — merged centroids blur slightly — so a heavily sharded index carries a little more error than the same data in one shard. This is also why the result depends on the shard layout: reindexing the same data into a different shard count can shift a reported p99 marginally, which looks like a bug and is not one. ## The percentile arithmetic trap Because percentiles are not additive, you cannot average them. Given hourly p99 latencies, the daily p99 is *not* the mean of the twenty-four hourly figures, and it is not their maximum either. The only correct way to get a daily p99 is to compute it over the day's documents in a single aggregation, or to merge the underlying sketches. Dashboards that average percentiles across buckets are producing a number with no statistical meaning, and spotting that in a review is a strong senior signal. ## percentile_ranks: the inverse question `percentile_ranks` takes the values you care about and reports what share of documents sit at or below each: ```json { "aggs": { "within_sla": { "percentile_ranks": { "field": "latency_ms", "values": [200, 500] } } } } ``` A result of 96.4 for 500 means roughly 96.4% of requests completed within 500 ms — the natural shape for an SLA statement, where the threshold is fixed and the compliance figure is what you need. It uses the same sketch and carries the same approximation. ## Practical guidance Use percentiles rather than mean and standard deviation for latency, price and any skewed distribution. Report p50 alongside the tail so readers can see the skew. Raise `compression` only if you have shown the default is insufficient for the percentiles you actually publish, and remember that memory is per bucket per shard, so a percentiles aggregation nested under a wide `date_histogram` multiplies. Finally, if a stakeholder needs a number they can defend to the decimal, say plainly that this aggregation does not produce one.

  • Given hourly p99 latencies, how do you get the p99 for the whole day?
    Not by averaging or maxing the hourly figures — percentiles are not additive and neither operation has a statistical meaning. Compute the percentile over the full day's documents in one aggregation, or merge the underlying sketches. A dashboard that averages percentiles across buckets is publishing a number that corresponds to nothing.
  • When would you switch the percentiles aggregation from TDigest to the hdr implementation?
    When values sit in a known, bounded, non-negative range — typical of latency in milliseconds — and you want a guaranteed relative precision rather than TDigest's tail-weighted accuracy. HDR sizes its memory up front from number_of_significant_value_digits and cannot record values outside the configured range, so it is a deliberate trade, not a default.
  • Why do developers often believe the percentiles aggregation is exact?
    On small data sets it effectively is. When the number of values is small enough that each occupies its own centroid, nothing is approximated and the reported percentile matches the real one. Test fixtures with a few hundred documents therefore agree perfectly with a hand calculation, and the divergence only appears at production volume.

TDigest is like a ruler with millimetre markings at both ends and centimetre markings in the middle: it measures the extremes precisely because that is where the interesting measurements happen, and gives up resolution where nobody looks closely.

saying these in an interview costs you the question

  • Says percentiles are exact because a small test matched
  • Believes accuracy is best around the median
  • Averages per-bucket p99 values to get an overall p99
  • Confuses percentile_ranks with asking for arbitrary percentiles
  • Thinks compression sets how many percentiles are returned

context