Build your own distributed search engine (Elasticsearch / OpenSearch style)
Five million books, a search box, and a 100 ms budget. Build the engine from the inverted index up — segment, refresh, shard, replica, scatter-gather, BM25 — and feel why every guarantee that lives across shards is paid for in either an extra round trip or a small lie about the rankings.
- Scenes
- 12 interactive scenes
- Time
- about 84 minutes
- Topic
- Search, Indexing & Retrieval
What you are building, and why
You are designing the simplest thing that could possibly work as a search engine: a client sends a JSON query, the engine returns documents that match — ranked by relevance — in under 100 ms across a corpus that does not fit on one machine. Five million books, a search box, and a query that says "distributed systems."
The naive plan is SELECT * FROM books WHERE body LIKE '%distributed systems%'. It is wrong before you even count the rows. LIKE '%…%' is unanchored, so no B-tree applies; throwing 64 cores at it cuts wall-clock by 64× and still does not approach 100 ms; and there is no concept of relevance in the answer. The architectural response is not faster hardware. It is to flip the relation between document and term and serve queries from a precomputed map.
That map is an inverted index: each word in your corpus points to the list of document IDs that contain it. A search becomes a hash lookup followed by a sorted-list intersection — the cost collapses from O(documents) to O(matches). Every later mechanism — segments, refresh/flush, shards, replicas, scatter-gather, BM25 — exists to make that one structure correct under mutability, durability, and scale.
Resist the urge to "describe Elasticsearch." Make decisions yourself, defend them, and let the design push back. The point is to feel why each price is paid: why segments are immutable, why num_primary_shards is fixed for life, why replicas help reads and not writes, why deep pagination has a hard wall at 10 000, why BM25 is approximate across shards.
What you will be able to explain afterwards
- inverted index — term → posting list
- Lucene segments — many small immutable indexes + background merges
- refresh / flush / translog — three cadences for one write
- primary shards — hash(routing) % N, fixed for life
- replicas — read scale + HA, no write speedup
- scatter-gather — query then fetch, the from+size=10 000 wall
- BM25 — IDF × saturating-TF × length-norm
- distributed scoring drift — per-shard IDF, dfs_query_then_fetch
- aggregation approximation — shard_size and the ghost bucket
- operational sharp edges — refresh storm, mapping explosion, hot shard, ILM
- 01Five million books, under 100 ms — SQL LIKE scans and the document-term flip5M books, a search box, and a 100ms budget — the naive scan-every-row plan never crosses the line, no matter the cores or the disk.~7 min
- 02The inverted index — term to docsFlip the relation: a precomputed map from each word to the doc IDs it appears in turns O(documents) Boolean queries into O(matches) sorted-list merges.~7 min
- 03Segments — many small immutable indexesLucene appends a fresh tiny inverted index per write batch and merges in the background — readers never lock, deletes are bit flips, mutability is an asynchronous receipt.~7 min
- 04Refresh, flush, translog — three cadencesRefresh = visible to search; flush = survives a crash; translog bridges the gap. Three cadences, three durability properties, one famous source of confusion.~7 min
- 05Shards — hash(routing) mod NEach document lives on shard = hash(routing) mod number_of_primary_shards — and that mod is exactly why you cannot change the shard count after you create the index.~7 min
- 06Replicas help reads, not writesReplicas linearly add read capacity and HA but do not speed up writes — every primary fans every write out to all in-sync replicas before acking.~7 min
- 07Scatter, reduce, fetch — the two-phase query-then-fetch search_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.~7 min
- 08BM25 — TF saturation × IDF × length normBM25 ranks each match by IDF × saturating-TF / length-norm — and the saturation knob k1 is exactly what stops a keyword-stuffed doc from winning the page.~7 min
- 09Distributed scoring is approximateEvery shard computes IDF from its own local corpus, so the same document scores differently depending on where it lives — dfs_query_then_fetch buys correctness for one extra round trip.~7 min
- 10Top-N aggregations can miss the winnerEach shard returns its own top shard_size buckets and the coordinator merges — a term that is 11th on every shard is invisible in a top-10 result, and shard_size is the knob.~7 min
- 11Operational sharp edges — refresh storms, mapping explosions and hot shardsRefresh storm, mapping explosion, hot shard, deep pagination — the four most expensive Elasticsearch outages are misuses of knobs the earlier scenes already introduced.~7 min
- 12Design your search clusterCapstone: pick e-commerce, logs, or security analytics and configure shards, replicas, refresh, scoring, and ILM — the verifier traces every ✓/✗ back to the scene that earned it.~7 min
Where you'll use this
Product designs whose trade-offs turn on what this curriculum teaches.
More in Search, Indexing & Retrieval
Finding a needle: inverted indexes, distributed search, vector and ANN retrieval, crawling, and the query box itself.
- Build Build a vector database (Pinecone / Weaviate / pgvector style)Approximate nearest-neighbor over a billion 1536-dim vectors in 10 ms. Build the index from scratch (HNSW, IVF, PQ), pay the recall-vs-latency tax explicitly, support filtered + hybrid search, and feel why every LLM stack in 2026 has a vector store next to its KV store.
- Web CrawlerPolitely traverse the web at scale. Don't crawl yourself in circles.
Prefer to design it yourself?
The same subject as a staged workspace: draw the architecture, and a simulator traces requests through the boxes you drew.
Open the Build a distributed search engine (Elasticsearch / OpenSearch style) workspace