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.
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.
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.
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
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