Top-N aggregations can miss the winner
A terms aggregation asks each shard for its local top-shard_size buckets and merges; if the global #1 was rank 6 on every shard and shard_size is 5, the coordinator never sees it — wrong by omission. shard_size is the knob that buys accuracy with per-shard memory.
Scoring drift was approximation per document — same data, reranked. The same per-shard truncation reappears in aggregations, with a sharper consequence: not just reordered, but missing.
Scene 10
Top-N aggregations can miss the winner
- Watch
- Try it
- Predict
- Capture
We're running a terms agg on the tags field across 5 shards (~100k books each), asking for the top 5 most-popular tags globally. By default each shard returns its local top 5. Watch the merge: 5 shards × 5 buckets = 25 candidates, coordinator picks the global top 5. Looks clean — but one tag is silently missing.
Highlighted lines are the ones running in the diagram right now.
def terms_agg(field, size, shard_size):per_shard = []for shard in cluster.shards_for(index):per_shard.append(shard.local_top_k(field, k=shard_size))merged = merge_buckets(per_shard)return merged[:size]
def local_top_k(field, k):counts = Counter()for doc in self.docs:for term in doc[field]:counts[term] += 1return counts.most_common(k) # tail truncated here
def merge_buckets(per_shard):totals, seen = {}, {}for shard_idx, buckets in enumerate(per_shard):for term, count in buckets:totals[term] = totals.get(term, 0) + countseen.setdefault(term, set()).add(shard_idx)out = []for term, count in sorted(totals.items(), key=-count):# error: shards that did NOT return this term could hide# up to their smallest-returned count for it.err = sum(per_shard[s][-1].countfor s in range(len(per_shard))if s not in seen[term])out.append({term, count, doc_count_error_upper_bound: err})return out
Where this sits in Build a distributed search engine (Elasticsearch / OpenSearch style)
Scene 10 of 12. Each shard returns its own top shard_size buckets and the coordinator merges — a term that is 11th on every shard is invisible in a top-10 result, and shard_size is the knob.
Up next. Every approximation so far is mathematical. The next class of failures is operational — knobs from earlier scenes turned past where their defaults still hold.
All 12 scenes in Build a distributed search engine (Elasticsearch / OpenSearch style) · Every curriculum