skip to content

What happens during the query phase and the fetch phase of an Elasticsearch search?

level: middleimportance: must knowfreq 65%

answer

  1. A search takes two round trips
  2. The first round trip moves no documents
  3. Each shard returns from+size ids and scores
  4. One node merges into a global queue
  5. Only the winners' _source is fetched

basics

~20 s

In the query phase each shard runs the search locally and returns only document ids plus sort values or scores; the coordinating node merges these into one globally sorted list. The fetch phase then retrieves the _source of just the winning documents.

solid answer

~50 s

A search arrives at any node, which acts as the **coordinating node**. It resolves the target indices to shards, picks one copy of each shard, and broadcasts the query. In the **query phase** every shard executes the query against its own Lucene segments and builds a local priority queue of size `from + size`, returning only doc ids and the sort values or `_score` — no document bodies move. The coordinating node merges those per-shard lists into a single sorted list of `from + size` entries and throws the rest away. In the **fetch phase** it issues multi-get requests for just the surviving entries to the shards that own them, which return `_source`, highlighting, stored fields and script fields. The coordinating node assembles the final response. Two round trips, and document bodies cross the network only for hits the user will actually see.

code

json · 6 lines
json
{
  "from": 90,
  "size": 10,
  "query": { "match": { "title": "elasticsearch" } },
  "highlight": { "fields": { "title": {} } }
}

go deeper

for a junior

Recall that a search fans out to every relevant shard and that the results are merged by one node before being returned to the client.

for a middle

Be ready to name both phases, state that the query phase returns only ids and sort values, and explain the from+size priority queue on each shard and on the coordinator.

for a senior

Show the operational consequences: the slowest shard sets the latency floor, deep pagination multiplies queue work across shards, and expensive highlighting is a fetch-phase cost bounded by page size.

for a principal

Frame the fan-out as the real capacity model — per-query cost scales with shards touched, so index layout, pre-filtering and routing are the levers that decide what your search tier can serve.

## Two round trips, by design Elasticsearch's default search execution is called `query_then_fetch`, and the name is the algorithm. Sorting a distributed result set requires seeing candidates from every shard, but shipping full documents from every shard would be enormously wasteful when the client asked for ten hits. So the engine separates *deciding which documents win* from *materialising them*. ## The coordinating node Any node can receive a search; for that request it becomes the coordinating node. Its first job is resolution: expand index patterns and aliases into concrete indices, apply alias filters, and produce the list of shards the search must touch. For each shard it then picks **one copy** — primary or replica — to serve the request. That choice is load-aware rather than round-robin, and it means two identical searches may execute on different physical copies. Before the query phase, the coordinator may run a **can_match** pre-filter round trip: a cheap question to each shard asking whether it could possibly hold a match, based on the range of values it actually stores. On time-based indices this lets a search over a month of daily indices skip most shards outright. Shards that cannot match are dropped and never enter the query phase. ## The query phase The coordinator sends the query to every remaining shard concurrently, subject to `max_concurrent_shard_requests`, which caps how many shard-level requests one search issues per node at a time. Each shard: 1. Executes the query against its own segments, using **its own local term statistics** to compute scores. 2. Maintains a priority queue of size `from + size`, keeping only the best entries. 3. Returns that queue to the coordinator — **doc ids and sort values or scores only**, plus its own hit count and any aggregation results. The crucial detail is what does *not* come back: no `_source`, no highlighted fragments, no stored fields. The query phase moves metadata. ## The merge The coordinating node merges the per-shard queues into one global priority queue of size `from + size`, sorted by score or by the requested sort keys. If a search asks for `from: 90, size: 10` across five shards, each shard sends its top **100** entries and the coordinator merges 500 candidates down to 100, then discards the first 90 and keeps 10. Aggregation results are reduced here too, shard result by shard result, into the final aggregation tree. This merge is where deep pagination gets expensive: the per-shard queues and the coordinator's queue both grow with `from + size`, and the cost is paid on every shard, for every request. Hit counts are also summed here. Since Elasticsearch 7.0 the coordinator stops counting exactly once it is certain the total exceeds 10,000, reporting `"relation": "gte"`; `track_total_hits` restores an exact count at a cost. ## The fetch phase With the winners known, the coordinator issues multi-get requests to the shards that own them — usually a small subset of the shards that participated in the query phase. Those shards load `_source` for the requested doc ids, run highlighting, evaluate script fields and stored fields, and return complete hits. The coordinator assembles them in the merged order and responds to the client. ## Why the split matters in practice - **Network cost is proportional to `size`, not to the number of matches.** A query matching ten million documents with `size: 10` transfers ten documents. - **Scoring happens per shard with local statistics.** Each shard computes term frequencies and document frequencies from its own data, which is why relevance can be slightly inconsistent across shards and why an alternative search type exists to gather global statistics first. - **Highlighting and `_source` decompression cost is bounded** by the page size, so an expensive highlighter on a big page is a fetch-phase problem, not a query-phase one. - **The slowest shard sets the latency.** The query phase cannot finish until every participating shard replies, so one hot node or one oversized shard determines p99 for the whole search. ## What interviewers listen for The answer that lands names both phases, states explicitly that the query phase returns ids and sort values rather than documents, describes the coordinating-node merge into a queue of `from + size`, and can then explain a real consequence — deep pagination cost, tail latency from the slowest shard, or per-shard scoring — without being prompted.

  • Why does from: 10000, size: 10 cost far more than from: 0, size: 10?
    Every participating shard must build and return a priority queue of 10,010 entries, and the coordinating node merges shard_count × 10,010 candidates before discarding the first 10,000. Both memory and network grow with the offset, on every shard, on every request — the work is repeated rather than resumed.
  • What does the can_match phase do before the query phase?
    It asks each shard, cheaply, whether it could contain a match at all — typically by checking whether the shard's stored range for the sorted or filtered field overlaps the query. Shards that provably cannot match are skipped entirely. On time-based indices this removes most shards from a search over a narrow time window.
  • Are aggregation results computed in the query phase or the fetch phase?
    In the query phase. Each shard computes its own partial aggregation over its matching documents and returns that structure alongside the top-hit metadata; the coordinating node reduces the partials into the final result. The fetch phase only materialises hit documents, so a `size: 0` aggregation search skips it entirely.

saying these in an interview costs you the question

  • Says every shard returns full documents to be sorted centrally
  • Thinks one shard runs the search and the others idle
  • Cannot say what the query phase actually transfers
  • Believes the coordinating node re-scores the merged hits
  • Claims scores are computed globally across all shards by default

context