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
- Watch
- Try it
- Predict
- Capture
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.
Highlighted lines are the ones running in the diagram right now.
def search(query, from, size):# phase 1 — scatter (doc_id, score) onlyresults = []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 onlybodies = [shard_of(h).fetch(h.doc_id) for h in top]return assemble(top, bodies)
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
def merge(per_shard, from, size):# buffer = num_shards * (from + size)# this is the from + size <= 10 000 wallpool = []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