Why can two identical documents get different BM25 scores on a multi-shard Elasticsearch index?
answer
- Each shard scores with what it can see
- The coordinating node merges, it does not rescore
- idf depends on a document count that is local
- One search_type value fixes it, at a price
- Small indices suffer most
basics
~20 sBM25's term statistics are gathered per shard, not cluster-wide. Each shard computes idf and average field length from its own documents, so identical documents on different shards get different scores. Running the search with search_type=dfs_query_then_fetch gathers global statistics first.
solid answer
~50 sSearch executes per shard, and each shard scores using only its own term statistics: the number of documents containing the term, the number of documents with the field, and the average field length. Those differ from shard to shard, so a duplicate document on shard 0 and shard 3 gets two different scores, and the coordinating node simply merges the already-scored results. The skew is worst on small or unevenly-routed indices, where one shard may hold two documents containing the term and another two hundred. Fixes, in order of preference: use a single shard for small indices, since the statistics are then genuinely global; let large indices average the skew out, which they generally do; or pass `search_type=dfs_query_then_fetch`, which adds a preliminary round trip collecting global term statistics before the query phase, at the cost of extra latency.
code
bash · 4 linesGET /articles/_search?search_type=dfs_query_then_fetch
{
"query": { "match": { "title": "elasticsearch" } }
}go deeper
Know that Elasticsearch searches each shard separately and that relevance scores are computed there, not centrally.
Explain that idf and average field length come from per-shard statistics, and name search_type=dfs_query_then_fetch as the setting that collects them globally first.
Diagnose it in the wild: flaky relevance tests, duplicates ranking inconsistently, skew from custom routing, and drift from deleted-but-unmerged documents. Prefer one shard for small indices over a DFS phase everywhere.
Own the tradeoff explicitly — Elasticsearch sacrifices globally consistent scoring for scale-out, and deciding where to buy that consistency back with latency is a capacity and product-quality call, not a default to flip.
## Why the statistics are local An Elasticsearch search fans out to one copy of each shard. Each shard runs the query against its own Lucene index and returns its top hits **already scored**. The coordinating node merges those ranked lists; it does not rescore anything. That means every score was computed from the statistics available on one shard: - `n` — how many documents *on that shard* contain the term. - `N` — how many documents *on that shard* have the field. - `avgdl` — the average length of that field *on that shard*. BM25's idf term is `log(1 + (N - n + 0.5) / (n + 0.5))`, so a term that is rare on shard 0 and common on shard 3 gets a large idf on one and a small idf on the other. Two byte-identical documents therefore score differently depending on where routing put them, and the one on the shard where the term happens to be rarer wins. ## When it actually bites - **Small indices.** With five shards and 50 documents, per-shard counts are tiny and the variance is enormous. This is why relevance behaves "wrongly" in tests and demos far more often than in production. - **Uneven routing.** Custom routing that concentrates a tenant on one shard makes that shard's statistics unrepresentative of the corpus. Default `_id`-hash routing spreads documents evenly enough that the law of large numbers takes over once shards hold enough documents. - **Skewed content per shard.** If routing correlates with content — for example routing by language or by category — then term distributions genuinely differ per shard, and no amount of volume averages it out. - **Duplicates.** Near-duplicate documents that land on different shards will rank inconsistently, which users notice as "the same product appears twice, in different places". A related source of drift: term and collection statistics include documents that are deleted but whose segments have not yet been merged. Delete a large batch and scores shift slightly until merging catches up; on a static index a force merge stabilises them. ## dfs_query_then_fetch The explicit remedy is the `search_type=dfs_query_then_fetch` query parameter. It inserts a **distributed frequency search (DFS)** phase before the query phase: the coordinating node asks every shard for its term and document-frequency statistics, sums them into global values, and passes those to the shards so they all score with the same idf and average field length. Ranking then matches what a single-shard index would produce. The cost is one extra round trip to every shard before any scoring happens, which shows up directly in latency and grows with shard count. It is not a free correctness switch. Typical use: turn it on for relevance debugging and for automated relevance tests so results are deterministic; leave it off for high-volume production traffic unless you have measured the latency and decided the ranking stability is worth it. ## The pragmatic ordering of fixes 1. **Give small indices one primary shard.** If the whole corpus fits comfortably in a single shard, the statistics are global by construction and the problem disappears with no query-time cost. Most catalogue-sized indices are in this category. 2. **Let volume do the work.** On a large, evenly-routed index, per-shard statistics converge and the skew becomes negligible. Do not add a DFS phase for a theoretical problem you cannot measure. 3. **Use dfs_query_then_fetch deliberately.** For relevance test suites, for debugging a reported ranking bug, or for genuinely skewed corpora where a business requirement demands consistent ordering. 4. **Do not fight it with boosts.** Compensating for shard skew by tweaking field boosts hides the cause and makes the ranking worse when the data distribution changes. ## Interview framing The test here is whether you understand that Elasticsearch trades global scoring correctness for scalability by default, and that this is a deliberate design choice rather than an oversight. The strong answer names the per-shard statistics explicitly, explains that the coordinating node merges pre-scored results, gives `dfs_query_then_fetch` as the correction with its latency cost, and adds that shrinking to one shard is usually the better fix for the small indices where the problem is actually visible.
- What exactly does the DFS phase cost you?One extra round trip to every shard before scoring begins, plus the work of gathering and summing term and document-frequency statistics. Latency grows with shard count and it is paid on every request. That is why it is normally reserved for relevance tests and debugging rather than enabled for production search traffic.
- Why do relevance problems from shard skew show up in tests but rarely in production?Test fixtures hold a handful of documents spread over the index's shards, so per-shard document counts are tiny and idf varies wildly between them. A production index with millions of evenly routed documents gives each shard a nearly identical term distribution. Pinning test indices to one primary shard removes the flakiness at the source.
- Can custom routing make the skew permanent rather than transient?Yes. If routing correlates with content — one tenant, language or category per shard — the shards genuinely hold different term distributions, so extra volume never averages them out. Those corpora are the real case for a DFS phase, or for splitting into separate indices whose statistics are meant to be separate.
- Why do scores change slightly after deleting a batch of documents?Deleted documents are only marked, not removed, so their terms still count toward the document frequency and average field length until the segments holding them are merged. Until then idf and length normalisation reflect documents no longer visible. A force merge on a static index removes them and stabilises the scores.
It is like several judges scoring contestants in separate rooms and then merging the sheets: each judge calibrates against only the contestants in front of them, so identical performances get different marks.
saying these in an interview costs you the question
- Believes idf is computed cluster-wide by default
- Says the coordinating node rescores the merged hits
- Enables dfs_query_then_fetch everywhere with no latency thought
- Blames replicas or analyzers for the score difference
- Compensates for shard skew by tuning field boosts