Why can an approximate percentile over unchanged data return a slightly different p99 each run?
answer
- exact percentiles need order, sketches do not keep it
- the summary discards points as it fills
- which points survive depends on arrival order
- parallel plans do not fix the merge order
- the bound is on rank, not on the value
basics
~20 sApproximate percentile functions build a compact quantile summary whose retained points depend on the order values arrive and how partial summaries merge. Parallel execution varies that order between runs, so the compacted summary and its estimate shift slightly.
solid answer
~50 sAn exact percentile needs the values ordered, which means a full sort or a full copy of the column. Approximate percentile functions instead maintain a bounded-size **quantile summary** — commonly a t-digest, which keeps weighted centroids packed more finely toward the tails, or a KLL-style sketch, which keeps a sampled hierarchy with a uniform rank-error bound. Both summarise by discarding and compacting points, and which points survive depends on the order values were absorbed and on how partial summaries from different workers were merged. A distributed engine does not guarantee that order: split boundaries, file ordering, worker count and concurrency all vary, so the compaction lands differently and the estimate moves. Note the contrast with distinct-count sketches, whose merge is a maximum and therefore order-independent and reproducible. Fixes are to raise the accuracy parameter, compute exactly for small groups, or stop treating a percentile as reproducible to the last digit.
go deeper
Know that percentile functions prefixed with approx are estimates built from a compact summary, and that small run-to-run movement is expected rather than a bug.
Explain the mechanism: a bounded-size summary compacts points as it fills, and the compaction depends on arrival and merge order, which a parallel plan does not fix.
Show you can act on it — raise the accuracy parameter, compute exactly for small groups, keep alert thresholds outside the error band, and merge stored summaries rather than averaging finished percentiles.
Own the reporting contract: which percentiles are declared approximate, what tolerance is published to consumers, and how SLA or contractual numbers get an exact path that is auditable.
## What an approximate percentile actually computes A percentile is defined by rank: the p99 of a column is the value with 99% of the data below it. Computing it exactly requires putting the values in order, which means either a full sort — expensive and memory-hungry at scale — or retaining the whole column. Approximate percentile functions replace the ordered data with a **quantile summary**: a bounded-size structure that can answer "what value sits at rank q?" within a stated error. Two families dominate. **t-digest** groups values into weighted centroids, each holding a count and a mean. The permitted centroid weight is a function of the centroid's position in the distribution: tiny near rank 0 and rank 1, larger in the middle. That deliberate asymmetry gives fine resolution at the tails — exactly where p99 and p999 live — and coarser resolution around the median, at a compact fixed size controlled by a compression parameter. **KLL** and related sampling-based summaries keep a hierarchy of buffers with items promoted and halved as levels fill. Their guarantee is a uniform bound on **rank error**: the returned value is guaranteed to lie between the (q − ε) and (q + ε) quantiles of the data, for a stated ε. ## Why the answer moves between runs Both families are **lossy by compaction**. When a buffer or centroid list is full, the structure merges neighbouring points and discards detail. Which points end up merged depends on the sequence in which values arrived and, in a distributed plan, on how the workers' partial summaries were merged together. A warehouse gives no stable guarantee about that sequence. The number of workers assigned to a query can vary with cluster size and concurrency; the set of files or blocks a worker reads and the order it reads them can vary between runs; and the merge tree that combines partial summaries follows whatever shape the scheduler produced. Feed the same multiset of values through two different compaction paths and you get two summaries that are both valid within their error bound but not identical — hence a p99 that shifts by a small amount on refresh. This is worth contrasting with distinct-count sketches. A HyperLogLog merge is a per-register maximum, which is commutative, associative and idempotent, so the merged sketch is bit-identical no matter how the work was split. Distinct-count estimates are reproducible; quantile estimates generally are not. ## Rank error is not value error The most commonly misread part of the guarantee is what the error applies to. The bound is on **rank**, not on the value. If the returned value sits at true rank 0.9897 instead of 0.9900, that is a tiny rank error — but on a long-tailed latency distribution the gap between the 98.97th and 99th percentile in milliseconds can be substantial. On a flat part of the distribution the same rank error is invisible. So the practical accuracy of an approximate p99 depends on the shape of the data as much as on the sketch. This is precisely why t-digest allocates its resolution to the tails: a uniform rank error is a poor deal at extreme quantiles where the value gradient is steepest. ## What to do about it **Raise the accuracy parameter** where the engine exposes one (compression for a digest, `k` for a KLL). More retained points means less compaction and a tighter bound, paid for in memory and merge cost. **Compute exactly when the group is small.** Percentiles over a few million values in one group sort cheaply; the approximation earns its keep on very large groups or on many groups at once. A hybrid — exact for small partitions, approximate above a threshold — is a defensible design but introduces its own discontinuity in the reported number. **Do not build hard thresholds inside the error band.** An SLA alert that fires when p99 crosses 250 ms will flap if the estimate wobbles by a few milliseconds around that line. Either widen the trigger with hysteresis, require several consecutive breaches, or compute that particular number exactly. **Publish the tolerance.** If a dashboard's percentile is approximate, say so next to it. Analysts who compare today's refresh with yesterday's screenshot and see a 0.3% shift will otherwise open an incident. **Beware combining percentiles.** Averaging daily p99s, or taking a percentile of percentiles, is not a valid operation regardless of approximation. If you need a monthly p99, merge the retained quantile summaries — the same rollup pattern used for distinct counts — rather than aggregating the finished numbers.
- Why is an approximate distinct count reproducible when an approximate percentile is not?Because their merge operations differ. A HyperLogLog register keeps a maximum, which is order-independent and idempotent, so any split of the work yields the identical sketch. A quantile summary compacts by merging neighbouring points, and which neighbours meet depends on arrival and merge order, so different plan shapes produce different — though equally valid — summaries.
- Can you compute a monthly p99 by averaging the daily p99s?No, and that has nothing to do with approximation: percentiles are not additive or averageable, since a busy day with high latency contributes more mass than a quiet one. The correct approach mirrors distinct counts — store the day's quantile summary in the rollup and merge the summaries, then extract the monthly percentile once.
- When is an approximate percentile clearly the wrong tool?When the number is contractual or the decision sits on a knife edge — an SLA credit triggered by a p99 threshold, a regulatory latency report, or a comparison between two variants whose true difference is smaller than the error band. In those cases either compute exactly over the relevant slice or widen the decision rule so the error cannot flip it.
saying these in an interview costs you the question
- Thinks the function samples a different subset of rows each run
- Reads the error bound as a percentage of the returned value
- Averages daily percentiles to get a monthly one
- Sets an alert threshold inside the estimate's error band
- Assumes higher accuracy is free rather than paid in memory