Why can an Elasticsearch terms aggregation report wrong doc_count values on a multi-shard index?
answer
- Aggregations run per shard, then merge
- Shards return candidates, not everything
- Truncation happens before the coordinating node sees it
- Two response fields quantify what was lost
- shard_size is the accuracy dial
basics
~20 sEach shard independently returns only its own top candidate terms, and the coordinating node sums those partial lists. A term ranked low on one shard contributes nothing from it, so counts can be undercounted or the term missed entirely.
solid answer
~50 sA `terms` aggregation runs per shard and then merges. Each shard sorts its own terms and returns only its top `shard_size` candidates — it never ships the full term list, because that would defeat the point. The coordinating node adds up what it received, so if a term is globally in the top 10 but sits at rank 900 on one busy shard, that shard's contribution is silently dropped and the merged `doc_count` is too low. Elasticsearch reports this honestly: `doc_count_error_upper_bound` bounds how far off the counts may be, and `sum_other_doc_count` says how many matching documents fell outside the returned buckets. The dials are `shard_size` (default derived from `size`, roughly `size * 1.5 + 10`) and `size`. Counts are exact when the index has a single shard, or when routing puts every document for a given term on one shard.
code
json · 13 lines{
"size": 0,
"aggs": {
"by_customer": {
"terms": {
"field": "customer_id",
"size": 10,
"shard_size": 500,
"show_term_doc_count_error": true
}
}
}
}go deeper
Recall that aggregation counts on an index with several shards are approximate by default, and that the response includes fields telling you how approximate.
Explain the per-shard candidate truncation end to end, and distinguish sum_other_doc_count from doc_count_error_upper_bound without hesitating.
Diagnose a real mismatch: check the error bound, raise shard_size with a cost argument, and know when routing alignment or a composite aggregation is the correct structural fix rather than more tuning.
Decide where approximate grouping is acceptable at all. Own the split between exploratory dashboards that tolerate error and reconciliation paths that must be exact, and the pre-aggregation strategy that serves the latter.
## Two phases, and where the truth is lost A search that carries a `terms` aggregation is broadcast to one copy of every shard. Each shard computes the aggregation over *its own* documents and returns a partial result. The coordinating node then reduces those partials into the final response. Nothing about that is unusual — the problem is what the shard is allowed to send. Sending every distinct term with its count would make the aggregation as expensive as shipping the field itself, so each shard sorts its terms by the requested order and returns only the top `shard_size` of them. That truncation is the entire source of the inaccuracy. If term `X` is the 3rd most common term globally but ranks 900th on shard B because shard B holds a different slice of traffic, and `shard_size` is 40, then shard B contributes **zero** to `X`'s count. The merged number is wrong, and it is wrong in one direction: counts can only be too low, never too high. ## The two bookkeeping fields Elasticsearch does not hide this. Two fields in the response describe it, and interviewers love the distinction: - **`sum_other_doc_count`** — how many matching documents belong to terms that were *not returned* to you. It is about the tail you did not see. It is exact, and it is not an error measure. - **`doc_count_error_upper_bound`** — an upper bound on how far the *returned* counts might be from the truth, computed from the count of the last term each shard returned: any term a shard omitted must have had at most that many documents on that shard. If the value is `0`, the returned counts are exact. Elasticsearch can only compute this bound when the aggregation is ordered by descending doc count; ordered by term or by a sub-aggregation, the bound is not meaningful. Adding `"show_term_doc_count_error": true` also emits a per-bucket error, so you can see which specific bucket is untrustworthy rather than a single global bound. ## When counts are exact - **One shard.** A single-shard index sees everything, so there is no merge and no truncation; the error bound is 0. - **Routing alignment.** If documents are routed so that every document sharing the aggregated term lands on the same shard — for example routing by `customer_id` and aggregating on `customer_id` — then each term is fully counted on exactly one shard, and merging cannot lose anything. - **Low cardinality.** If the field has fewer distinct values than `shard_size`, no shard ever truncates. ## The dials, and what they cost `shard_size` is the number of candidate terms each shard returns. It defaults to a value derived from `size` — roughly `size * 1.5 + 10` — which is deliberately a little larger than `size` so the merge has slack. Raising it is the direct accuracy fix: more candidates per shard means fewer silently dropped contributions. The cost is real but bounded — more per-shard sorting, more network payload, more memory on the coordinating node during the reduce — and it is far cheaper than raising `size` itself, because `shard_size` affects only the intermediate stage, not the returned buckets. A practical recipe: keep `size` at what the user actually needs, push `shard_size` up until `doc_count_error_upper_bound` drops to 0 or to a level you can defend, and check `sum_other_doc_count` to see whether the tail matters at all. ## Ordering by a sub-aggregation makes it worse `{"order": {"avg_price": "desc"}}` asks each shard to pick its local candidates by a metric computed from local documents only. A term whose global average is high but whose local sample on one shard is unremarkable never becomes a candidate there, and its documents are lost from the merge — so both the ordering *and* the counts can be wrong, and Elasticsearch cannot even bound the error. Treat sub-aggregation ordering as approximate by construction, and if the ranking must be right, either enlarge `shard_size` substantially or compute over a candidate set you have already narrowed with a filter. ## When approximation is not acceptable If you need exhaustive, exact grouping — a reconciliation report, a billing rollup — stop using `terms` for it. The `composite` aggregation walks *all* buckets in composite-key order with `after_key` paging, so nothing is truncated; you pay by streaming every bucket rather than getting a cheap top-N. Alternatively pre-aggregate: a rollup or a transform that materialises per-key totals into their own index turns an approximate query into an exact lookup. ## The interview shape The question is usually posed as a bug report: "the dashboard's top-10 customers don't match the finance export." The expected answer names the per-shard truncation, points at `doc_count_error_upper_bound` and `sum_other_doc_count` as evidence, offers `shard_size` as the tuning knob and routing or `composite` as the structural fixes — and, crucially, does not claim the counts are randomly wrong. They are systematically *under*-counted.
- What exactly does doc_count_error_upper_bound measure, and what does a value of 0 mean?It is an upper bound on how far the returned counts could be from the true global counts, derived from the doc count of the last term each shard returned — anything a shard dropped had at most that many documents there. Zero means no shard truncated in a way that could affect the returned buckets, so those counts are exact. It is only meaningful when ordering by descending doc count.
- How does sum_other_doc_count differ from doc_count_error_upper_bound?They answer different questions. `sum_other_doc_count` is exact and tells you how many matching documents belong to terms you were not shown — the size of the tail. `doc_count_error_upper_bound` is a bound on how wrong the counts of the terms you *were* shown might be. A big tail with zero error means the numbers are right but incomplete; a small tail with non-zero error means the opposite.
- Why is raising shard_size usually cheaper than raising size?`size` grows the buckets that are returned, sorted and held for the whole response, and it also drags `shard_size` up with it. `shard_size` alone only widens the intermediate candidate list each shard produces and the coordinating node merges — extra transient sorting and network payload, discarded once the reduce finishes. So you buy accuracy without growing the final result set.
saying these in an interview costs you the question
- Claiming terms counts are always exact
- Thinking sum_other_doc_count is the error estimate
- Believing adding replicas improves aggregation accuracy
- Assuming counts can be over-reported as well as under-reported
- Ordering by a sub-aggregation and trusting the ranking