Distributed scoring is approximate

BM25's IDF is computed per shard, not cluster-wide, so the same document scores differently depending on where it lives — dfs_query_then_fetch fixes this with one extra round trip that pre-aggregates df across all shards.

Previously

Phase 1 of every search asked each shard to BM25-rank its local matches. The IDF those scores depend on is also LOCAL — and that's where 'same data, same query, same answer' quietly stops being true.

Scene 09

Distributed scoring is approximate

  1. Watch
  2. Try it
  3. Predict
  4. Capture
COORDINATORCoordinatorphase 2 · searchfetchfetchfetchshard-0top-K candsLOCAL STATSn=1000000 df=80000doc 10219.4doc 10315.8shard-1top-K candsLOCAL STATSn=1000000 df=1000doc 10253.1doc 20749.3shard-2top-K candsLOCAL STATSn=1000000 df=50000doc 31124.2doc 31813.6GLOBAL TOP-N · merged at coordinatorscore#1 doc 10253.1 shard-1#2 doc 20749.3 shard-1#3 doc 31124.2 shard-2#4 doc 10219.4 shard-0#5 doc 10315.8 shard-0SEARCH LATENCYphase 1 only42 ms42 ms scatter-gather · no preflightEach shard scores with LOCAL (N, df). Same doc text, different shards → different scores.
What to watch for

Three shards. Each one BM25-ranks its local matches for 'distributed' using its OWN df. Look at the per-shard df badges — shard 0 has 80k matches, shard 1 has just 1k, shard 2 has 50k. Same query, same term, three different IDFs.

Implementation

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

Shard.localScore
BM25 with PER-SHARD (N, df) — the source of drift
def local_score(doc, query, shard_stats):
# PER SHARD: N_local, df_local — never cluster-wide by default
N_local = shard_stats.N
df_local = shard_stats.df[query.term]
idf = ln((N_local - df_local + 0.5) / (df_local + 0.5) + 1)
tf_norm = doc.tf / (doc.tf + 1.2)
return tf_norm * idf

Where this sits in Build a distributed search engine (Elasticsearch / OpenSearch style)

Scene 09 of 12. Every shard computes IDF from its own local corpus, so the same document scores differently depending on where it lives — dfs_query_then_fetch buys correctness for one extra round trip.

Up next. Score drift is the per-document version of a more general problem. Ask each shard for its top-10 of anything and merge — the global #1 can be invisible if it was #11 on every shard.

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