Limited time: AI code review, hints, mock interviews, whiteboard analysis, and all Pro features are unlocked. Enroll
⏱️ 33 min read

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:

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


4. Functional Requirements

Core (Top 3)

  1. Return top-k suggestions for a prefix - as the user types each character, return the 10 most relevant completions in under 100ms
  2. Rank by popularity and freshness - suggestions reflect both historical popularity and real-time trending queries
  3. 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


5. Non-Functional Requirements

Core

Below the Line


6. Core Entities


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:

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:

  1. User types how; the browser sends the prefix to the load balancer
  2. Any replica will do β€” every instance holds an identical copy of the trie and a lookup mutates nothing
  3. The service walks three edges, h β†’ o β†’ w, which is three pointer hops
  4. From that node it walks the subtree beneath it, gathering the complete queries that hang below
  5. 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:


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:

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:

  1. A user submits a search β€” not a keystroke, a real submitted query
  2. The Query Logger appends {query, user_id, timestamp} to query_events and returns immediately
  3. Once an hour the Aggregation Job reads the events since its last run
  4. It groups by query string and adds the counts into query_counts
  5. 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:


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:

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:

  1. After the Aggregation Job writes new totals, the Trie Builder reads the full query_counts table
  2. It builds a fresh trie containing every query, with each node’s counts attached
  3. It writes the finished structure out as a versioned, read-only snapshot
  4. Each Trie Service replica downloads it and flips a pointer from the old trie to the new one
  5. 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:


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:

  1. Breaking news: β€œearthquake tokyo” starts trending β†’ thousands of searches hit Kafka topic
  2. Flink streaming job: tumbling 5-minute window counts queries. Detects count_5min / count_7d_avg > 10x β†’ marks as trending
  3. Trending queries get a score boost: effective_score = historical_score + (trend_multiplier Γ— velocity)
  4. Next Trie rebuild (every 15 min): trending query’s boosted score pushes it into top-k lists at relevant prefix nodes
  5. For faster surfacing: a separate β€œtrending” overlay β€” a small Redis sorted set of trending queries checked alongside the Trie
  6. 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?

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:

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:

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):

GDPR compliance checklist:


12. Design Self-Audit


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
  1. Browser debounces keystrokes (100-150ms) to avoid flooding the server
  2. Request hits any Trie Service replica (stateless β€” all hold the same snapshot)
  3. Redis cache absorbs ~80% of traffic for popular prefixes
  4. On cache miss, Trie lookup is O(prefix_length) β€” effectively constant time
  5. Trending Cache merge ensures fresh viral queries appear without waiting for rebuild
  6. 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
  1. Query Logger is fire-and-forget β€” doesn’t block the search response
  2. Kafka provides durability; if Flink lags, events buffer safely
  3. Flink maintains in-memory state of windowed counts, flushes periodically
  4. Trie Builder runs every 15 minutes, reads aggregated counts, produces a new snapshot
  5. 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):

  1. User types a prefix β€” keystroke sent to CDN (static prefix cache for top queries)
  2. CDN cache miss β€” request routed through Load Balancer to the Trie Service (sharded replicas)
  3. Trie Service looks up suggestions β€” checks Redis Cache for hot prefixes, falls back to in-memory trie traversal
  4. Trending Cache checked β€” Redis Sorted Set injects trending/breaking queries that haven’t aged into the main trie yet
  5. Top-K suggestions returned β€” ranked by popularity score, served in <50ms

How it works end-to-end (update path):

  1. Query Logger captures searches β€” every completed search published to Kafka
  2. Stream Aggregator tallies counts β€” rolling window aggregation updates Cassandra (Popularity Store)
  3. Trending Detector identifies spikes β€” velocity detection flags sudden surges, updates Trending Cache in real-time
  4. Trie Builder rebuilds periodically β€” batch job reads Popularity Store, constructs new trie version, swaps into Trie Service (blue-green)

Discussion

Newest first
You

Free system design + DSA prep. If it helped you crack an interview, consider supporting.

SensAI SensAI
Beta
Listening...
Tap mic to stop voice mode

Shape what we build next

Every piece of feedback is read by the team and directly influences our roadmap.

What type of feedback?

Install SystemCraft

Add to your home screen for instant access, offline reading, and a distraction-free experience.

Offline reading Faster loads No browser tabs App-like feel

Unlock AI Features

One click to activate - no payment, no credit card. Just sign in and you're in.

AI code review and hints
SensAI chat assistant
AI mock interviews
Whiteboard analysis
100% free during early access