AI System Design
← All chapters
Chapter 13 5 min read

Design a Search Autocomplete System

Search autocomplete suggests the most popular completions for whatever a user has typed so far, within about a hundred milliseconds of each keystroke. The design separates an offline data-gathering pipeline that builds a trie with cached top-k results from a fast query service that only reads it.

Architecture at a glance
  1. Client (debounced AJAX)
  2. Load balancer
  3. API servers
  4. Trie cache
  5. Trie DB
  6. Data gathering pipeline (logs → aggregators → workers)

Problem & requirements

As the user types, the system returns the top suggestions that start with the current prefix. Clarify whether matching is only at the beginning of the query (assume yes), how many suggestions to show (assume five), how ranking works (by historical query frequency), whether spelling correction is needed (no), the language (English first), and scale (10 million DAU).

The non-functional requirements carry the design. Responses must be fast, ideally under about 100 ms, or the UI stutters. Suggestions must be relevant to the prefix and sorted by popularity. The system must be scalable and highly available, since it sits in front of every search.

Back-of-the-envelope estimation

Assume each user performs 10 searches a day and a query is about 20 bytes (four words of five characters). Every character typed sends a request, so one search triggers about 20 requests. That gives 10M × 10 × 20 / 86,400 ≈ 2,000,000,000 / 86,400 ≈ 24,000 QPS, with peaks around 48,000 QPS.

If 20% of daily queries are new, new data is 10M × 10 × 20 bytes × 0.2 = 0.4 GB per day. Check: 100M queries × 20 B = 2 GB of raw queries; a fifth of that is 0.4 GB. Storage is small; the challenge is serving tens of thousands of lookups per second with tight latency.

High-level design: gathering vs querying

The system splits into two services. The data gathering service collects user queries and aggregates them into frequencies. The query service answers 'given this prefix, what are the top five queries?' A naive version stores a frequency table and runs a SQL query with LIKE 'prefix%' ordered by frequency, which works for tiny data sets but becomes a bottleneck at scale.

Separating the two lets each optimise for its job: gathering can be slow, batched, and eventually consistent, while querying must be fast and read-only. In practice the suggestions are rebuilt periodically rather than updated in real time, because instant updates on every search would overwhelm the serving structure and barely change the results.

Figure 1Offline gathering, online serving
Offline gathering, online servingDATA GATHERING (BATCH)QUERY SERVICEwrite trieloadprefixtop 5sampledAnalytics logsquery + time, append-onlyAggregatorscount per query per weekAggregated countsquery, week, frequencyWorkersrebuild trie on scheduleLoad balancerAPI serversTrie cachein-memory snapshotTrie DBprefix → node, top-kSearch boxAJAX per keystroke
The two halves meet only at the trie DB: the top row rebuilds it in batches, while the query path reads an in-memory copy and never waits on aggregation. Query logging is sampled and asynchronous.

The trie and why it is optimised

A trie (prefix tree) stores strings so that each node represents a prefix and children extend it by one character. Storing each query's frequency on its terminal node lets the system find completions by walking down to the prefix node and exploring its subtree.

The naive lookup costs O(p) to find the prefix, then O(c) to traverse all c children below it, then O(c log c) to sort them. That is too slow for short prefixes like 'a', which cover enormous subtrees. Two optimisations fix it. First, limit the prefix length, since almost nobody types very long queries, making the first step O(1). Second, cache the top-k queries at every node, so the answer is read directly from the prefix node in O(1). The price is extra memory, which is an acceptable trade for fast responses.

Figure 2Trie with cached top-k
Trie with cached top-kbeiesrttlookupUser types "be"read node, done: O(1)rootbeer 40 · best 35bbeer 40 · best 35bebeer 40 · best 35bibit 18beefreq 15 · beer 40, bee 15besbest 35bitfreq 18beerfreq 40bestfreq 35
Each prefix node stores its best completions (here k = 2), so answering be is one node read: beer, best, with no subtree walk or sort. The price is that a count change must update every ancestor's list.

Data gathering pipeline

Raw analytics logs record every query with a timestamp; they are append-only and not indexed. Aggregators roll these up into frequencies per query per time window. The window depends on the use case: a general search engine might aggregate weekly, while a real-time product like Twitter needs much shorter windows.

Workers are servers that run asynchronous jobs at regular intervals to build the trie from aggregated data and store it in the trie DB. The trie can be stored as a serialised snapshot in a document store, or mapped into a key-value store where each prefix is a key and its node data is the value. A trie cache keeps the trie in memory for reads and takes a fresh snapshot after each weekly rebuild.

Query service and trie operations

A request goes from the load balancer to an API server, which reads the trie cache and builds the response; on a cache miss it refills from the trie DB. Several tricks shave latency. The client uses AJAX so the page never reloads. The browser can cache results using a Cache-Control header, since suggestions for a prefix rarely change within an hour. Data sampling logs only one of every N requests, because logging every keystroke wastes processing and storage.

The trie supports three operations. Create is done by workers from aggregated data. Update can rebuild the whole trie weekly, replacing the old one, or update individual nodes, though that is slow because changing a leaf means updating cached top-k lists on every ancestor. Delete handles hateful, violent, or dangerous suggestions through a filter layer in front of the cache, so removal is immediate; the bad entries are then purged physically during the next rebuild.

Scaling storage with sharding

When the trie no longer fits on one server, shard it. A simple scheme splits by first letter: 'a' to 'm' on one server and 'n' to 'z' on another, extending to 26 shards, and to further levels by second letter if needed. The problem is skew: far more English words start with 'c' than with 'x', so letter-based shards are badly uneven.

The fix is to analyse the historical distribution and assign ranges by load, not by alphabet. A shard map manager keeps the lookup table from prefix range to shard, for example putting all of 's' on one shard and all of 'u' through 'z' on another, and the query service consults it to route each lookup.

Figure 3Load-based trie shards
Load-based trie shardsprefixwhichshard?top-k for "ca"other prefixesother prefixesShard map managerprefix range → shardAPI serverQuery "ca"Shard 1prefixes a–bShard 2c only (heavy)Shard 3u–z (light, combined)
Shards follow measured query volume, not the alphabet: a heavy letter like c gets a shard to itself while light letters u–z share one. The API server asks the shard map manager where a prefix lives, then reads only that shard.

Trade-offs & wrap-up

Caching top-k at each node trades memory for speed, and weekly rebuilds trade freshness for simplicity. For real-time trending queries such as breaking news, a weekly build is too slow; options include shrinking the window, giving recent queries more weight, and using stream processing systems such as Kafka with Spark Streaming or Storm. Multi-language support means storing Unicode characters in trie nodes, and different countries may need different tries stored in country-specific CDNs.

A good answer clearly explains the trie, the two key optimisations and their cost, the offline pipeline, and sharding by actual traffic distribution rather than by alphabet.

Key numbers

Latency target
≤ 100 ms per keystroke
Suggestions shown
Top 5
Average QPS
≈ 24,000
Peak QPS
≈ 48,000
New data per day
≈ 0.4 GB

Key terms

Trie
A tree in which each node represents a string prefix and each edge adds one character.
Top-k cache
The k most frequent completions stored directly on a trie node to make lookups constant time.
Data gathering service
The offline pipeline that turns raw query logs into frequency data and builds the trie.
Aggregator
A job that rolls raw log entries up into counts per query per time window.
Filter layer
A component in front of the trie cache that removes unwanted suggestions before they are returned.
Shard map manager
A service that maps prefix ranges to shards based on observed traffic distribution.
Data sampling
Logging only a fraction of requests to reduce processing and storage cost.

Common mistakes

  • Updating the trie synchronously on every search request.
  • Traversing the whole subtree for each prefix instead of caching top-k results per node.
  • Sharding strictly by first letter and ignoring the skewed letter distribution.
  • Forgetting a way to remove offensive suggestions immediately.
  • Logging every keystroke without sampling.
  • Ignoring client-side debouncing and browser caching when computing QPS.

Further study

  • Edward Fredkin, Trie Memory (1960)
  • Twitter typeahead.js
  • LinkedIn Cleo open-source typeahead
  • Apache Kafka and Spark Streaming

Now practise it