Mid to Senior Engineer

System Design Interview Prep

A structured path from the interview framework through core concepts, key technologies and patterns to eighteen full problem breakdowns, each with diagrams and weak, solid and excellent answers to every deep dive.

Chapter 11 of 36Key technologies · Elasticsearch and Search Engines

Key Technology: Elasticsearch and Search Engines

When a design needs full-text search, relevance ranking, filtering with facets, or flexible queries over a lot of documents, the tool is a search engine, and Elasticsearch (and the engines related to it, built on the same underlying library) is the one interviewers mention. This chapter explains how it works, how to keep it in sync with your database, and how it fails, so that you can add a search box to a design without hand-waving.

Features and defaults vary between versions and products, so confirm specifics against the documentation of the version you would use.

1. What a search engine is for

A relational database is excellent at finding rows by key and by exact conditions. It is poor at questions like "find products whose description mentions waterproof hiking boots, ranked by relevance, filtered by size and price, with counts per brand, tolerant of typos". A search engine is built for exactly this.

Use it for:

  • Full-text search with relevance ranking.
  • Filtering and faceting across many attributes at once.
  • Autocomplete and fuzzy matching.
  • Log and event search and aggregation, which is another large use of the same technology.
  • Geospatial search combined with text and filters.

Do not use it as your primary database. It is a derived read model: the source of truth lives elsewhere, and the search index is built from it.

2. The inverted index

The core data structure is the inverted index. A normal index maps a document to its words. An inverted index maps each term to the list of documents that contain it, with positions and frequencies.

Suppose three documents:

doc 1: "waterproof hiking boots"
doc 2: "leather hiking shoes"
doc 3: "waterproof rain jacket"

The inverted index looks like this:

TermDocuments
boots1
hiking1, 2
waterproof1, 3
leather2
shoes2
rain3
jacket3

A query for "waterproof hiking" looks up two terms and intersects or scores their document lists. The cost depends on the lists touched, not on the total number of documents, which is why search stays fast at huge scale.

Analysis. Before indexing, text goes through an analyzer: split into tokens, lowercase, remove stop words, stem words to a root ("running" to "run") and handle synonyms and accents. The same analysis runs on the query so that terms match. Choosing analyzers per language and field is where much of the quality comes from.

3. Relevance ranking

Documents that match are scored. The classic scoring idea is TF-IDF: a term matters more if it appears often in the document (term frequency) and rarely across all documents (inverse document frequency). The modern default in these engines is a refinement called BM25, which also dampens the effect of very long documents and repeated terms.

Practical levers you can mention:

  • Field boosts: a match in the title counts more than one in the body.
  • Filters versus queries: filters (exact conditions such as category or price range) do not affect scoring and can be cached, so use them for constraints.
  • Business signals: blend in popularity, recency or rating.
  • Learning to rank: a trained model re-scores the top candidates for further relevance gains, as a later stage.

Say that relevance is a product and measurement problem, not a switch, and that you would evaluate changes with click-through and judged queries.

4. Distribution: shards and replicas

An index is split into shards, and each shard is a self-contained index. Documents are assigned to a shard by hashing an identifier. Each primary shard has replica copies on other nodes, for fault tolerance and extra read capacity.

<!--fig:cluster-->
merge Search request Coordinator fan out, merge Node 1 P0 (primary) R1 (replica) Node 2 P1 (primary) R2 (replica) Node 3 P2 (primary) R0 (replica) Top results scored and merged Primary shard count is fixed when the index is created. Figure 1. An index is split into primary shards, each copied to replicas on other nodes. A search fans out to one copy of every shard and merges the results.

A search request goes to a coordinating node, which sends the query to one copy (primary or replica) of every shard, collects each shard's top results and merges them into the final ranking. A document write goes to its primary shard and is copied to the replicas.

Important consequences:

  • The number of primary shards is fixed at index creation. You cannot simply add shards later. If you outgrow the shard count, you create a new index with more shards and reindex. Plan for this, and use time-based indexes for logs so that old indexes can be dropped.
  • Too many small shards waste resources, and too few huge shards limit scale and recovery speed. A common guideline keeps shards in the tens of gigabytes. Treat it as a starting point to measure.
  • A query touches every shard, so cost grows with shard count. Routing by a key (for example by tenant) can limit a query to one shard.
  • Replicas increase read throughput and survive node loss, and they cost storage and indexing work.

Documents are not searchable the instant they are written. New documents are written to an in-memory buffer and a transaction log, and become searchable after a periodic refresh, by default about a second. That is why this is called near real time. Changes are made durable through the transaction log, and the segments are merged in the background.

Implications to state: a freshly written document may not appear in search for about a second, so do not use the search engine to read back what a user just wrote. Read your own writes from the database, and let the search index catch up.

6. Keeping the index in sync with the database

Because the index is a read model, you must feed it. The reliable pattern:

<!--fig:ingest-->
index query Primary database source of truth Change capture log or outbox Queue Indexer bulk requests Index refresh ~1 s Search API A new or changed record becomes searchable after seconds.If the index is lost, rebuild it from the database. Figure 2. Search is a read model: the database stays the source of truth and changes flow into the index asynchronously.
  1. Write to the primary database only. That is the source of truth.
  2. Capture the change, either by reading the database's change log or by writing an event to an outbox table in the same transaction as the data change.
  3. Publish the change to a queue.
  4. An indexer consumes changes and sends them to the search engine in bulk requests, which are far more efficient than one request per document.

Why not write to both from the application? A dual write can fail halfway: the database commit succeeds and the index update fails, or the reverse, leaving the two inconsistent with no record of what went wrong. A change log or outbox makes the pipeline reliable and replayable.

Further points:

  • Idempotent indexing. Use the document's identifier as the index identifier, so reprocessing a change overwrites the same document and is harmless. Use a version number so that an older update cannot overwrite a newer one.
  • Rebuilding. Because the database is the truth, you can rebuild the index from scratch by reading all data, writing to a new index and switching an alias to it when ready. This is how you change mappings or analyzers without downtime.
  • Lag. Monitor the delay between a database change and its appearance in the index, and alert on it.

7. Querying at scale

Pagination problem. Asking for page 1,000 of results (an offset of tens of thousands) forces every shard to produce and sort that many results, which is slow and memory-heavy. Limit the maximum depth, and use search-after style pagination that continues from the last result's sort values, which stays cheap. For exporting everything, use a scrolling or snapshot mechanism, not deep offsets.

Aggregations. The engine computes counts, averages and facets over matching documents, which powers filters like "brand (12), colour (5)". These are fast on structured fields prepared for it, and expensive on high-cardinality fields.

Mappings. Define field types explicitly. Letting the engine guess types can lead to a mapping explosion, where unbounded dynamic field names create a huge number of fields, which exhausts memory. Disable dynamic mapping for fields you do not control.

Autocomplete. Use prefix-friendly analysis, such as n-grams of terms at index time, or dedicated suggesters. A separate precomputed structure is often better for the highest-traffic suggestion box, as in the autocomplete chapter.

8. Failure modes and operations

  • Split brain and cluster state. Modern versions use a consensus-based master election with a majority, so run an odd number of master-eligible nodes (three is typical).
  • Node loss. Replicas of its shards are promoted to primaries, and new replicas are rebuilt on other nodes, which uses network and disk. Leave capacity headroom.
  • Heap pressure and long garbage collections. Large aggregations, very wide documents and too many shards cause memory trouble. Monitor memory and limit expensive queries.
  • Slow queries. Wildcard and regular-expression queries, very deep pagination and huge aggregations can hurt the whole cluster. Set timeouts and limits, and use separate nodes or clusters for heavy analytical queries.
  • Not durable like a database. Do not treat it as the only copy. Take snapshots to external storage, and rely on rebuild from the source of truth.
  • Disk watermarks. When disks fill, the engine stops allocating shards and can mark indexes read-only. Alert early.
  • Reindexing and mapping changes need a new index and an alias switch.

9. Where it fits in designs, with the reasoning

E-commerce product search. Products in the relational database, a search index with analyzers, facets and boosts, fed by change capture, with the cart and checkout reading the database.

Search over user content such as messages or notes: index per tenant or routed by tenant, with permission filters applied in the query.

Log search. Time-based indexes, with retention by dropping old ones, and tiered storage for older data.

Autocomplete and "did you mean". Suggesters or precomputed prefix structures for speed.

When not to add it. For modest data and simple search, a relational database's built-in text search may be enough, and avoids running another system and the sync pipeline.

10. Interview questions and model answers

Q: Why not just use SQL LIKE queries? They scan the data, do not rank by relevance, do not handle stemming, typos or synonyms, and do not scale to rich filtering. An inverted index finds matching documents quickly and ranks them.

Q: How do you keep the search index in sync with the database? Write to the database only, capture changes through the change log or an outbox, publish to a queue, and have an indexer bulk-index them idempotently. I avoid dual writes, because they can leave the two stores inconsistent.

Q: Is Elasticsearch your source of truth? No. It is a derived read model that can lag by seconds and can be rebuilt from the database. I read my own writes from the database.

Q: How does a search scale? The index is split into shards across nodes, with replicas for availability and read capacity. A query fans out to every shard and the coordinating node merges the results. The number of primary shards is fixed at creation, so I size it up front or plan to reindex.

Q: How do you handle deep pagination? I cap the depth and use search-after style pagination that resumes from the last result, since large offsets make every shard sort too many results.

Q: How do you change the analyzer without downtime? Create a new index with the new settings, reindex from the source of truth, and switch an alias to the new index.

11. Common mistakes

  • Using it as the primary database.
  • Dual writes from the application to the database and the index.
  • Reading your own writes from the search index immediately.
  • Deep offset pagination.
  • Too many tiny shards, or too few huge ones, and ignoring that the shard count is fixed.
  • Letting dynamic mapping create unbounded fields.
  • Putting sensitive permission checks only in the client instead of in the query filters.
Header Logo