Scatter, reduce, fetch — the two-phase query-then-fetch search

Every _search is two round trips — phase 1 scatters to one copy of every shard for top-(from+size) (doc_id, score) tuples and phase 2 fetches bodies for just the winners — which is why from + size is bounded at 10 000.

Scene 07

Scatter, reduce, fetch

  1. Watch
  2. Try it
  3. Predict
  4. Capture
COORDINATORcoord · "distributed systems"phase 0 · searchs0top-K candsbook 11.2book 50.7cand 10.1cand 20.1s1top-K candsbook 40.4cand 10.1cand 20.1s2top-K candscand 10.1cand 20.1s3top-K candsbook 23.1cand 10.1cand 20.1s4top-K candsbook 32.8cand 10.1cand 20.1GLOBAL TOP-N · merged at coordinatorscore(waiting for shard responses)FROM + SIZE · candidate bytesfrom=0 size=10≈ 5,000 bytes shipped to coordCoordinator idle — 5 shards, from=0, size=10
What to watch for

The receiving node becomes the conductor. Watch phase 1 fan out 'distributed systems' to one copy of every shard, the coordinator merge the local top-Ks into a global top-N, then phase 2 fetch _source ONLY from the shards that actually contributed a winner.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Coordinator.search
phase 1 scatter for top-(from+size); phase 2 fetch winners
def search(query, from, size):
# phase 1 — scatter (doc_id, score) only
results = []
for shard in all_primary_shards:
copy = adaptive_replica_selection(shard)
results.append(copy.localSearch(query, from + size))
top = merge(results, from, size)
# phase 2 — fetch _source for global winners only
bodies = [shard_of(h).fetch(h.doc_id) for h in top]
return assemble(top, bodies)
Shard.localSearch
BM25-rank local matches; ship (doc_id, score) only
def localSearch(query, top_k):
pq = PriorityQueue(top_k)
for doc_id in postings_for(query):
score = bm25(query, doc_id)
pq.offer((doc_id, score))
return pq.drain() # no _source yet
Coordinator.merge
buffer N × (from+size) candidates, then slice
def merge(per_shard, from, size):
# buffer = num_shards * (from + size)
# this is the from + size <= 10 000 wall
pool = []
for hits in per_shard:
pool.extend(hits)
pool.sort(key=score, reverse=True)
return pool[from : from + size]

Where this sits in Build a distributed search engine (Elasticsearch / OpenSearch style)

Scene 07 of 12. _search is two round trips: scatter top-(from+size) doc-ids+scores from every shard, reduce, then fetch only the global winners — which is where the 10,000 page wall comes from.

All 12 scenes in Build a distributed search engine (Elasticsearch / OpenSearch style) · Every curriculum

Built with Arqly
Every scene in Build a distributed search engine (Elasticsearch / OpenSearch style) builds on the one before it.All 12 Build a distributed search engine (Elasticsearch / OpenSearch style) scenes