System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Search Autocomplete: the data pipeline, serving and follow-ups


The cached trie from Search Autocomplete: scope, scale and the top-k trie is fast to read and painful to modify. The resolution is to make the serving index immutable and rebuild it on a schedule.

This lesson builds that pipeline, then serves the result at 25,000 requests per second, then takes the five extensions interviewers ask for — and sorts the cheap ones from the ones that are a different system.

The data-gathering pipeline

USERTypingQuery servicereads onlyTrie cachein memory, shardedTrie builderhourly jobAggregatorcounts per termDATAQuery logsTrie snapshotsprefixtop 10log the querybatchcountspublishloadthe read path never writes: a suggestion isserved from a snapshot, never from live counts
Rebuilding the trie on a schedule, not per keystroke, is what keeps the read path at a few milliseconds.

Notice that every arrow points away from the serving nodes until the last one: nothing on the read path ever writes to the index.

1. Collection. Every executed search — the query the user actually ran, not every prefix they typed — is written to a distributed log (Kafka being the common implementation). Volume is 100 million events a day, around 1,000 per second, which is small.

2. Aggregation. A batch or stream job counts occurrences per query over a decay window: typically a weighted blend of the last day, last week, and last month, so a term that spiked yesterday ranks above one that was popular last year. Output is a table of query → weighted count, filtered to remove terms below a minimum count and anything on the unsafe list.

3. Build. A builder reads the aggregated table, constructs the trie, computes the top-k list at every node in a single bottom-up pass (each node's list is the merge of its children's lists plus its own terminal entry — this is the reason building is cheap and updating is not), and serialises the result into a compact, read-only file.

4. Distribution. The snapshot goes to blob storage. Serving nodes download it, build the in-memory structure, and swap the new index in behind an atomic pointer change. The old index is freed once in-flight requests drain.

Why the serving path never writes

Three reasons, and all three matter:

  • No locking. A read-only structure needs no synchronisation, so every core can traverse it concurrently at full speed.
  • Predictable latency. No lock contention, no rebalancing, no garbage from mutation. The 99th percentile stays close to the median, which is what a 50 ms budget requires.
  • Trivial rollback. A bad index is reverted by pointing back at the previous snapshot. A mutable index that has been corrupted by a bad update has to be rebuilt anyway.

The cost is staleness, and you should quantify it rather than wave at it.

How fresh the suggestions actually are

With an hourly rebuild, the worst case for a brand new query is: up to 60 minutes waiting for the next aggregation window, plus build time (minutes for tens of millions of terms), plus distribution and load time on every serving node (a few minutes for a multi-gigabyte snapshot).

Realistic end-to-end freshness: 60 to 90 minutes.

For most searches that is invisible — the popular completions of weath do not change hourly. For breaking news it is a visible failure, and the honest answer is that a second path is needed: a small, separate, real-time layer holding trending terms from the last few minutes, merged into the results at request time. The trending follow-up below returns to this. Keeping it separate from the main trie preserves the immutability that makes the main path fast.

Serving at scale

25,000 requests per second at peak against an in-memory index. The work here is distribution and cache placement, not algorithms.

Four places a prefix is answeredBrowser prefix cacheCDN or edge cacheTrie replica in memoryHourly rebuild job
The whole index is a few gigabytes, so replicate it everywhere rather than sharding and creating hot letters.

Does it even need sharding?

Start by checking. Requirements and scale put the raw term data at around 600 MB and the trie design put the trie with cached lists at a few gigabytes. That fits on one machine.

So the first answer is: replicate, do not shard. Run N identical serving nodes, each with a full copy of the index, behind a load balancer. Every node can answer every request, capacity scales linearly with node count, and losing a node costs capacity but not coverage. For an index of a few gigabytes this is the right design, and saying so is better than reflexively sharding a dataset that does not need it.

Sharding becomes necessary when the index outgrows a machine — many languages, personalised segments, or a much larger k.

Sharding by prefix, and the imbalance it creates

The obvious split is by first letter: node 1 serves a–d, node 2 serves e–k, and so on.

It distributes badly, because letter frequency is not uniform. In English query logs, initial letters such as s, c, p, and m carry far more traffic than x, z, and q — the skew between the busiest and quietest initial letter is large enough that an even alphabetical split leaves some nodes idle while others saturate. (The precise distribution depends on the language and the product, and should be measured rather than assumed.)

Two fixes:

  1. Build the shard map from observed traffic. Assign prefix ranges so that each shard gets a roughly equal share of requests, not an equal share of the alphabet. Since the index is rebuilt on a schedule anyway, the shard map can be recomputed at the same time from the previous window's traffic.
  2. Split hot prefixes deeper. A shard responsible for s alone may need splitting on the second character.

A router in front holds the current map. Because the map changes only at rebuild time, it can be distributed with the snapshot.

Caching in front

The distribution of prefixes is extremely skewed — a small number of short prefixes account for a large share of requests — which makes caching unusually effective here.

  • Browser cache. Set a cache lifetime of a few minutes on the response. A user typing res, deleting to re, and typing res again should not produce a second request.
  • Content delivery network / edge cache. Suggestions for a given prefix are identical for every non-personalised user, so they cache at the edge. This also removes the network round trip that the latency budget identified as the largest single component of the latency budget.
  • Server-side cache. A small in-process cache of the hottest prefixes, though with a read-only in-memory trie the lookup is already fast enough that this adds little.

Cache lifetime is a freshness trade-off: a 10-minute edge cache on top of a 60-minute rebuild adds 10 minutes to the worst-case staleness computed above.

The request pattern

Requests are issued asynchronously from the browser and must never block typing. Three client rules:

  • Debounce by about 50 ms, as costed in Requirements and scale.
  • Cancel in-flight requests when a newer keystroke arrives, so a slow response for re cannot overwrite the displayed suggestions for resta. Out-of-order responses overwriting newer ones is the most common client bug in this feature.
  • Fail silently. If the request errors or times out at, say, 200 ms, show nothing. The search box must still work.

Follow-ups

Five extensions, each of which changes something structural. The value here is knowing which ones are cheap and which ones are a different system.

Which extensions are cheapAfter the trieTypo toleranceTrending queriesPersonalisationNon-Latin scriptsUnsafe-word filtering
Personalisation is the expensive one: it turns a shared, cacheable index into a per-user computation.

Typo tolerance

Prefix matching on a trie cannot find restaurant from resturant, because the paths diverge at character five.

Options, cheapest first:

  1. A precomputed misspelling map. Mine the logs for pairs where a user typed one query, got few results, and immediately typed a similar one. Store resturant → restaurant and consult it before the trie. Cheap, high precision, and it covers the misspellings people actually make. Recommend this first.
  2. Edit-distance search over the trie. Traverse with a budget of one or two edits, which multiplies the nodes explored substantially — feasible for short prefixes, and a latency risk for the most common short queries.
  3. An n-gram index alongside the trie. Index every 3-character substring, retrieve candidates that share enough n-grams, then rank. This is a different index and a different service; it handles fuzziness properly and it costs what a second system costs.

Trending and time-sensitive queries

The 60-to-90-minute staleness from the pipeline above is unacceptable for breaking news. Run a small real-time layer: a streaming job over the last few minutes of query logs, keeping counts for terms whose rate has risen sharply, held in an in-memory store of at most a few thousand entries. At request time, merge that short list with the trie's answer, boosting trending terms.

Keeping it separate is deliberate. The main trie stays immutable and fast; the volatile data lives in a structure that is small enough to rebuild every minute.

Personalisation

Blending the user's own history into results breaks the property that made edge caching work: responses stop being identical across users.

The usual compromise is a two-tier blend: the global trie result is fetched and cached normally, and a small per-user list — recent queries, stored client-side — is merged into the top positions on the device. Personalisation costs nothing on the server and the edge cache survives. Full server-side personalisation means per-user state on every request and no shared caching, which is a large price for a modest gain.

Multi-language and non-Latin scripts

Three distinct problems:

  • Separate indexes per language. A user's language is known, so route to the right index. This multiplies index size by the number of languages, which is the point at which the sharding described above becomes necessary.
  • Scripts where a character is not a keystroke. In Chinese and Japanese, users type romanised input that an input method editor converts. The prefix the server receives may be Latin characters that must map to a non-Latin query, so the index needs both forms.
  • Languages without spaces or with rich morphology need a normalisation step before indexing, and the details are language-specific enough that a reviewer with domain knowledge should check any claim made about a specific language.

Filtering unsafe suggestions

Apply the blocklist at build time, not at serve time. Filtering during the build means the term never enters the trie and costs nothing at query time; filtering at serve time adds a check to the hot path and, worse, can return fewer than k results after filtering.

Keep a serve-time check as a second layer for terms blocked since the last build, and accept that it is a stopgap until the next rebuild. Also filter the combination: a safe prefix can have an unsafe completion, and that pairing is the failure mode that becomes a news story.