How Elasticsearch works
Elasticsearch and OpenSearch internals for system design: inverted indexes and analyzers, Lucene segments, refresh and near real-time search, shards and replicas, scatter-gather queries, relevance scoring with BM25, indexing pipelines and common scaling mistakes.
Reading is half of it. See this used in a real interview: walk through Design Yelp →
Whenever a design needs full-text search, filtering with facets, or log search, the box usually says Elasticsearch (or its fork, OpenSearch). Knowing how it works inside explains its behaviour: why new documents appear after a second, why too many shards hurt, why deep pagination is expensive, and why it should not be your source of truth.
The inverted index
Search engines index terms, not rows. For each term, an inverted index stores the list of documents containing it (a postings list), with positions and frequencies. A query for "cheap pizza" looks up both terms and combines their postings lists, which is fast no matter how many documents there are. See search and indexing.
Analyzers turn text into terms at index and query time: tokenising, lowercasing, removing stop words, stemming ("running" to "run"), handling synonyms and languages. The same analysis must apply at both ends, or queries miss documents.
Structured fields (numbers, keywords, dates, geo points) use other structures (BKD trees, doc values in columnar form) for range filters, sorting and aggregations.
Lucene segments
Each shard is a Lucene index made of immutable segments:
- New documents are buffered in memory and written to the translog (a write-ahead log) for durability.
- A refresh (every second by default) writes the buffer as a new small segment and makes it searchable. This is why Elasticsearch is near real-time: a document is searchable about a second after indexing.
- Segments are never modified. Updates are a delete (a marker) plus a new document; deletes are bitmaps.
- Background merges combine small segments into larger ones and purge deleted documents, much like LSM compaction. See storage engines.
Heavy update workloads therefore cost more than they seem: every update rewrites the whole document and creates merge work.
Shards and replicas
- An index is split into primary shards (fixed at creation, unless you reindex or split), distributed across nodes.
- Each primary has replica shards on other nodes for availability and read throughput.
- Documents are routed to shards by hashing the id (or a custom routing key, such as tenant id, to keep related documents together).
How a search runs
- A coordinating node receives the query and sends it to one copy of every shard (scatter).
- Each shard finds its top N matches and returns ids and scores.
- The coordinator merges them into the global top N (gather), then fetches those documents.
Consequences:
- Latency follows the slowest shard. See tail latency.
- Deep pagination (page 1,000) makes every shard return thousands of results to merge. Use
search_aftercursors instead of large offsets. See pagination. - Custom routing lets a query hit one shard instead of all.
Relevance
Scoring defaults to BM25: documents score higher when the query terms are frequent in them, rare across the collection, and in shorter fields. Real products add business signals (ratings, distance, popularity, freshness) with function scores, and often a second-stage machine-learning re-ranker over the top results. See Design Yelp.
Keeping it in sync
Elasticsearch should be a derived store, rebuilt from a source of truth:
- Stream changes from the primary database via change data capture into an indexer that writes in bulk.
- Make indexing idempotent with document ids and versions so replays are safe.
- Reindex into a new index and switch an alias for mapping changes, with zero downtime.
Scaling and common mistakes
- Too many shards: each shard has overhead; thousands of tiny shards waste memory and slow the cluster. Aim for shards in the tens of gigabytes.
- Too few shards for a growing index: a single shard cannot spread across nodes.
- Time-based indices for logs and metrics (one per day), with lifecycle policies that move old indices to cheaper nodes and delete them. See time-series data.
- Mapping explosions: dynamic mapping with arbitrary keys creates thousands of fields. Define mappings explicitly.
- Expensive queries: leading wildcards, huge aggregations and scripts; test with realistic data.
- Bulk indexing with refresh disabled or slowed during large loads.
Autocomplete
Elasticsearch supports prefix search with edge n-gram analyzers and completion suggesters, though very high-QPS typeahead is often served from a precomputed prefix table instead. See tries and autocomplete.
Checklist
- Inverted index with consistent analyzers at index and query time.
- Near real-time via refresh; merges reclaim deleted documents.
- Primary shard count planned; replicas for availability and reads.
- Scatter-gather awareness: routing,
search_after, slow-shard latency. - BM25 plus business signals and re-ranking.
- Derived from the source of truth via CDC; aliases for reindexing.
- Right-sized shards, explicit mappings, lifecycle for time-based indices.