skip to content

questions

5

In a search engine whose index is split across many shards, how does scatter-gather serve a query and return the global top 10?

level: juniorimportance: must knowfreq 60%

answer

  1. coordinator plus parallel sub-requests
  2. each shard answers only for itself
  3. bounded lists, not all matches
  4. heap or k-way merge
  5. a winner is a winner locally

basics

~10 s

A coordinator sends the query to every shard in parallel. Each shard scores its own documents and returns only its local top 10. The coordinator merges those lists and keeps the best 10 overall.

solid answer

~40 s

A **coordinator** node forwards the query to one replica of every shard in parallel (scatter). Each shard runs the query against its own slice of the index and returns its **local top-k**: document IDs, scores and sort values, not full documents. The coordinator merges the `S` partial lists with a bounded heap or a k-way merge and keeps the global top-k (gather). If scores are comparable across shards, the merge is exact: a document outranked by 10 others on its own shard is also outranked by those 10 globally, so every global winner appears in some shard's list. The fan-out has a price. Each query sends one sub-request per shard, the coordinator merges up to `S·k` candidates, and latency is set by the slowest shard it waits for.

code

pseudocode · 10 lines
pseudocode
function search(query, k):
  lists = parallel for shard in shards:
            shard.topK(query, k)      // (score, docId) pairs
  heap = empty min-heap ordered by (score, docId)
  for list in lists:
    for hit in list:
      heap.push(hit)
      if heap.size > k:
        heap.popMin()                 // drop the weakest candidate
  return heap.drainDescending()       // global top k, best first

go deeper

for a junior

Recall the two phases in order: the coordinator scatters the query to every shard, and each shard returns only its local top k. Then the coordinator gathers the lists and merges them into one global top k.

for a middle

Explain why the per-shard top-k is enough for an exact global top-k, name the heap or k-way merge and its cost, and describe why query-then-fetch avoids shipping full documents for candidates that lose.

for a senior

Show you know what fan-out costs in production: one sub-request per shard, coordinator memory proportional to shards times k, and latency set by the slowest shard. Mention routing keys that let a query skip shards.

for a principal

Frame shard count as a trade-off between smaller, faster shards and a costlier, more tail-sensitive fan-out. Explain when routing, fewer shards or a two-phase fetch changes the economics for a given query mix.

## Why a search index is split into shards Eventually one machine can no longer hold the whole **inverted index** (the map from each term to the documents containing it) or answer queries over it fast enough. The usual fix is **document partitioning**. The corpus is divided into **shards**, and each shard is a self-contained index over its own subset of documents. Each shard is copied to several **replicas** for availability and read throughput. Any shard might hold a matching document, so a plain keyword query has to consult all of them. The pattern for serving such a query is **scatter-gather**. ## The scatter phase A **coordinator** (also called a root, broker or aggregator) receives the query and runs these steps: 1. Parse the query once and pick one replica of every shard, choosing by health, load or recent latency. 2. Send the query to all of those replicas **in parallel**, with `k` (the number of results the caller wants) and a deadline. 3. Each shard evaluates the query against only its own documents, scores the matches and keeps its **local top-k** in a bounded heap. A shard does not return every match. A common term can match millions of documents, and shipping them all would saturate the network and the coordinator. Each shard returns at most `k` compact entries: a document ID, a score and any sort values the query needs. ## The gather phase: merging partial top-k lists The coordinator now holds `S` lists of at most `k` entries each, where `S` is the number of shards. There are two common ways to merge them into one global top-k: - **Bounded min-heap** of size `k`. Push every candidate and evict the weakest whenever the heap grows past `k`. This costs about `S·k·log k`. - **k-way merge.** Each list arrives already sorted, so the coordinator can keep a heap of the `S` list heads and pop only `k` times. This costs about `k·log S`. With 20 shards and `k = 10`, the coordinator sees at most 200 candidates and keeps 10. ```pseudocode function search(query, k): lists = parallel for shard in shards: shard.topK(query, k) heap = empty min-heap ordered by (score, docId) for list in lists: for hit in list: heap.push(hit) if heap.size > k: heap.popMin() return heap.drainDescending() ``` ## Why the local top-k is enough The merge is **exact**, not an approximation, as long as scores from different shards are comparable. Take any document `d` in the true global top 10, and suppose `d` were missing from its own shard's top 10. Then 10 documents on that shard would outrank it. Those same 10 documents also outrank it globally, so `d` could not be in the global top 10, which contradicts the starting point. Every global winner is therefore in some shard's list. How engines make scores comparable across shards is a ranking question and is outside this mechanism. ## Query-then-fetch: two phases instead of one Only `k` of the `S·k` candidates survive the merge, so having every shard send full documents would waste most of the bytes. Many engines therefore serve a query in two phases: 1. **Query phase.** Shards return only IDs, scores and sort values, and the coordinator merges them. 2. **Fetch phase.** The coordinator asks only the shards that own the final `k` winners for stored fields, snippets or highlights. The cost is one extra round trip. The gain is that the heavy payload is fetched for `k` documents instead of `S·k`. When `k` is tiny and documents are small, a single phase can be cheaper. ## What grows with the number of shards | Cost | How it scales | |---|---| | Sub-requests per query | one per shard (`S`) | | Candidates merged | up to `S·k` | | End-to-end latency | set by the slowest shard it waits for, not the average | | Coordinator memory | proportional to `S·k` | More shards keep each shard small and fast, but every query pays for all of them, and the chance that at least one shard is slow rises with `S`. Some systems place documents by a routing key, such as a tenant ID. When a query names that key, the coordinator can skip every shard that cannot hold a match, and the fan-out drops to one shard.

  • In scatter-gather search, why do shards usually return only IDs and scores in the first round?
    Only `k` of the `S·k` candidates survive the merge, so sending full documents from every shard wastes network and coordinator memory. In query-then-fetch serving, the first phase returns only IDs, scores and sort values. A second phase then fetches stored fields, snippets or highlights from the shards that own the final `k` winners. The price is one extra round trip, which is worth paying when documents are large or `k` times the shard count is big.
  • In a sharded search index, when can the coordinator skip shards instead of querying all of them?
    It can skip shards when documents are placed by a routing key and the query names that key. For example, if every document is routed by tenant ID and each query is scoped to one tenant, only that tenant's shard can hold matches, so the coordinator sends one sub-request. Without such a key, document-partitioned shards give no hint about where matches live, so every shard must be asked.

It is like asking every branch of a library chain for its ten best books on a subject and then picking the ten best from those short lists. No branch has to send its whole catalogue.

saying these in an interview costs you the question

  • Each shard returns all of its matches so the coordinator can sort them.
  • Merging each shard's top 10 can miss a document from the true global top 10.
  • The coordinator queries shards one after another until it has 10 results.
  • More shards reduce the number of sub-requests each query sends.
  • Shards must return full documents in the first round.
open as a page

In a search cluster that fans each query out to 100 shards, how do hedged requests tame tail latency caused by rare per-shard slowness?

level: seniorimportance: must knowfreq 55%

basics

~20 s

If each shard is slow 1% of the time, a query that waits on 100 shards is slow about 63% of the time. A hedged request sends a backup copy to another replica once the first is late and uses whichever reply arrives first.

open as a page

In a search engine sharded across 50 index shards, why does jumping to results page 1,000 cost far more than page 1?

level: middleimportance: should knowfreq 45%

basics

~20 s

Any shard could supply results 9,991-10,000, so each shard must return its own top 10,000, and the coordinator merges 500,000 entries to keep 10. Cursor paging avoids this by resuming after the last hit's sort key.

open as a page

For an interactive search box served by scatter-gather over many shards, how would you decide whether to return partial results when some shards miss the deadline?

level: principalimportance: should knowfreq 35%

basics

~20 s

Return partial results only if the missing shards hold little of the likely answer and the consumer can tolerate gaps. Flag the response as partial, keep it out of the result cache, and fail instead where completeness is the requirement.

open as a page

When partitioning a search index across machines, how does a document-partitioned layout differ from a term-partitioned one?

level: middleimportance: nice to knowfreq 30%

basics

~20 s

A document-partitioned index gives each shard a complete index over some of the documents, so every query hits every shard. A term-partitioned index gives each shard whole postings lists for some of the terms, so a query hits only its terms' shards.

open as a page