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.
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
- Watch
- Try it
- Predict
- Capture
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.
Highlighted lines are the ones running in the diagram right now.
def local_score(doc, query, shard_stats):# PER SHARD: N_local, df_local — never cluster-wide by defaultN_local = shard_stats.Ndf_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