System Design Problem

Design Google Typeahead / Autocomplete

Commonly Asked By:GoogleMicrosoftAmazonMeta

Interview Setup

Interview Prompt

Design a typeahead/autocomplete system like Google Search suggestions or Amazon product search. As the user types, return the top 5-10 suggestions ranked by popularity within 50ms.

Clarifying Questions (ask before designing)

QuestionWhy it matters
What's the query volume and corpus size?A volume of ~350K average QPS (peaking at ~1.7M) with 100M terms fits within an in-memory sharded trie (6 shard tiers x ~50 GB each x 3 replicas = ~900 GB aggregate serving RAM), whereas 1B terms requires multi-character prefix sharding (such as 256 to 4,096 shards) and coordinated blue-green rebuilds.
Static corpus or dynamic (updates from user queries)?Dynamic corpora require streaming frequency updates, where maintaining a mutable trie in the serving path introduces lock contention and operational complexity, so we serve from an immutable snapshot with an hourly batch rebuild plus a 5-minute prefix-aware trending overlay.
Exact prefix match only, or fuzzy/typo tolerance?Fuzzy matching introduces edit-distance computation, which typically routes through a dedicated spell-check path rather than the latency-sensitive trie lookup.
Personalized suggestions or global top-K?Personalization adds user context to ranking as a secondary lookup stage evaluated after the initial prefix match.

Scope

In scope

  • Prefix-based autocomplete with top-10 results (top-K generalized)
  • Trie or equivalent prefix index structure
  • Frequency-based ranking with approximate updates
  • Low-latency serving architecture (<50ms p99)
  • Prefix sharding for horizontal scale

Out of scope (state explicitly)

  • Full web search results (only suggestions)
  • Natural language understanding / intent detection
  • Sophisticated multi-language morphology and algorithmic stemming (language and country-specific trie partitions and filtering are supported, but deep language-specific grammatical morphology is out of scope)
  • Frontend UI implementation

Functional Requirements

Ask your interviewer about prefix matching, ranking, and latency targets. Confirm whether personalization and fuzzy matching are in scope.

  • As the user types, suggest the top 5 to 10 matching search queries in real time.
  • Rank suggestions according to historical search frequency, decayed popularity, and trending velocity.
  • Support prefix matching so that typing "sys" returns completions such as "system design" and "system architecture".
  • Update suggestions dynamically based on emerging queries within 5 minutes via a trending overlay, and refresh broader popularity shifts through hourly batch snapshot rebuilds.
  • Handle multi-language queries with dedicated language and country filtering (sophisticated grammatical morphology and stemming are out of scope).
  • Provide optional personalized suggestions derived from the individual user's recent search history.

Non-Functional Requirements

Sub-50ms suggestions at scale are the bar. Your interviewer cares about index size, hot-prefix sharding, and whether personalization is in scope for this round.

  • Ultra-Low Latency: Backend autocomplete p99 latency must remain strictly under 50ms, ensuring total client-perceived response time (including network transport and UI rendering) stays within 50 to 100 ms of keystroke entry.
  • High Availability: Maintain 99.99% uptime for all lookup endpoints.
  • Scalability: Support ~350K average autocomplete requests per second, peaking at ~1.7M requests per second.
  • Freshness: Surface breaking trending queries within 5 minutes via the trending overlay, while reflecting broader corpus popularity shifts through hourly core trie snapshot rebuilds.
  • Consistency: Eventual consistency is acceptable because slightly stale suggestions from an hourly snapshot do not compromise user experience.
  • Fault Tolerance: The system degrades gracefully by continuing to serve previous immutable trie snapshots if background update pipelines stall.

Capacity Estimations

Query volume and dictionary size tell you whether a trie fits in memory and how many shards you need.

MetricCalculationValue
DAUGiven500M
Searches / day500M DAU x 105B
Autocomplete requests / searchGiven6 (avg 6 keystrokes before selecting)
Autocomplete requests / day5B searches x 630B
Autocomplete requests / sec (avg)30B ÷ 86400~350K
Autocomplete requests / sec (peak)350K x 5x peak factor~1.7M
Unique query termsGiven100M
Avg query lengthGiven20 characters
Trie storage (top-10 pruned)6 shard tiers x ~50 GB each~300 GB logical (~900 GB aggregate across 3 replicas)

Latency budget (p99 < 50 ms)

StageBudgetNotes
API gateway + prefix router5 msAuthentication, rate limiting, and consistent-hash shard routing
Trie lookup + top-K return15 msO(prefix length) lookup querying main immutable trie and trending overlay, returning top-10
Personalization overlay (optional)10 msRedis recent-search boost, skipped upon reaching timeout threshold
Network + serialization10 ms
  • Same-region transport
  • cross-region network RTT can consume a significant portion of the latency budget, so requests should normally be routed to the nearest healthy region
Headroom10 msReserved for garbage collection pauses and hot-spot retry backoff

Architecture Diagram

In the room: state the latency budget before picking Elasticsearch, because a sub-50ms backend p99 requirement necessitates an in-memory trie with precomputed top-10 completions at each node.

The architecture separates into two decoupled pathways: an online read path where keystrokes route through a prefix router to an in-memory trie shard and trending overlay returning the top 10 completions by popularity, and an offline update path where search logs flow through Kafka, Flink streaming aggregation, and offline batch rebuilds to publish immutable trie blobs deployed via S3. These paths are intentionally decoupled because slightly stale suggestions are acceptable whereas slow lookups violate the <50ms p99 latency contract.

Each trie node caches its precomputed top-10 completions so lookup operates in O(prefix length) time rather than performing a subtree scan. Shards partition by prefix range across 6 canonical shard tiers (~50 GB each, replicated 3x for ~900 GB aggregate serving RAM). A prefix-aware trending trie overlay bridges the gap between hourly trie rebuilds and breaking-news velocity.

Loading...

Component Deep Dives

Index structure and ranking are the two decisions that matter most here. Start with the trie and top-K precomputation, then walk the offline pipeline that keeps popularity scores fresh without blocking reads.

Trie Data Structure: The Core Algorithm

A trie node stores a character, a count of how many times the full path forms a top query, and child pointers. For the prefix "sy", the path traverses root to s, then to y, exposing illustrative completions (truncated to 3 for illustration; concrete production configuration stores top-10):

  • "system": count 50,000
  • "system design": count 20,000
  • "system architecture": count 5,000

Each node caches its precomputed top-10 completions (the general algorithm supports top-K; top-10 is the canonical production setting) so a prefix lookup runs in O(prefix length) time rather than executing a subtree traversal. During offline index construction, frequency scores are propagated upward and the top-10 lists are computed for each node. The resulting trie is published as an immutable snapshot, ensuring serving nodes query the trie with zero lock contention and no in-place mutation overhead.

Root
└── s
    └── y
        └── s
            └── t
                └── e
                    └── m (★ "system" count=50000)
                        └── d → e → s → i → g → n (★ "system design" count=20000)
                        └── a → r → c → h (★ "system architecture" count=5000)

Optimization 1: Store Top-10 results at each node (top-K generalized). Precompute and cache the top 10 most popular completions directly at each trie node. When a user types "sys", the service traverses from node 's' to 'y' to 's' and immediately returns the cached top-10 list. This avoids traversing the entire subtree at query time and guarantees an O(prefix_length) lookup.

YAML
node: "sys"
top_10:
  - query: "system design"
    count: 20000
  - query: "system architecture"
    count: 5000
  - query: "system requirements"
    count: 3000

Optimization 2: Compressed Trie (Radix Tree / Patricia Trie). Merge single-child node chains so the sequence from s through y, s, t, e, and m compresses into a single "system" node. This structural compression reduces memory consumption by 50% to 70%.

Optimization 3: Trie Serialization. Serialize the trie into a flat byte array for compact storage and rapid deserialization, utilizing memory-mapped files to achieve instantaneous service startup.

Offline Data Pipeline

The offline data pipeline transforms raw search events into updated, ranked immutable trie snapshots without placing load on the latency-sensitive read path.

  • Search Logs: Search logs record every query along with its timestamp, anonymized user identifier, and geographic location.
  • Kafka and Data Lake Archive: Kafka buffers incoming search events with a 7-day retention window primarily for streaming replay and pipeline recovery. Concurrently, a sink connector archives sampled search events to durable object storage (data lake) for long-term historical batch processing.
  • Flink Streaming Aggregation and Spark Batch Processing: Stream and batch workloads are clearly separated. Apache Flink continuously consumes live Kafka search events across 1-hour sliding windows (sliding every 5 minutes) to evaluate query velocity against a 24-hour baseline, outputting trending terms into Redis prefix buckets. Apache Spark reads historical logs from durable object storage to compute decayed popularity scores across daily and weekly windows, score = Σ (count_i × decay^(age_in_hours)), filtering noise (< 5 occurrences), profanity, and PII. Spark outputs aggregated query-score pairs to the Aggregated Query database.
  • Trie Builder: The builder reads aggregated query-score pairs from the Aggregated Query database, constructs the trie in memory, propagates scores upward, and computes the top-10 descendants at each node. It serializes the structure into an immutable binary blob and uploads it to S3. Rebuild Schedule: The core stable trie is rebuilt and published hourly via offline batch processing, while the trending overlay updates continuously on a 5-minute cadence from the streaming aggregation pipeline.

Autocomplete Service

The autocomplete serving fleet handles high-throughput prefix queries with strict sub-50ms latency guarantees.

  • Stateless Instances: Each autocomplete server operates statelessly, loading the immutable trie snapshot into local memory upon startup.
  • Trie Refresh: When a new hourly trie blob is published to S3, instances poll or receive an invalidation signal and atomically swap active pointers between double-buffered memory slots with zero downtime.
  • Prefix Lookup: Invoking trie.search(prefix) on the main immutable trie and trending_trie.search(prefix) on the local trending overlay returns candidate completions in O(prefix_length) time, which are merged, boosted, and deduplicated into the final top 10 suggestions.
  • Horizontal Scaling: The canonical baseline organizes 6 shard tiers of ~50 GB each with 3 replicas (~900 GB aggregate serving RAM). For extreme scale, the cluster can expand into 256 to 4,096 fine-grained prefix shards.

CDN Caching

Edge caching is what makes global delivery feasible.

  • Edge Caching for Hot Prefixes: High-frequency prefixes such as "how to", "what is", and "why" are cached at edge points of presence.
  • Cache Duration and Freshness Interaction: CDN caching is used primarily for stable, global autocomplete responses with a 1-hour time-to-live because popularity distributions for common short prefixes remain stable throughout the day. Responses that depend on the 5-minute trending overlay or personalization use a short TTL (≤5 minutes) or bypass the CDN entirely so that edge caching does not violate the 5-minute trending freshness guarantee.
  • Cache Key Structure: Responses are stored under keys formatted as autocomplete:{prefix}.
  • Offload Impact: Edge caching absorbs 30% to 50% of total autocomplete query traffic before reaching origin shards.

Client-Side Optimizations

Client applications implement caching and prefetching to eliminate redundant network traffic and provide instantaneous perceived response times.

  • Debouncing: Clients wait for 100 to 200 ms of inactivity before dispatching an HTTP request, cutting keystroke query volume by approximately 70%.
  • Local Browser Cache: The client caches previous responses so typing from "syst" to "syste" can immediately reuse results client-side without an extra network call.
  • Predictive Prefetching: Upon rendering results for "sys", the client asynchronously prefetches likely next prefixes such as "syst" and "sysa" during idle cycles.

Trending Overlay (Flink + Redis + In-Memory Trending Trie)

Pre-computed rankings live here so reads never sort at query time.

  • Near-Real-Time Overlay: Because hourly trie rebuilds cannot capture breaking news spikes, a lightweight trending overlay is refreshed every 5 minutes.
  • Velocity Detection: Apache Flink consumes search-events and evaluates 1-hour tumbling window frequencies against a 24-hour historical baseline.
  • Prefix-Aware Redis Storage: A standard flat Redis sorted set cannot efficiently answer prefix queries such as matching "sys". Redis stores trending queries partitioned into prefix buckets (such as trending:prefix:{prefix_2char}) alongside a global top-trending set with a 1-hour expiration.
  • Local Trending Trie: Each autocomplete serving node maintains a lightweight in-memory Trending Trie synchronized from Redis prefix buckets every 5 minutes, allowing sub-millisecond prefix searches on trending terms.
  • Result Blending and Deduplication: During query execution, the service queries both main_trie.search(prefix) and trending_trie.search(prefix). It merges the two result sets, boosts trending scores according to freshness weights, deduplicates identical terms, and returns the top 10 results within the 15ms lookup budget.

Event Bus Design (Kafka)

The event bus decouples producers from consumers and buffers traffic spikes.

YAML
topic: search-events
partitions: 64
partition_key: normalized_query # preserves per-query ordering for frequency aggregation
retention: 7 days # streaming replay and recovery; long-term data archived to object storage
replication_factor: 3
min_insync_replicas: 2

producer:
  collector: search-log-collector
  sampling_rate: 1-in-10 sampled at API gateway
  schema:
    event_id: string
    query: string
    user_id_hash: string
    timestamp: int64
    country: string
    language: string
    result_count: int32

consumer_groups:
  flink_aggregator:
    window: 1-hour sliding window with 5-minute slide
    sink: Aggregated Query DB (frequency deltas)
  trending_detector:
    metric: velocity compared against 24-hour baseline
    sink: Redis prefix buckets backing in-memory Trending Trie overlay
  trie_builder:
    schedule: hourly batch job (Apache Spark / offline builder)
    action: reads aggregates, rebuilds immutable trie binary blob, and uploads to S3

serving_path:
  loader: Autocomplete Service loads immutable trie snapshot from S3
  merge_strategy: prefix query on main trie union prefix query on trending trie (top-10 combined)

dead_letter_queue:
  topic: search-events-dlq
  max_retries: 3
  alert_condition: Flink consumer lag exceeds 30 minutes (warning alert; lag >2 hours triggers severe incident escalation)

API Design

A single GET with query prefix and limit is enough for the core API. Show the response shape.

Autocomplete Suggestions

TYPESCRIPT
interface AutocompleteRequest {
  prefix: string;      // input characters typed by user
  limit?: number;      // default: 10 suggestions
  language?: string;   // default: "en"
  country?: string;    // ISO country code for regional filtering
}

interface SuggestionItem {
  query: string;       // suggested query completion
  score: number;       // popularity score derived from historical frequency
}

interface AutocompleteResponse {
  suggestions: SuggestionItem[];
}

// Retrieves the top-K query suggestions matching the requested prefix
function getAutocompleteSuggestions(params: AutocompleteRequest): Promise<AutocompleteResponse>;
HTTP
GET /api/v1/autocomplete?q=system+des&limit=10&lang=en HTTP/1.1
Host: api.example.com

HTTP/1.1 200 OK
Content-Type: application/json

{
  "suggestions": [
    { "query": "system design", "score": 20000 },
    { "query": "system design interview", "score": 15000 },
    { "query": "system design primer", "score": 8000 },
    { "query": "system design patterns", "score": 5000 },
    { "query": "system design basics", "score": 3000 }
  ]
}

Log Search (Internal Ingestion)

TYPESCRIPT
interface SearchLogPayload {
  query: string;       // executed search query string
  userId: string;      // anonymized user identifier hash
  timestamp: string;   // ISO 8601 creation timestamp
  location: string;    // ISO country code
  language?: string;   // language locale
  resultCount?: number;// count of matching search results
}

// Ingests an anonymized search query event for offline popularity aggregation
function ingestSearchLog(payload: SearchLogPayload): Promise<{ status: "queued" }>;
HTTP
POST /api/v1/search-log HTTP/1.1
Host: internal-api.example.com
Content-Type: application/json

{
  "query": "system design",
  "user_id": "anon-hash",
  "timestamp": "2026-03-13T10:00:00Z",
  "location": "US"
}

Common Error Responses

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
504 Gateway Timeout: search index shard responded slowly, narrow query parameters or retry

Data Model

Trie nodes and aggregated database records separate real-time memory lookups from asynchronous updates based on write frequency.

In-Memory Trie Node

TYPESCRIPT
interface SuggestionEntry {
  query: string;       // completion query string
  score: number;       // precomputed popularity score
}

class TrieNode {
  children: Map<string, TrieNode>; // character mapped to child node, or fixed array[26] for lowercase ASCII
  isEnd: boolean;                  // indicates whether this node completes a valid query
  topK: SuggestionEntry[];         // pre-cached top-K suggestions sorted by popularity score
}

Aggregated Query DB (Cassandra / ScyllaDB)

SQL
-- Stores aggregated query frequencies with time decay
CREATE TABLE query_frequency_aggregates (
  query_string text PRIMARY KEY,
  hourly_count bigint,
  daily_count bigint,
  weekly_count bigint,
  score double,            -- weighted popularity score with exponential decay
  last_updated timestamp
);

Kafka Topic: search-events

JSON
{
  "query": "system design",
  "user_id": "anon-hash-uuid",
  "timestamp": "2026-03-13T10:00:00Z",
  "country": "US",
  "language": "en",
  "result_count": 1250000
}

Trie Blob (Serialized Binary Format)

YAML
header:
  version: uint32
  node_count: uint32
  total_queries: uint64
  built_at: int64

nodes:
  - char: uint8
    is_end: boolean
    child_count: uint16
    child_offsets: uint32[]
    top_k:
      - string_offset: uint32
        score: float32

string_table:
  encoding: utf-8
  format: contiguous null-terminated strings referenced by byte offset
  • Stored in S3 with sizes of ~50 GB per shard tier (~300 GB logical aggregate across the 6-tier baseline).
  • Memory-mapped directly from local disk for near-instantaneous startup times.

Fault Tolerance

Index rebuild lag, cache stampede on trending queries, and stale popularity scores.

ConcernSolution
Trie server crashNodes are stateless so load balancers route traffic to healthy instances, while any newly provisioned instance downloads the latest trie snapshot from S3.
Trie build failureServing instances continue serving the active trie snapshot while alerting on-call engineers, falling back to the last known healthy trie artifact.
Stale suggestionsAcceptable during transient delays because users continue receiving relevant suggestions from the hourly core trie snapshot while recent spikes are covered by the 5-minute trending overlay.
S3 unavailableThe trie snapshot is cached locally in memory and on local disk across all serving instances, allowing lookups to survive S3 outages.
Data pipeline lag
  • Consumer lag exceeding 30 minutes triggers an early operational warning alert, while lag exceeding 2 hours triggers severe incident escalation
  • during pipeline delays, the system safely falls back to suggestions from the active core trie snapshot.

Specific: Trie Hot-Swap Without Downtime

  1. A new trie binary blob is compiled offline and uploaded to S3.
  2. Each serving instance maintains two memory slots (Slot A and Slot B) using double buffering.
  3. The server loads the newly downloaded trie into the inactive slot (B) while continuing to serve active queries from Slot A.
  4. The server executes an atomic pointer swap redirecting active traffic to Slot B, allowing the old trie in Slot A to be garbage collected.
  5. The transition achieves zero downtime with no latency degradation during the refresh cycle.

Additional Considerations

Personalization, typo tolerance, and trending query injection are stretch topics.

Handling Trending / Breaking News

  • Hourly batch rebuilds cannot capture breaking news queries in real time.
  • Trending Overlay: Maintain a lightweight in-memory trending trie overlay on serving nodes, refreshed every 5 minutes from Flink streaming window outputs backed by Redis prefix buckets.
  • Query-Time Merge: Autocomplete nodes query both the primary immutable trie and the trending trie by prefix, merge candidate suggestions, apply a trending score boost, and deduplicate results.
  • Window Aggregation: The trending index evaluates query velocity over the preceding 1-hour window compared against a 24-hour baseline to isolate sudden spikes from normal background volume.

Personalized Suggestions

  • Global popularity suggestions are augmented with the user's individual search history.
  • The client application stores the user's last 100 queries in local device storage and searches them locally before issuing remote requests.
  • The user interface blends 2 to 3 personalized matches with 7 to 8 globally ranked suggestions.
  • Privacy Preservation: Performing personalized matching client-side avoids transmitting and storing sensitive user search histories on backend servers.

Multi-Language Support

  • Maintain dedicated trie indices per supported language and filter by client country/locale headers (sophisticated multi-language morphology and stemming are out of scope for this design).
  • Languages without whitespace boundaries, including Chinese, Japanese, and Korean, use n-gram indexing structures rather than standard character tries to handle ideographic segmentation.

Handling Offensive / Sensitive Queries

  • A curated blocklist filters offensive, hateful, or explicit suggestions during index generation and serving.
  • Automated classifiers prevent suggestions containing personally identifiable information such as email addresses, phone numbers, and Social Security numbers.
  • Administrative APIs support real-time removal of specific query strings to satisfy legal mandates and DMCA compliance requests.

Sampling for Scale

  • High-volume search engines avoid recording every individual query event into the aggregation pipeline.
  • Canonical Baseline: Apply deterministic or probabilistic 1-in-10 request sampling at the API gateway (configurable from 1-in-10 to 1-in-100 based on traffic volume and accuracy targets). This operates as fixed-rate stream event sampling rather than reservoir sampling.
  • Downstream aggregators scale recorded frequencies by 10x (the inverse sampling rate) to derive accurate global popularity estimates.
  • Sampling introduces estimation error, so the sampling rate is tuned against ranking-quality metrics and can be increased for low-frequency or high-accuracy workloads.
  • This sampling strategy cuts data pipeline storage and compute overhead by 10x (and up to 100x when configured for extreme traffic reduction).

Spell Correction Integration

  • When an exact-prefix lookup yields zero suggestions, the service invokes a synchronous but tightly bounded spell-correction fallback.
  • The service queries a precomputed 1-edit-distance index (such as SymSpell with pre-indexed deletion neighborhoods) with a strict 5 to 10 ms internal timeout budget.
  • If correction completes within budget, matching suggestions for the corrected prefix are returned; otherwise, the request gracefully degrades and returns an empty list rather than violating the <50ms p99 autocomplete latency SLO.
  • For example, an input of "systm desi" corrects to "system desi" within 8ms to surface suggestions for "system design".

Related Problems

Autocomplete suggestions feed the primary query input described in Search Engine, whereas this article covers prefix top-K retrieval rather than full document search ranking. In addition, real-time frequency estimation and stream ranking overlap with patterns discussed in Top-K Rankings.

Interview Walkthrough

  • 25-minute cut

    Skip arch50/arch75 depth unless staff.

    • Latency budget (<50ms backend p99) and trie vs database tradeoffs (3 min)
    • Trie structure with precomputed top-10 candidates per node (6 min)
    • Prefix sharding (canonical 6-shard baseline vs fine-grained 4,096 shards) (6 min)
    • Offline pipeline: search logs through Flink streaming and Spark batch to rebuilds (5 min)
    • CDN edge cache for common prefixes (5 min)
  • State the latency target first with a backend p99 threshold strictly under 50ms, which rules out database prefix scans and necessitates an in-memory trie.
  • Explain the trie layout where each node caches precomputed top-10 completions for its prefix ranked by decayed popularity score.
  • Partition the trie across servers using the canonical 6-shard baseline of ~50 GB each with 3 replicas (~900 GB aggregate serving RAM), and provision additional replicas for hot shards to distribute lookup volume.
  • Maintain suggestions asynchronously through a decoupled log pipeline: hourly batch snapshot rebuilds for the core index alongside a 5-minute prefix-aware trending trie overlay.
  • Introduce CDN edge caching for the most frequent prefix queries such as "a", "th", and "the" to absorb global traffic spikes.
  • Incorporate personalization as a secondary signal blended with global popularity scores, retaining the standard trie as a fallback for cold users.
  • Address common pitfalls such as querying SQL databases using LIKE 'prefix%' on each keystroke, because full table scans fail sub-50ms latency requirements under load.

Engineering Trade-offs

Evaluating an in-memory trie against Elasticsearch or compressed prefix trees depends primarily on dictionary scale, the required <50ms p99 latency SLA, and frequency update cadence.

Flink Streaming Aggregation for Trending Detection

Problem: Hourly trie rebuilds are too slow for breaking news events. For instance, if an unexpected event occurs at 2:15 PM, related search suggestions must surface by 2:20 PM via the 5-minute trending overlay.

Delivery and Replay Semantics: Flink uses checkpointed state (every 30 seconds) and replayable Kafka offsets. The downstream aggregation sinks in the Aggregated Query DB and Redis prefix buckets utilize idempotent, versioned writes (such as upserts keyed by query and window timestamp with monotonic sequence numbers), so that replaying Kafka offsets after a recovery does not double-count event totals. Where Kafka transactional producers are used, they guarantee exactly-once message delivery between Kafka topics, while sink idempotency ensures consistent downstream state.

Loading...

Decay Function Scoring: Worked Example

Recent searches carry greater weight than historical queries to reflect real-time interest. The ranking pipeline applies exponential decay using the formula score = Σ (count_i × decay^(age_in_hours)), configured with a default decay factor of 0.97 so queries lose approximately 3% of their weight each hour. This approach provides smooth degradation compared to linear decay models that introduce abrupt cutoffs at fixed intervals.

The decay rate is tuned per domain. Setting decay to 0.99 provides slow decay suitable for stable corpora, 0.97 acts as the standard moderate baseline, and 0.90 enables rapid decay for fast-moving verticals. For instance, breaking news queries use a fast decay factor of 0.90 to surface emerging stories quickly, whereas e-commerce product catalogs utilize a slow decay factor of 0.99 to preserve steady search demand.

Trie Sharding for Very Large Query Corpuses

Approach 1: Sharding by Initial Characters. Partitioning the trie based on the first two or three characters simplifies routing logic, but vocabulary skew leads to uneven shard sizes and prevents efficient cross-shard suggestion aggregation unless hot shards are split or replicated.

Approach 2: Full Trie Replicated on Every Server (Alternative for Aggressively Pruned Corpora). In this alternative configuration, each serving instance loads the entire trie into memory, enabling any node to answer any autocomplete prefix without a prefix router. For an aggressively pruned ~50 GB corpus, allocating 50 GB across 20 servers requires 1 TB of aggregate RAM, offering high operational simplicity. For our canonical ~300 GB baseline, however, replicating the full trie on every node is memory-prohibitive, which is why the 6-shard partitioned fleet (~50 GB per shard tier x 3 replicas = ~900 GB RAM) is the canonical production choice.

Approach 3: Two-Level Hierarchical Trie. A global Level 1 cache replicated across all servers hosts the top 100,000 queries within 1 GB of RAM to absorb more than 80% of total traffic. Meanwhile, a sharded Level 2 tier holds the long-tail dictionary. This hybrid model maximizes throughput for common keystrokes while accepting slightly higher latency for rare queries.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...