Skip to content

Design a Distributed Search Engine (Elasticsearch-style)

Case Study: Design a Distributed Search Engine

Section titled “Case Study: Design a Distributed Search Engine”

Not “search over a users list” (that’s a DB index lookup) — full-text search over unstructured documents, ranked by relevance. Think Elasticsearch/Solr, or the search box behind an e-commerce catalog or a docs site.


Functional:

  • Index a document (JSON with text fields); it becomes searchable within seconds
  • GET /search?q=... returns matching documents ranked by relevance, not just presence
  • Support phrase queries, filters (e.g. category:electronics), and pagination
  • Delete/update a document; it’s removed/refreshed from search results

Non-functional:

  • Index 10K+ documents/sec sustained
  • Search query p99 < 100ms
  • Index survives node failure with no data loss
  • Scale to billions of documents — no single node holds the full index

MetricValue
Documents1B
Avg document size2 KB → 2 TB raw text
Inverted index size (with term stats)~30-50% of raw text → ~700GB-1TB
Writes/sec10,000 docs/sec
Search QPS50,000

A forward index maps doc_id → text. Search needs the opposite: term → list of doc_ids containing it. That’s the inverted index — the core data structure of every full-text search engine.

Document 1: "the quick brown fox"
Document 2: "the lazy dog sleeps"
Document 3: "quick fox jumps"
Inverted index:
"the" → [1, 2]
"quick" → [1, 3]
"brown" → [1]
"fox" → [1, 3]
"lazy" → [2]
"dog" → [2]
"sleeps"→ [2]
"jumps" → [3]

Each posting list also stores term frequency (how often the term appears in that doc) and position (for phrase queries), not just the doc_id — relevance ranking and “exact phrase” matching both need this.

// Simplified posting list entry
{
term: "fox",
postings: [
{ docId: 1, termFreq: 1, positions: [3] },
{ docId: 3, termFreq: 1, positions: [1] },
]
}

flowchart LR
Client["📱 Client"] --> API["Search API"]
API --> Coord["Query Coordinator"]
Coord --> S1[("Shard 1<br/>Inverted Index")]
Coord --> S2[("Shard 2<br/>Inverted Index")]
Coord --> S3[("Shard N<br/>Inverted Index")]
Indexer["Indexer Service"] --> Queue["Write Queue"]
Queue --> S1
Queue --> S2
Queue --> S3
style Client fill:#7c3aed,color:#fff
style API fill:#4f46e5,color:#fff
style Coord fill:#6366f1,color:#fff
style S1 fill:#8b5cf6,color:#fff
style S2 fill:#8b5cf6,color:#fff
style S3 fill:#8b5cf6,color:#fff
style Indexer fill:#059669,color:#fff

A single node can’t hold a billion-document inverted index in memory. Documents are partitioned across shards (typically by hash(doc_id) % num_shards), each shard holding a complete inverted index over its own subset of documents — not a partition of the vocabulary.

sequenceDiagram
participant C as Client
participant Q as Query Coordinator
participant S1 as Shard 1
participant S2 as Shard 2
participant S3 as Shard 3
C->>Q: search("quick fox")
par scatter to all shards
Q->>S1: search("quick fox")
Q->>S2: search("quick fox")
Q->>S3: search("quick fox")
end
S1-->>Q: top 10 local matches + scores
S2-->>Q: top 10 local matches + scores
S3-->>Q: top 10 local matches + scores
Q->>Q: merge, re-rank globally, take top 10
Q-->>C: final ranked results

This is the scatter-gather pattern: every query fans out to every shard (each shard independently has documents matching any term), and the coordinator merges partial results. Because each shard scores its own local matches independently, per-shard term statistics (like document frequency) are approximations of the global statistic unless the coordinator does a global stats round-trip first — a real accuracy/latency trade-off, not a free scaling win.


“Contains the word” isn’t enough — results must be ranked by relevance. BM25 (used by Elasticsearch/Lucene) scores a document for a query based on term frequency, inverse document frequency, and document length normalization:

score(D, Q) = Σ IDF(term) × [ tf(term, D) × (k1 + 1) ]
─────────────────────────────────
[ tf(term, D) + k1 × (1 - b + b × |D|/avgDL) ]
IDF(term) = log( (N - n(term) + 0.5) / (n(term) + 0.5) + 1 )
N = total documents
n(term) = documents containing the term
tf = term frequency in this document
|D| = this document's length
avgDL = average document length across the index
k1, b = tuning constants (typical: k1=1.2, b=0.75)
function bm25Score(term, doc, index, k1 = 1.2, b = 0.75) {
const N = index.totalDocs;
const n = index.docFreq(term); // docs containing this term
const idf = Math.log((N - n + 0.5) / (n + 0.5) + 1);
const tf = doc.termFreq(term);
const norm = 1 - b + b * (doc.length / index.avgDocLength);
return idf * (tf * (k1 + 1)) / (tf + k1 * norm);
}

Intuition: a term that’s rare across the whole index (high IDF — e.g. “quixotic”) is a much stronger relevance signal than a term that’s everywhere (low IDF — e.g. “the”). Length normalization (b) stops long documents from winning purely by containing every term somewhere.


BottleneckSolution
Scatter-gather means every query touches every shardBound shard count sensibly; a query’s latency is the slowest shard’s latency — tail latency matters more than average
Real-time indexing vs. query throughput contentionBuffer writes into in-memory segments, periodically merge into larger on-disk segments (Lucene-style segment merging) instead of updating one shared structure per write
Deletes are expensive in an append-only index structureSoft-delete (tombstone bit) at write time, physically purge during background segment merges
Global term statistics (IDF) differ from per-shard estimatesAccept per-shard approximation for low-latency queries, or pay for a global-stats round-trip when ranking accuracy matters more than speed
Hot shard from a skewed document distributionRebalance shards by size/load, not just document count, and consider routing by a smarter key than plain hash if access patterns are skewed

Q: Why can’t a single node just hold the whole inverted index if we throw enough RAM at it? Even ignoring memory limits, single-node search throughput is capped by that node’s CPU for scoring and I/O for reading posting lists — sharding isn’t just about fitting data, it’s about parallelizing the scoring work itself across many nodes for a single query.

Q: If per-shard IDF is only an approximation of the true global IDF, doesn’t that make ranking wrong? It makes ranking approximately right, which in practice is good enough for most use cases as long as documents are randomly distributed across shards (not clustered by content) — a document isn’t systematically favored or penalized because its shard happens to have slightly different term statistics than the global average.

Q: How do you handle a document update (not just a fresh insert) in an append-heavy index structure? Treat an update as a delete-then-reinsert: tombstone the old version so it’s excluded from scoring/results, write the new version as a fresh document in the current in-memory segment, and let the background merge process physically reclaim the tombstoned space later.

Q: What happens to a search query if one shard is slow or temporarily unavailable? The coordinator can apply a timeout per shard and return partial results (marking the response as partial) rather than blocking the whole query on the slowest or dead shard — replicating each shard (so a replica can serve if the primary is down) is the actual fix for availability, the timeout is just damage control.

Q: How would phrase search (“quick fox” as an exact phrase, not just both words present) work against the posting list structure shown? Because each posting stores term positions within the document, a phrase query intersects the posting lists for “quick” and “fox” and additionally checks that “fox“‘s position is exactly one greater than “quick“‘s position in the same document — position data is what turns “both words present” into “these words are adjacent.”


  • Full-text search is built on an inverted index — term → documents containing it, not the forward doc → text mapping a normal DB uses.
  • Billions of documents means sharding by document, with every query scattered to every shard and results merged — not partitioning the vocabulary.
  • BM25 ranks by how rare and how frequent a term is in a document, normalized by document length — “contains the word” isn’t ranking.
  • Writes go into small in-memory segments merged into larger ones in the background, so indexing throughput and query latency don’t fight over the same structure.