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.

Previously

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

  1. Watch
  2. Try it
  3. Predict
  4. Capture
COORDINATORcoordinator · merge top bucketsphase 2 · aggshard0top-N bucketsLOCAL STATSn=100000 df=?classic · 80kchildren · 78ksystems · 76kfiction · 74kdatabases · 72kshard1top-N bucketsLOCAL STATSn=100000 df=?systems · 80kfiction · 78knetworking · 76kjava · 74ksearch · 72kshard2top-N bucketsLOCAL STATSn=100000 df=?fiction · 80kchildren · 78ksearch · 76kjava · 74knetworking · 72kshard3top-N bucketsLOCAL STATSn=100000 df=?systems · 80kdatabases · 78knetworking · 76ksearch · 74kjava · 72kshard4top-N bucketsLOCAL STATSn=100000 df=?fiction · 80kchildren · 78kdatabases · 76kjava · 74knetworking · 72kGLOBAL TOP-N · merged at coordinatordoc_count#1 fiction · 312k · err ≤ 72k#2 networking · 296k · err ≤ 72k#3 java · 294k · err ≤ 72k#4 systems · 236k · err ≤ 144k#5 children · 234k · err ≤ 144k(no pagination meter)terms agg on `tags` · size=5 · shard_size=5.
What to watch for

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.

Implementation

Highlighted lines are the ones running in the diagram right now.

Coordinator.termsAgg
fan out for top-shard_size; merge; slice top-size
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]
Shard.localTopK
count + sort + slice — anything below k is invisible
def local_top_k(field, k):
counts = Counter()
for doc in self.docs:
for term in doc[field]:
counts[term] += 1
return counts.most_common(k) # tail truncated here
Coordinator.mergeBuckets
sum counts; sort; emit doc_count_error_upper_bound
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) + count
seen.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].count
for 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

Built with Arqly
Every scene in Build a distributed search engine (Elasticsearch / OpenSearch style) builds on the one before it.All 12 Build a distributed search engine (Elasticsearch / OpenSearch style) scenes