skip to content

When does search_type=dfs_query_then_fetch change an Elasticsearch result ordering, and what does it cost?

level: seniorimportance: should knowfreq 45%

answer

  1. Shards score independently by default
  2. Document frequency is a per-shard number
  3. One extra round trip before the query phase
  4. Global term statistics gathered, then scoring
  5. Matters most on tiny or skewed shards

basics

~20 s

By default each shard scores using only its own term statistics, so identical documents can score differently per shard. dfs_query_then_fetch adds a preliminary round trip that gathers term and document frequencies from all participating shards so scoring uses global statistics.

solid answer

~50 s

In the default `query_then_fetch`, every shard computes relevance from **its own** term statistics — how many documents on that shard contain the term, and how long that shard's fields are on average. The coordinating node then merges scores that were produced against different statistical baselines. When documents are spread evenly over well-populated shards the baselines converge and nobody notices. When they are not — tiny shards, few documents overall, custom routing concentrating a tenant, or one shard carrying many undeleted deletions — the same term looks rare on one shard and common on another, and the merged ranking is wrong. `search_type=dfs_query_then_fetch` inserts a **DFS** round trip before the query phase: the coordinator asks each participating shard for its term and document frequencies, sums them into global statistics, and hands those back so the query phase scores every shard against the same baseline. The cost is one extra network round trip to every shard on every search, plus the coordinating work to aggregate the statistics.

code

bash · 5 lines
bash
# same query, global term statistics: use as a diagnostic, not a default
GET /products/_search?search_type=dfs_query_then_fetch
{
  "query": { "match": { "description": "waterproof jacket" } }
}

go deeper

for a junior

Know that a search is scored on each shard separately and that an alternative search type exists which collects statistics across shards first.

for a middle

Explain that document frequency is computed per shard, so the same term can look rare on one shard and common on another, and that the DFS phase sums those frequencies before scoring.

for a senior

Diagnose it: recognise the small-corpus and custom-routing signatures, use the search type as a one-request experiment, and prefer fixing shard count or distribution over paying a round trip on every query.

for a principal

Own the tradeoff between ranking fidelity and per-query latency, and set the policy on where relevance-critical data lives — a single-shard index removes the problem entirely rather than paying to work around it.

## The problem it solves Relevance scoring rewards rare terms. The rarity component is derived from **document frequency**: how many documents in the corpus contain the term, relative to the corpus size. In a distributed index there is no cheap global corpus — each shard is its own little Lucene index with its own term dictionary and its own document count. So in the default execution, `query_then_fetch`, each shard answers the question "how rare is this term?" locally. Shard A might hold 40 documents containing *kubernetes* out of 10,000; shard B might hold 3 out of 10,000. Shard B concludes the term is far rarer and scores its matches much higher. The coordinating node merges the two lists as if the numbers were comparable. They are not. ## When the skew is visible - **Small corpora.** With a few hundred documents spread over five shards, per-shard frequencies are dominated by noise. This is the classic case where a developer tests relevance locally, sees nonsense rankings, and concludes the query is broken. - **Uneven document distribution.** Custom routing deliberately concentrates a routing key's documents on one shard. If that key's documents use a distinctive vocabulary, that shard's statistics diverge sharply from the rest. - **Uneven deletions.** Deleted documents still count in a shard's statistics until the segments holding them are merged away. A shard that just absorbed a bulk update can carry a large deleted population and skew its own baseline. - **Time-based indices with different volumes.** Yesterday's index holding ten million documents and today's holding fifty thousand produce very different frequencies for the same term. ## What the DFS phase actually does DFS stands for *distributed frequency search*. Setting `search_type=dfs_query_then_fetch` on the search request turns the two-phase execution into three: 1. **DFS phase.** The coordinating node asks every participating shard for the term-level statistics needed by the query — document frequencies for the query terms, plus the collection-level counts. It sums them into a global view. 2. **Query phase.** The same fan-out as usual, except each shard is handed the global statistics and scores its documents against them. Now every shard's scores share one baseline, so the merge is meaningful. 3. **Fetch phase.** Unchanged: the winners' `_source` is retrieved from the shards that own them. Only scoring changes. Which documents *match* is identical; only their relative order can move. ## The cost One additional round trip to every shard the search touches, before any real work begins. On a search that fans out to a handful of shards on a fast network this is small; on a search that fans out to hundreds of shards, or across a cluster with cross-datacentre latency, it adds a full network hop to every query's latency, plus coordinating-node work to aggregate statistics. Because the extra phase is per-request, the cost is paid on every search, forever — it is not amortised or cached. That is why it is not the default. Elasticsearch's designers bet that with reasonably sized, reasonably balanced shards, local statistics are close enough to global ones that the ranking difference is undetectable, and the round trip is not worth paying. ## When to reach for it, and what to reach for instead Use it when you have measured a ranking problem attributable to statistical skew and you cannot fix the distribution. Before that, consider the cheaper structural fixes: - **Fewer shards.** An index that comfortably fits one shard has exactly one set of statistics and no skew at all. Many relevance-sensitive catalogues are small enough for this. - **Even distribution.** Do not use custom routing on an index whose relevance ranking matters across routing keys. - **Merge away deletions** on a static index so deleted documents stop inflating frequencies. - **Test relevance on a single-shard index**, so you are debugging your query rather than your shard layout. A useful diagnostic habit: if a query ranks correctly on a one-shard copy of the data but badly on the sharded production index, statistical skew is the prime suspect, and running the same search with `search_type=dfs_query_then_fetch` confirms it in one request. ## What interviewers listen for A strong answer states plainly that scoring is per-shard by default, names document frequency as the statistic that diverges, describes the extra round trip concretely, and — most importantly — treats DFS as a diagnostic and a last resort rather than a setting to turn on everywhere. A candidate who proposes enabling it cluster-wide "for better relevance" has not costed the round trip.

  • Does dfs_query_then_fetch change which documents match, or only their order?
    Only their order. The matching logic is unchanged — the same documents satisfy the query either way. What changes is the statistical baseline used to score them, so relative ranking can shift and, with a page size smaller than the result set, different documents can appear on page one.
  • You get sensible rankings in a one-shard test index but poor ones in production with ten shards. What is your first hypothesis?
    Per-shard statistical skew. The single-shard index has one global set of term frequencies, so scores are directly comparable; ten shards give ten baselines that the coordinating node merges as if they matched. Confirm by re-running the production search with `search_type=dfs_query_then_fetch` and comparing the ordering.
  • Why do deleted documents affect scoring at all?
    A delete is recorded as a marker; the document stays in its segment until a merge rewrites that segment without it. Until then it still contributes to the shard's document count and term document frequencies, so a shard carrying many pending deletions reports a distorted rarity for its terms.

saying these in an interview costs you the question

  • Says dfs_query_then_fetch returns extra matching documents
  • Proposes enabling it globally for better relevance
  • Cannot name document frequency as the skewed statistic
  • Thinks the extra phase fetches document bodies
  • Believes scores are already global across shards by default

context