Designing Search Autocomplete / Typeahead
Difficulty: Intermediate Topics: Trie, Prefix Matching, Ranking, Caching, Real-time Trending Asked at: Google, Amazon, Microsoft, LinkedIn, Uber, Flipkart Prerequisites:Caching, Database Indexing, and Scalability
1. Understanding the Problem
Search autocomplete predicts what a user is about to type and suggests completions in real-time as they press each key. It powers the dropdown under every search bar β Google, Amazon product search, YouTube, LinkedIn people search. The core challenge: return the top-k most relevant suggestions for any prefix in under 100ms, while continuously learning from billions of new queries to keep suggestions fresh and trending-aware.
Real examples: Google Search Suggestions, Amazon product typeahead, YouTube search, LinkedIn search, Spotify song search.
2. Naive First Cut
flowchart LR
USER["User types prefix"]:::client
API["API Server"]:::service
DB[("SQL DB<br/>all queries + counts")]:::data
USER --> API
API -->|"SELECT WHERE query LIKE 'pre%'<br/>ORDER BY count DESC LIMIT 10"| DB
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Store every query with its count in a SQL table. On each keystroke, run a LIKE prefix query ordered by count.
Why this breaks:
- LIKE βprefix%β on billions of rows is too slow even with B-tree indexes (full scan for short prefixes like βaβ)
- A keystroke fires every 50-100ms - SQL canβt keep up at millions of concurrent users
- No real-time trending - counts only update in batch, so todayβs viral topic wonβt surface for hours
- Single DB becomes the bottleneck - no horizontal scaling for read-heavy prefix lookups
- No personalization - everyone sees the same suggestions regardless of context
- Network round-trip to DB on every keystroke adds unacceptable latency
The rest of the doc evolves this into a Trie-based service with in-memory prefix lookups, async count aggregation, and a caching layer that serves most requests without hitting the data tier.
3. Prior Art Weβre Drawing From
- Google Autocomplete - Uses a combination of query popularity, freshness, and user context. Suggestions update in real-time for trending queries using a streaming pipeline separate from the batch popularity index. (Google Blog)
- LinkedIn Typeahead - Built a distributed Trie service called βGaleneβ that serves prefix-based entity search (people, companies, jobs) with sub-50ms P99. Uses a two-level architecture: coarse-grained sharding by prefix + fine-grained in-memory Tries per shard. (LinkedIn Engineering)
- Facebook Unicorn (Social Graph Search) - Typeahead over a social graph combines prefix matching with social proximity scoring (friends-of-friends rank higher). The ranking signal isnβt just popularity but personalized affinity. (Facebook Engineering)
- Elasticsearch Completion Suggester - Uses FST (Finite State Transducer) data structure internally for prefix lookups with weighted suggestions, serving sub-5ms responses from an in-memory structure. Common choice for product search typeahead.
4. Functional Requirements
Core (Top 3)
- Return top-k suggestions for a prefix - as the user types each character, return the 10 most relevant completions in under 100ms
- Rank by popularity and freshness - suggestions reflect both historical popularity and real-time trending queries
- Update suggestions with new queries - when users search for something new (a breaking event, a new product), it should appear in suggestions within minutes, not hours
Below the Line
- Personalized suggestions (based on user history)
- Spell correction / fuzzy matching
- Multi-language support
- Offensive content filtering
- Category-aware suggestions (products vs pages vs people)
5. Non-Functional Requirements
Core
- Low latency: P99 < 100ms (users expect instant response on each keystroke)
- High availability: 99.99% - autocomplete is on the critical search path
- Scale: 100K+ prefix lookups per second; 10B+ queries/day feeding the popularity model
- Freshness: Trending queries surface within 5-15 minutes
Below the Line
- Eventual consistency is acceptable (a few minutes stale is fine)
- Multi-region serving (CDN-friendly for static prefix results)
- Graceful degradation under load (return cached stale results rather than fail)
6. Core Entities
- Query - a search string submitted by a user (the raw event)
- PrefixNode - a node in the Trie representing a character in a prefix path
- Suggestion - a complete query string with its popularity score and metadata
- TrendingQuery - a query whose recent velocity exceeds its historical baseline
- QueryAggregate - a time-bucketed count for a query string (used for ranking)
7. API / System Interface
GET /v1/suggestions?prefix=<string>&limit=10
Authorization: Bearer <token> (optional - for personalization)
Response:
{
"prefix": "how to des",
"suggestions": [
{"text": "how to design a url shortener", "score": 9842},
{"text": "how to design distributed systems", "score": 7231},
{"text": "how to design uber", "score": 6890}
],
"trending": ["how to design ai agents"]
}
POST /v1/queries (internal - logs a completed search)
Body: {"query": "how to design uber", "userId": "u123", "timestamp": 1720000000}
Response: 202 Accepted
Security notes: rate-limit prefix lookups per IP/session to prevent scraping the full suggestion index. Filter offensive and legally restricted terms server-side before returning suggestions.
8. High-Level Design
FR1: Return top-k suggestions for a prefix
Given the characters typed so far, return the ten best complete queries that start with them.
A relational LIKE 'how%' can answer that, and it has to scan, so the structure that makes
prefix lookup natural is a trie: a tree where each edge is a character and the path from
the root spells the prefix.
π‘ Trie = prefix tree. Walking βhβ then βoβ then βwβ lands you
on the node for βhowβ, and everything completable from βhowβ hangs below that node.
Build one trie, in memory, in one service. No cache, no sharding, and no pre-computation beyond the tree itself β those answer the latency and scale targets in the non-functional list, so each is earned in a deep dive.
New components:
- Trie Service: holds the whole trie in memory. Walk the prefix character by character, then walk the subtree below that node to collect completions and pick the ten with the highest stored counts.
flowchart LR
USER["User Browser"]:::client
LB["Load Balancer"]:::edge
TRIE["Trie Service<br/>whole trie in memory"]:::service
USER -->|"1. Ask for prefix how"| LB
LB -->|"2. Route to any replica"| TRIE
TRIE -->|"3. Return top ten"| USER
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#38bdf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
| Color | Meaning |
|---|---|
| Purple | Client |
| Blue | Edge / Load Balancer |
| Green | Application Service |
Flow:
- User types
how; the browser sends the prefix to the load balancer - Any replica will do β every instance holds an identical copy of the trie and a lookup mutates nothing
- The service walks three edges,
h β o β w, which is three pointer hops - From that node it walks the subtree beneath it, gathering the complete queries that hang below
- It sorts those by their stored counts, takes ten, and returns them
Why hold it in memory rather than query a database per keystroke? Autocomplete fires on every character, so a five-letter word is five requests. Anything involving a disk seek per request loses the 100ms budget on the first syllable. Keeping the structure in RAM makes the lookup a few pointer dereferences.
What we have deliberately left broken. For a corpus of a million queries this is a complete, fast autocomplete. At the scale in the requirements it has three separate problems:
- The trie does not fit in one machineβs RAM. The requirements imply a corpus of over a billion distinct queries feeding the popularity model. A trie over that, with a node per character and child pointers, runs well past what a single node can hold, so βone service holding the whole trieβ stops being deployable. That is Deep Dive 1.
- Step 4 is the slow part, and it gets slower as the prefix gets shorter. Walking to the node is trivial; walking everything underneath it is not. For
hthat subtree is a large fraction of the entire corpus, so the cheapest possible keystroke β the first one β is the most expensive query we serve. Deep Dive 1 removes the walk; Deep Dive 4 deals with the traffic skew that short prefixes also cause. - Every keystroke does this work from scratch. At 100K lookups/sec with no cache anywhere, identical prefixes are recomputed continuously, and a handful of prefixes account for most of the traffic. That is Deep Dive 3.
FR2: Rank by popularity
FR1 sorted by βstored countsβ without saying where they came from. This is where. Ranking by popularity needs one thing: a record of what people actually searched for, counted.
New components:
- Query Logger: records every search a user issues. It sits off the response path, so logging never slows a search down.
query_eventstable: the raw log β one row per search, with the query text, the user, and a timestamp.- Aggregation Job: runs hourly, groups the raw events by query string, and writes the totals to
query_counts.
flowchart LR
USER["User Browser"]:::client
QLOG["Query Logger"]:::service
EVENTS[("query events<br/>raw log")]:::data
JOB["Aggregation Job<br/>hourly"]:::async
COUNTS[("query counts<br/>totals per query")]:::data
USER -->|"1. Submit a search"| QLOG
QLOG -->|"2. Append a row"| EVENTS
JOB -->|"3. Read the last hour"| EVENTS
JOB -->|"4. Write totals"| COUNTS
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
classDef async fill:#3a2a4c,stroke:#c084fc,color:#e2e8f0
Flow:
- A user submits a search β not a keystroke, a real submitted query
- The Query Logger appends
{query, user_id, timestamp}toquery_eventsand returns immediately - Once an hour the Aggregation Job reads the events since its last run
- It groups by query string and adds the counts into
query_counts - Those totals are the numbers FR1βs step 5 sorts by
Why count submitted searches rather than every keystroke? A user typing weather produces
seven prefixes but wants one thing. Counting keystrokes would rank w, we and wea as
enormously popular queries, which they are not β they are fragments of one intention. Counting
at submission is the only point where what we record matches what somebody meant.
What we have deliberately left broken. The counts are right and how we get them is not:
- An hour is the wrong unit. The non-functional list asks for trending queries to surface within 5-15 minutes, and this design cannot beat its own schedule: a query that spikes at 12:01 is invisible until the 13:00 run. Averaged over an hour, a sharp spike also barely moves a total that is dominated by all-time popularity, so a breaking event may not rank even after the job runs. That is Deep Dive 2.
- We are storing every search next to the identity of who made it.
query_eventsis, by construction, a per-user search history in plain text, retained indefinitely because the Aggregation Job never deletes anything. That is a serious privacy liability rather than a scaling one, and it is Deep Dive 5.
FR3: Get new queries into the suggestions
A query nobody has ever searched for is not in the trie, so FR1 can never suggest it. Something has to carry new queries and changed counts from FR2βs totals into the structure FR1 serves.
The trie is read on every keystroke by every replica, so mutating it in place while it is being walked means locking on the hottest path in the system. Build a new one instead and swap it.
New components:
- Trie Builder: reads
query_counts, constructs a complete trie from scratch, and publishes it as an immutable snapshot. Each Trie Service replica loads the new snapshot and switches to it between requests, which needs no locking because the old trie is never modified.
flowchart LR
COUNTS[("query counts")]:::data
BUILD["Trie Builder"]:::async
SNAP[("Trie snapshot<br/>immutable")]:::data
TRIE["Trie Service"]:::service
USER["User Browser"]:::client
COUNTS -->|"1. Read all totals"| BUILD
BUILD -->|"2. Publish a snapshot"| SNAP
SNAP -->|"3. Replicas load it"| TRIE
TRIE -->|"4. Serve from the new trie"| USER
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
classDef async fill:#3a2a4c,stroke:#c084fc,color:#e2e8f0
Flow:
- After the Aggregation Job writes new totals, the Trie Builder reads the full
query_countstable - It builds a fresh trie containing every query, with each nodeβs counts attached
- It writes the finished structure out as a versioned, read-only snapshot
- Each Trie Service replica downloads it and flips a pointer from the old trie to the new one
- Requests in flight finish against the old trie and are unaffected; the old snapshot is dropped once nothing references it
Why rebuild the whole thing rather than update the nodes that changed? Because the readers never block. An immutable snapshot means a lookup needs no lock, no version check, and no retry, and a bad build can be rolled back by pointing at the previous snapshot. The cost is that the smallest possible change requires a full rebuild, which is a trade worth making while rebuilds are cheap and worth revisiting when they are not.
What we have deliberately left broken. New queries do now reach suggestions, and the delay is the problem:
- The pipeline is two batch stages deep. A new query waits for the hourly Aggregation Job, then for the next build, so worst case it is over an hour before anyone can be suggested it. The requirements ask for minutes. That is Deep Dive 2.
- Rebuilding everything to add one query gets absurd at scale. A full rebuild over a billion-query corpus is a large job, and its cost is the same whether one query changed or a million did β which is exactly what stops us from simply running it more often to fix the point above. Deep Dive 1βs structure is what makes an affordable build possible.
9. Technology Choices
| Tier | Purpose | Stores | Access Pattern | Primary Pick | Alternatives |
|---|---|---|---|---|---|
| Prefix index | In-memory prefix lookups | Top-k completions per prefix node | Point lookup by prefix string | Custom distributed Trie service | Elasticsearch Completion Suggester / Redis sorted sets |
| Query log store | Raw query event stream | Every search query with timestamp | Append-only writes | Kafka / Kinesis | Pulsar / Redpanda |
| Aggregation | Count queries over time windows | Query frequency per time bucket | Streaming aggregation | Flink / Kafka Streams | Spark Structured Streaming |
| Popularity store | Aggregated query counts | query -> count + trend score | Batch read for Trie rebuild | Cassandra / DynamoDB | Postgres (if scale is moderate) |
| Cache | Hot prefix results | prefix -> top-10 suggestions | Key-value lookup | Redis / Memcached | Cloudflare Workers KV |
| Analytics DB | Historical query analytics | Long-term query logs | OLAP queries | ClickHouse / BigQuery | Redshift / Snowflake |
Why a custom Trie over Elasticsearch? For pure prefix completion at massive scale (100K+ QPS), an in-memory Trie with pre-computed top-k at each node is 10-50x faster than ES Completion Suggester because it avoids serialization and network hops. ES is the right call if you also need fuzzy matching, typo correction, and faceted search alongside autocomplete.
10. Data Modeling
Custom Distributed Trie (Prefix Index β in-memory):
Structure per node:
char: 'h'
children: { 'o' β node, 'i' β node, ... }
top_k: [("how to design uber", 9842), ("how to design systems", 7231), ...]
is_terminal: boolean
Key insight: top-k suggestions are PRE-COMPUTED at each node during Trie rebuild.
No traversal needed at query time β just navigate to the prefix node and return its top_k list.
Kafka (Query Log Store β raw search events):
Topic: search-queries (partitioned by query_hash % N)
Key: query_text_normalized
Value: { query, user_id, timestamp, result_clicked, session_id }
Cassandra / DynamoDB (Popularity Store β aggregated counts):
Table: query_counts
PK: query_text_normalized
Columns: total_count, count_1h, count_24h, count_7d, trend_score, last_updated
Redis (Cache β hot prefix results):
Key: "suggest:{prefix}" β JSON list of top-10 suggestions with scores
TTL: 5 minutes (popular prefixes cached; long-tail misses go to Trie)
Example: "suggest:how to d" β [{"text":"how to design uber","score":9842}, ...]
Access Patterns:
| Query | Data Source | How |
|---|---|---|
| Get suggestions for prefix | Redis β Trie | Check cache first; on miss, query Trie service, cache result |
| Log a completed search | Kafka | Produce to search-queries topic (fire-and-forget) |
| Update popularity counts | Flink β Cassandra | Streaming aggregation over tumbling windows β upsert to query_counts |
| Rebuild Trie (periodic) | Cassandra β Trie | Every 15 min: read top 10M queries by score, rebuild Trie in memory, swap atomically |
How a Trending Query Surfaces in Suggestions Within Minutes:
- Breaking news: βearthquake tokyoβ starts trending β thousands of searches hit Kafka topic
- Flink streaming job: tumbling 5-minute window counts queries. Detects
count_5min / count_7d_avg > 10xβ marks as trending - Trending queries get a score boost:
effective_score = historical_score + (trend_multiplier Γ velocity) - Next Trie rebuild (every 15 min): trending queryβs boosted score pushes it into top-k lists at relevant prefix nodes
- For faster surfacing: a separate βtrendingβ overlay β a small Redis sorted set of trending queries checked alongside the Trie
- User types βearthβ β Trie returns historical suggestions, PLUS Redis returns trending matches β merge and dedupe
11. Deep Dives
1) The trie will not fit on one machine, and short prefixes are the slowest lookups. Now what?
The problem: A Trie stores all possible search queries as a tree where each node is a character. For Google with 5B+ unique queries searched historically, a single Trie would need ~200-500GB of RAM. No single machine has that. Plus, one machine can serve ~50K QPS max β Google needs 500K+.
How a Trie works (simple version):
Insert: "system", "systems", "sync", "syntax"
Root
βββ s
βββ y
βββ s β t β e β m [β "system", popularity: 50000]
β βββ s [β "systems", popularity: 12000]
βββ n
βββ c [β "sync", popularity: 8000]
βββ t β a β x [β "syntax", popularity: 3000]
Lookup "sy": walk root β s β y β return top-k children:
["system" (50K), "systems" (12K), "sync" (8K), "syntax" (3K)]
The key optimization: Pre-compute top-k at each node.
Without pre-computation: lookup βsβ β must traverse ALL children (millions of words start with βsβ) β too slow.
With pre-computation: at build time, store the top 10 suggestions directly at each node:
Node "sy" stores: top_10 = ["system design", "system", "systems", "sync", "syntax", ...]
Lookup "sy" β jump to node β read pre-computed list β return immediately. O(L) where L = prefix length.
No traversal of children needed at query time.
Bad: What FR1 built β the whole trie in one process, walking the subtree below the prefix node on each request. It is genuinely excellent for a million queries and it fails twice over a billion. The corpus stops fitting: a node per character, with child pointers and counts, puts a billion-query trie well past any single machineβs RAM, so no instance size makes this deployable. And the cost is inverted β reaching the prefix node is a few pointer hops, but gathering what hangs beneath it means traversing a subtree that grows as the prefix gets shorter. The first character a user types is the most expensive lookup in the system, and it is the one every single session issues.
Good: Shard the Trie by first 1-2 characters of the prefix. Prefix βaβ goes to shard 1, βbβ to shard 2, etc. Each shard fits in memory (~50-100GB) and handles a subset of traffic. The load balancer routes based on the first character.
Great: Pre-compute the top-k suggestions at each Trie node during the build phase (not at query time). This means a lookup doesnβt need to traverse all children to find the best completions β theyβre already stored at the prefix node. Combined with sharding, this gives O(L) lookup with zero fan-out.
How the Trie is built (offline, not real-time):
1. Aggregation Job (every 15 min):
- Reads completed searches from Kafka
- Counts: {"system design": 50000, "system": 35000, "sync": 8000, ...}
2. Trie Builder:
- Creates fresh Trie from the counted queries
- At each node, computes and stores top-10 suggestions sorted by popularity
- Serializes to a binary format
3. Deploy:
- Push new Trie snapshot to all shard replicas (blue-green swap)
- Old Trie keeps serving until new one is loaded
- Atomic switch: old β new (zero downtime)
Why offline build, not live updates?
- Trie is a complex tree β concurrent writes cause locking and fragmentation
- Pre-computing top-k at every node requires a full traversal β canβt do incrementally
- Read performance is critical (<10ms). Live writes would add locks and slow reads.
- 15-minute freshness is fine for autocomplete (trending handles real-time separately)
LinkedInβs Cleo system uses this exact pattern to serve sub-50ms P99 at 100K+ QPS.
flowchart LR
REQ["prefix: 'sys'"]:::client
ROUTER["Prefix Router<br/>(route by first char)"]:::edge
S1["Shard 's'<br/>Trie in memory"]:::service
NODE["Node s-y-s<br/>top-10 pre-computed"]:::data
REQ -->|"1. Send prefix"| ROUTER
ROUTER -->|"2. Route to shard"| S1
S1 -->|"3. Walk trie to node"| NODE
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#38bdf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
2) A city starts searching one phrase at 12:01. Why is it invisible until 13:00?
The problem: When breaking news happens (India vs Australia match, earthquake, celebrity news), millions of people start searching for it immediately. If autocomplete only recomputes query counts in hourly batches, trending searches wonβt appear as suggestions for up to an hour β users miss the most relevant suggestions during peak interest.
The core question: How do you distinguish βtrending right nowβ from βalways popularβ? βWeatherβ is searched 300 times/hour EVERY day β thatβs not trending. βIndia vs Australia scoreβ is normally 100/hour but just spiked to 5000/5min β THATβS trending. Trending = fast-rising velocity, not high volume.
Bad: What FR2 built β an hourly batch job over the raw event log. The schedule is a hard floor on freshness: a phrase that starts spiking at 12:01 does not exist in any total until the 13:00 run, and then waits on FR3βs rebuild behind that. Against a requirement of 5-15 minutes, that is off by an order of magnitude. The averaging is the subtler half of the problem. An hourly total blends the spike into the other fifty-nine minutes, and it is then compared against all-time popularity, so a phrase with ten thousand searches in four minutes still ranks below everyday queries with steady large totals. Even after the job runs, the thing that is trending right now may not surface at all β the design cannot distinguish βpopularβ from βsuddenly popularβ, because a single cumulative count contains no notion of rate.
Good: Sliding window counts in Flink with 5-minute granularity. Compare current-window count against the 24-hour average. If current > 3x average, mark as trending. Latency: ~5 minutes. This works but detection only happens at window boundaries β if a spike starts at minute 2, you donβt know until minute 5.
Great: Use an exponential moving average (EMA) with a decay factor. Each new search event updates the EMA incrementally β no windows to maintain, no batch boundaries. A sudden spike causes the EMA to diverge sharply from the long-term average, triggering a trending alert within seconds.
How EMA works with a concrete example:
Suppose βindia vs australiaβ is normally searched ~100 times/hour (β0.03/sec baseline).
EMA formula: new_score = Ξ± Γ current_rate + (1 - Ξ±) Γ previous_score
Ξ± = 0.3 (decay factor β higher = more responsive to spikes)
Normal day (no match):
Score hovers around 0.03 β stable, no spike.
Match starts, people start searching:
Second 1: 5 searches β score = 0.3Γ5 + 0.7Γ0.03 = 1.52
Second 2: 8 searches β score = 0.3Γ8 + 0.7Γ1.52 = 3.46
Second 3: 15 searches β score = 0.3Γ15 + 0.7Γ3.46 = 6.92
Second 4: 20 searches β score = 0.3Γ20 + 0.7Γ6.92 = 10.84
Normal score: 0.03
Current score: 10.84
Ratio: 10.84 / 0.03 = 361x above baseline β TRENDING after just 4 seconds!
Why βweatherβ doesnβt falsely trigger:
"weather" β 300 searches/hour every day (β0.08/sec baseline):
Normal score: ~0.08
Today's score: ~0.08 (same as always)
Ratio: 0.08 / 0.08 = 1x β NOT trending (no spike, just steady volume)
Why EMA beats sliding windows:
- No fixed window boundaries β detects trends in seconds, not at 5-min intervals
- Captures velocity (how FAST searches are rising), not just total count
- Automatically decays β once the spike subsides, the EMA drops back to normal and the query is un-trended
- This is what Twitter uses for its Trends feature
How it integrates with the architecture:
User searches "india vs australia" β logged to Kafka
β
Trending Detector (Flink consumer):
- Reads every search event from Kafka
- Updates EMA for that query: score = 0.3 Γ rate + 0.7 Γ old_score
- Compares to 24-hour baseline
- If score / baseline > 3x β TRENDING
- Writes to Trending Cache (Redis Sorted Set):
ZADD trending_queries <score> "india vs australia score"
β
Next user types "ind" in the search bar:
- Trie returns standard suggestions: ["india population", "india news", ...]
- Trending Cache returns: ["india vs australia score" (score: 10.84)]
- MERGE: trending results get boosted to the top of suggestions
- User sees: "india vs australia score" as first suggestion β
3) At 100K lookups a second, how many times should we compute the same prefix?
The problem: Users type fast. Each keystroke triggers a suggestion request. βsystem designβ = 13 keystrokes = 13 requests in 3 seconds from ONE user. With 100K concurrent users typing, thatβs 500K+ requests/sec hitting the Trie service. Even with sharding, this is expensive.
Why caching helps: Most prefixes are repeated constantly. βhowβ, βwhatβ, βwhyβ, βbestβ are typed thousands of times per second by different users. If we cache these results, 80%+ of requests never hit the Trie service.
Bad: What FR1 built β no cache at all, so all 100K lookups/sec are computed from scratch. Query distributions are extremely skewed: a small set of prefixes accounts for most of the traffic, which means we are recomputing an identical answer for the same few thousand prefixes thousands of times a second. The answers are also nearly static, since popularity totals only change when FR3 publishes a new snapshot, so almost all of that work produces a byte-identical result to the one we produced a millisecond earlier. This is pure duplicated effort, and it is worst precisely on the short prefixes Deep Dive 1 already identified as the most expensive walks.
Good: Cache the top 10K prefixes in Redis with a 10-minute TTL. ~80% of lookups hit the cache. But when a hot key expires, thousands of requests simultaneously miss and slam the Trie service (cache stampede).
Concrete example of cache stampede:
Time 0:00 β "how" is cached in Redis (TTL = 10 min)
Time 9:59 β 5000 users type "how" per second, all served from cache. Trie service idle.
Time 10:00 β TTL expires. Redis key "prefix:how" disappears.
Time 10:00.001 β 5000 requests arrive. ALL miss cache. ALL hit Trie service simultaneously.
Trie service: 50 QPS normal β sudden 5000 QPS spike β overloaded, timeouts.
Great: Use probabilistic early expiration (cache stampede protection):
Each cached entry stores: { suggestions: [...], expiresAt: 10:00:00 }
Request arrives at 9:59:45 (15 seconds before expiry):
- Generate random number: rand() = 0.85
- Threshold: (time_remaining / TTL) = 15/600 = 0.025
- 0.85 > 0.025 β serve from cache normally (99% of requests)
Request arrives at 9:59:55 (5 seconds before expiry):
- Generate random number: rand() = 0.02
- Threshold: 5/600 = 0.008
- 0.02 > 0.008... but close! Some requests will trigger early refresh.
ONE lucky request gets selected β fetches fresh data from Trie β updates cache with new TTL
All other requests continue serving stale (but valid) data during this refresh.
Result: cache never expires for everyone simultaneously. One request refreshes proactively.
Additionally β request coalescing (single-flight):
If cache is empty and 100 requests arrive for βhowβ simultaneously:
- WITHOUT coalescing: 100 requests all hit Trie service β 100 duplicate lookups
- WITH coalescing: first request goes to Trie, other 99 wait for that ONE response, all get the same result
Implementation (Go's singleflight / Java's CacheLoader):
lock = mutex_per_key["prefix:how"]
if lock.tryAcquire():
result = trie.lookup("how")
cache.set("prefix:how", result, TTL=10min)
lock.release()
else:
result = lock.waitForResult() // wait for the one ongoing fetch
return result
4) Far more people type s than type x. How do we shard without one node melting?
The problem: When you shard the Trie by first character, some shards get crushed. In English, words starting with βsβ are 3x more common than words starting with βxβ or βzβ. The βsβ shard gets 3x the traffic β itβs a hot partition.
Concrete data (English word frequency by first letter):
Letter: s β 12% of all queries
Letter: c β 8%
Letter: p β 7%
Letter: a β 7%
...
Letter: x β 0.2%
Letter: z β 0.3%
If you shard A-Z (26 shards), the βSβ shard handles 12% of ALL traffic while βXβ handles 0.2%. Thatβs 60x imbalance.
Bad: Deep Dive 1βs sharding, taken literally β split on the first character, a* to one
shard, b* to the next. It divides the data and does not divide the load, because letter
frequency in real queries is nowhere near uniform: s, c and t start a large share of
English words while x, z and q start almost none. The shard owning s takes something like
an order of magnitude more traffic than the one owning x, so we size the whole fleet for the
worst shard and leave most of it idle. Splitting is also unhelpful in the wrong direction β the
hot shard is hot because of one popular letter, and no amount of extra shards changes which
shard that letter maps to.
Good: Shard by first 2 characters (βsaβ, βsbβ, β¦, βszβ). More even distribution (676 possible combinations), but still some skew β βstβ (start, stop, store, stock, streamβ¦) is hotter than βsxβ.
Great: Weighted consistent hashing based on observed traffic per prefix range:
Control Plane monitors QPS per shard:
"sa-sf" shard: 8000 QPS (hot β lots of "search", "service", "system")
"xa-xz" shard: 200 QPS (cold)
Action: split the hot shard further:
"sa-sc" β Shard A (with 3 replicas)
"sd-sf" β Shard B (with 3 replicas)
"xa-xz" β Shard C (with 1 replica β doesn't need more)
The router's routing table updates:
prefix "sa..." β Shard A (any of 3 replicas)
prefix "se..." β Shard B (any of 3 replicas)
prefix "xi..." β Shard C
Result: hot prefixes get more shards AND more replicas. Cold prefixes share resources.
Key insight: The routing isnβt static. A control plane continuously monitors QPS per shard and rebalances by splitting hot ranges or adding replicas. This is how Googleβs Bigtable tablets auto-split on hot keys.
5) We are keeping a plaintext search history per user forever. How do we not?
The problem: To build autocomplete suggestions, you need to know what people search for. But search queries are deeply personal β medical conditions, financial problems, relationship issues. Logging everything with user identity creates a privacy liability.
The tension: You NEED aggregate counts (βhow many people searched for X?β) to rank suggestions. But you DONβT need to know WHO searched for WHAT.
Bad: What FR2 built β {query, user_id, timestamp} appended for every search and never
deleted, because the Aggregation Job only ever reads forward. That table is a complete,
plaintext, permanently retained search history for every user, which is among the most sensitive
data a person generates: searches about health, finances, sexuality and legal trouble all sit in
it in the clear. Nothing in the design needs it. Ranking uses totals per query string, and the
user_id column is never read by anything, so we are carrying the entire liability for no
functional benefit. It is also a compliance problem with teeth β a GDPR erasure request means
locating and purging one personβs rows across the full history, and every derived aggregate
computed from them. Logging raw keystrokes rather than submitted searches, which FR2 explicitly
avoided, would multiply all of this by the length of the average query.
Good: Only log completed queries (when the user hits Enter or clicks a suggestion). Anonymize after 24 hours by stripping user IDs and keeping only aggregate counts:
Raw log (first 24 hours):
{ userId: "user_123", query: "diabetes symptoms", timestamp: "..." }
After 24 hours (anonymized):
{ query: "diabetes symptoms", count: 4521 } β no userId, just aggregate
Why keep userId for 24 hours? To deduplicate (same user searching same thing 10 times = count as 1).
After dedup, you don't need the userId anymore.
Great: Differential privacy at the aggregation layer β add calibrated random noise to query counts before theyβre used for ranking:
True count for "depression help": 847 searches today
Noise added: random(Β±50)
Published count used for ranking: 847 + 23 = 870 (or 847 - 31 = 816)
Why this matters:
- If count = exactly 1, you can infer ONE specific user searched it
- With noise, count=1 might actually be 0 (noise added) β can't determine if anyone searched it
- At high counts (1000+), noise is negligible β rankings stay accurate
- At low counts (1-10), noise dominates β individual privacy protected
Appleβs approach (local differential privacy):
- Noise is added ON THE USERβS DEVICE before sending to the server
- Server never sees the true query β only a βnoisyβ version
- Used for emoji suggestions, QuickType keyboard
- Even stronger guarantee: server literally cannot know what any individual searched
GDPR compliance checklist:
- Right to be forgotten β delete userId from raw logs on request (within 24h of anonymization, thereβs no userId left anyway)
- Data minimization β log only completed queries, not keystrokes
- Purpose limitation β counts used only for suggestion ranking, not advertising
- Retention limits β raw logs deleted after 24h, aggregates kept for ranking
12. Design Self-Audit
- Dedicated search index? The Trie IS the search index β purpose-built for prefix lookups. No need for a general-purpose search engine for this specific use case.
- Stale reads after writes? Yes β a newly searched query takes 3-15 minutes to appear in suggestions. Acceptable trade-off documented in FR3 (trending path reduces this to 3-5 min for viral queries).
- Single points of failure? Trie Service is replicated (multiple shards, each with replicas). Redis Cache has replicas. Kafka is multi-broker. Flink uses checkpointed state. No single-machine SPOF.
- Dead-letter / reconciliation? Kafka consumer offset tracking + Flink checkpoints. If processing fails, events replay from last checkpoint. No silent data loss.
- Cost at scale? The Trie is in-memory β at 1B unique queries with top-k pre-computed, expect 200-500GB total across shards. At cloud memory pricing (~$10/GB/month), thatβs $2K-5K/month for the Trie tier. Affordable for any company running a search product at scale.
13. Core Flows
Flow 1: Prefix Lookup (read path)
sequenceDiagram
participant U as User Browser
participant LB as Load Balancer
participant TS as Trie Service
participant RC as Redis Cache
participant TC as Trending Cache
U->>LB: GET /suggestions?prefix=earth
LB->>TS: route to any replica
TS->>RC: GET cache key "earth"
alt Cache Hit
RC-->>TS: top-10 suggestions
else Cache Miss
RC-->>TS: null
TS->>TS: Traverse Trie to node e-a-r-t-h and read pre-computed top-10
TS->>RC: SET "earth" with 10min TTL
end
TS->>TC: ZRANGEBYSCORE prefix match "earth*" LIMIT 5
TC-->>TS: trending matches
TS->>TS: Merge static + trending and re-rank
TS-->>LB: top-10 merged suggestions
LB-->>U: JSON response in 20-50ms
- Browser debounces keystrokes (100-150ms) to avoid flooding the server
- Request hits any Trie Service replica (stateless β all hold the same snapshot)
- Redis cache absorbs ~80% of traffic for popular prefixes
- On cache miss, Trie lookup is O(prefix_length) β effectively constant time
- Trending Cache merge ensures fresh viral queries appear without waiting for rebuild
- Response includes both static and trending results, clearly labeled
Non-obvious failure path: If the Trending Cache (Redis) is down, the Trie Service gracefully degrades β it returns only static Trie results. Suggestions are slightly stale but never unavailable.
Flow 2: Query Ingestion and Count Update (write path)
sequenceDiagram
participant U as User Browser
participant QL as Query Logger
participant K as Kafka
participant FA as Flink Aggregator
participant PS as Popularity Store
participant TB as Trie Builder
participant TS as Trie Service
U->>QL: POST /queries (search completed)
QL->>K: publish query event
K->>FA: consume batch (every 30s)
FA->>FA: Update windowed counts
FA->>PS: Write updated counts every 5min
TB->>PS: Read top-k per prefix (every 15min)
TB->>TB: Build new immutable Trie snapshot
TB->>TS: Hot-swap to new snapshot
- Query Logger is fire-and-forget β doesnβt block the search response
- Kafka provides durability; if Flink lags, events buffer safely
- Flink maintains in-memory state of windowed counts, flushes periodically
- Trie Builder runs every 15 minutes, reads aggregated counts, produces a new snapshot
- Hot-swap means the Trie Service atomically switches pointers β no downtime, no partial state
Non-obvious failure path: If Flink crashes mid-window, it replays from Kafka offset (exactly-once semantics via checkpointing). Counts may temporarily lag by one window but never lose data.
14. Final Architecture
flowchart TD
USER["User Browser"]:::client
CDN["CDN<br/>(static prefix cache)"]:::edge
LB["Load Balancer"]:::edge
TRIE["Trie Service<br/>(sharded replicas)"]:::service
CACHE["Redis Cache<br/>(hot prefixes)"]:::data
HOTCACHE["Trending Cache<br/>(Redis Sorted Set)"]:::data
QLOG["Query Logger"]:::service
STREAM["Kafka"]:::async
AGG["Stream Aggregator"]:::async
TREND["Trending Detector"]:::async
POPDB[("Popularity Store<br/>Cassandra")]:::data
BUILDER["Trie Builder<br/>(periodic)"]:::async
USER -->|"Type prefix query"| CDN
CDN -->|"Cache miss"| LB
LB -->|"Forward to trie svc"| TRIE
TRIE -->|"Lookup prefix in cache"| CACHE
TRIE -->|"Check trending cache"| HOTCACHE
USER -->|"Submit full search"| QLOG
QLOG -->|"Publish query event"| STREAM
STREAM -->|"Aggregate popularity"| AGG
STREAM -->|"Detect trending spikes"| TREND
AGG -->|"Update popularity store"| POPDB
TREND -->|"Refresh trending cache"| HOTCACHE
BUILDER -->|"Rebuild trie from data"| POPDB
BUILDER -->|"Push new trie version"| TRIE
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#38bdf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
classDef async fill:#3a2a4c,stroke:#c084fc,color:#e2e8f0
How it works end-to-end (query path):
- User types a prefix β keystroke sent to CDN (static prefix cache for top queries)
- CDN cache miss β request routed through Load Balancer to the Trie Service (sharded replicas)
- Trie Service looks up suggestions β checks Redis Cache for hot prefixes, falls back to in-memory trie traversal
- Trending Cache checked β Redis Sorted Set injects trending/breaking queries that havenβt aged into the main trie yet
- Top-K suggestions returned β ranked by popularity score, served in <50ms
How it works end-to-end (update path):
- Query Logger captures searches β every completed search published to Kafka
- Stream Aggregator tallies counts β rolling window aggregation updates Cassandra (Popularity Store)
- Trending Detector identifies spikes β velocity detection flags sudden surges, updates Trending Cache in real-time
- Trie Builder rebuilds periodically β batch job reads Popularity Store, constructs new trie version, swaps into Trie Service (blue-green)
Discussion
Newest first