MentorNode
Start free
Search, Ranking & DiscoveryHarddesign-search-engine

Design a Distributed Search Engine

Design the query side of search: an inverted index sharded across hundreds of nodes, scatter-gather retrieval, and ranking that runs in tens of milliseconds.

Inverted IndexScatter-GatherIndex ShardingTwo-Phase RankingSegment Merging
Traffic & Capacity Estimates:

10B documents · 100k queries/second · p99 under 200ms · index refresh within 60s

Functional Requirements

  • •Build and maintain an inverted index from term to posting list over billions of documents.
  • •Execute boolean and phrase queries across index shards and merge the results.
  • •Rank candidates with a cheap first pass and an expensive model over the top slice.
  • •Make newly indexed documents searchable within a bounded refresh window.

Non-Functional Requirements

  • •Query p99 under 200ms including ranking and snippet generation.
  • •A shard replica loss must degrade recall slightly, not fail the query.
  • •Index updates must not pause query serving.

Back-of-the-Envelope Math

  • 10B docs across 500 shards = 20M docs/shard; each query fans out to all 500 and merges.
  • Scatter-gather p99 is the slowest shard's p99 — with 500 shards, a 1-in-1000 slow node is hit by half of all queries.

Key Architectural Trade-offs

  • Shard by document (every query hits every shard, simple updates) vs shard by term (few shards per query, brutal for multi-term queries and updates).
  • Two-phase ranking — cheap BM25 retrieval then an expensive learned reranker on the top 1,000 — is what makes a heavyweight model affordable per query.
  • Near-real-time indexing means small, frequent segments and constant merging; batch indexing is far more efficient and leaves the index minutes behind.

Click or drag a component onto the canvas, then connect the handles to draw the data flow.

3 nodes · 2 edges

Components · 35

Client & Edge4
Compute & Gateway7
Storage & Caching11
Messaging & Streaming6
Coordination & Ops5
Intelligence2
Canvas overview