Why is Elasticsearch's cardinality aggregation approximate, and what does precision_threshold trade off?
answer
- exact distinct counting must remember every value
- the memory must not grow with uniqueness
- sketches from different shards combine cleanly
- one parameter buys accuracy with bytes
- there is a hard ceiling on that parameter
basics
~20 sThe cardinality aggregation estimates distinct counts with a HyperLogLog++ sketch of fixed size instead of tracking every value, so memory stays bounded as cardinality grows. precision_threshold trades memory for accuracy: higher means near-exact counts up to a larger unique count.
solid answer
~50 sCounting distinct values exactly means remembering every value you have seen, which costs memory proportional to the cardinality itself and does not merge cheaply across shards. Elasticsearch instead hashes each value into a **HyperLogLog++** sketch of fixed size and estimates the count from the hash distribution, so a field with a billion distinct values uses the same memory as one with a thousand. The `precision_threshold` parameter sets the sketch size: counts below it are expected to be very close to exact, above it error grows slowly. It defaults to 3000 and the maximum supported value is 40000, with memory roughly `precision_threshold × 8` bytes per counter. Because sketches merge, each shard builds its own and the coordinating node combines them without double counting. The estimate is deterministic for a given set of values, so a wrong number reproduces exactly on every run.
code
json · 11 lines{
"size": 0,
"aggs": {
"unique_users": {
"cardinality": {
"field": "user_id",
"precision_threshold": 40000
}
}
}
}go deeper
Know that cardinality gives an estimated distinct count, not an exact one, and that this is deliberate rather than a bug.
Explain the HyperLogLog++ sketch: fixed memory regardless of cardinality, mergeable across shards, and precision_threshold trading bytes for accuracy with a documented default and ceiling.
Bring operational judgment: memory is per bucket per shard, so nesting cardinality under a wide terms aggregation multiplies it, and there is no per-response error bound to inspect when a number is challenged.
Own the accuracy contract. Decide which consumers may be served approximate uniques, which need an exact pipeline, and how the platform labels and reconciles estimated figures so no one bills from one.
## Why exact distinct counting is expensive To count distinct values exactly you must remember which values you already saw — a hash set whose size grows with the cardinality, not with the data volume. On a distributed index that is worse than it sounds: each shard would have to ship its full set of values to the coordinating node, because two shards seeing the same value must not count it twice. Summing per-shard distinct counts is simply wrong. So an exact distinct count across shards means moving the values themselves, with memory and network cost proportional to cardinality. ## What HyperLogLog++ does instead The `cardinality` aggregation hashes each value and keeps a fixed-size array of small registers, each recording the maximum number of leading zero bits seen in the hashes that mapped to it. Long runs of leading zeros are rare, so seeing one is evidence that many distinct hashes passed through. From the distribution of register values the algorithm estimates the number of distinct inputs. HyperLogLog++ is Google's refinement of HyperLogLog, adding bias correction and a sparse representation that is near-exact for small cardinalities. Three properties matter for interviews: 1. **Fixed memory.** The sketch's size depends on the configured precision, never on how many distinct values arrive. 2. **Mergeable.** Two sketches combine by taking the element-wise maximum of their registers, which is exactly the sketch you would have built from the union. That is what makes distributed counting cheap: each shard builds a sketch, the coordinating node merges them, and duplicates across shards collapse naturally. 3. **Deterministic.** The estimate is a function of the set of hashed values. The same documents produce the same estimate on every run — the error is a stable bias, not run-to-run noise. A user who checks a number twice will see the same wrong answer, which is why "just retry it" is not a diagnosis. ## precision_threshold ```json { "aggs": { "unique_users": { "cardinality": { "field": "user_id", "precision_threshold": 40000 } } } } ``` `precision_threshold` is the unique count below which results are expected to be close to accurate. It defaults to 3000 and the maximum supported value is 40000; values above that are capped. Memory is roughly `precision_threshold × 8` bytes per counter, so the maximum setting costs on the order of a few hundred kilobytes — per aggregation, per shard. That last clause is where operators get hurt. If a `cardinality` aggregation is nested under a `terms` aggregation with ten thousand buckets, you are asking each shard for ten thousand sketches. Raising precision on a nested cardinality is how a harmless-looking dashboard query starts tripping circuit breakers. Raising the threshold does not buy exactness above it. With `precision_threshold: 40000` and 50,000 actual distinct users, the answer is still an estimate — the setting only moves where meaningful error starts. ## Where the error shows up Error is relative, so it looks small in percentage terms and large in absolute ones: a fraction of a percent on ten million uniques is still tens of thousands of users. There is no per-response error bound for `cardinality` — the response carries a single `value` with no confidence interval — so you cannot inspect a given result and tell how far off it is. The operational consequence is that approximate distinct counts are fine for trends, dashboards, anomaly spotting and relative comparison, and unsuitable as the number on an invoice, a compliance report, or anything a customer will reconcile against their own records. ## Practical notes - `cardinality` accepts `missing`, so documents lacking the field can be folded in under a substitute value. - It counts distinct values, and multi-valued fields contribute each of their values, so "distinct tags" across an array field behaves as you would expect. - It can run over a script or a runtime field rather than an indexed field, at the cost of evaluating the script per value. - It aggregates the *indexed* values, so if a `keyword` field was normalized at index time (lowercased, for instance), the distinct count is over the normalized forms, not the raw source strings. Counting distinct values of an analyzed `text` field is almost never what anyone means and requires fielddata, which you should not enable. ## The answer an interviewer wants Name the algorithm, explain the bounded-memory and mergeability motivation, state that `precision_threshold` trades memory for accuracy with a hard ceiling, and finish on judgment: know which consumers of the number can tolerate approximation and have a plan — exact enumeration, a precomputed entity-per-document index, or an offline pipeline — for the ones that cannot.
- Does setting precision_threshold to 40000 guarantee an exact count for 50,000 distinct users?No. The threshold is the unique count below which results are expected to be close to accurate; above it the estimate degrades gradually. 40000 is also the maximum supported value, so there is no setting that makes the aggregation exact for arbitrary cardinality. If exactness is required you need a different approach entirely, not a bigger sketch.
- Why can't you sum the cardinality results from each shard to get the cluster-wide distinct count?A value present on three shards would be counted three times. Elasticsearch avoids this by merging the shards' sketches — taking the element-wise maximum of their registers — which yields the sketch of the union rather than the sum of the parts. That mergeability is the main reason a sketch is used at all.
- Why does the same cardinality query return the identical wrong number every time?The estimate is a deterministic function of the hashed values, not a random sample. The same document set always produces the same sketch and therefore the same estimate. Error is a stable bias of the algorithm, so re-running a query is not a way to detect or work around it.
It is like estimating a crowd's size from the rarest ticket-stub pattern you spot rather than collecting every stub: constant effort no matter how big the crowd, and two gate staff can combine their observations without double counting anyone.
saying these in an interview costs you the question
- Claims cardinality is exact, just slow
- Thinks raising precision_threshold above 40000 has an effect
- Sums per-shard cardinality results to get a cluster total
- Says results vary randomly between runs
- Nests a high-precision cardinality under a huge terms aggregation without concern