System Design Problem

Design a Search Engine (Google)

Commonly Asked By:GoogleMicrosoftAmazon

Interview Setup

Interview Prompt

Design a web search engine like Google. Crawl and index billions of web pages. Given a search query, return the top 10 most relevant results with snippets in under 200ms.

Clarifying Questions (ask before designing)

QuestionWhy it matters
How many documents and what's the query QPS?A scale of 1B documents at 10K QPS drives the sharding strategy, whereas 100M documents can fit on fewer shards.
Full web crawl or search over a fixed corpus?A full web crawl requires crawler infrastructure, index freshness pipelines, and incremental updates, so we assume crawled documents already exist.
Simple keyword search or also phrase/boolean operators?Query parser complexity determines index lookups, because boolean queries like 'apple AND juice' require distinct execution compared to literal phrase queries like 'apple juice'.
Personalized ranking or global relevance only?Personalization requires user profile lookups during the ranking stage, which remains separate from the core retrieval stage.

Scope

In scope

  • Inverted index design and sharding
  • TF-IDF / BM25 scoring
  • Query parsing and execution
  • Ranking pipeline (retrieval + rerank)
  • Snippet generation
  • Capacity estimation

Out of scope (state explicitly)

  • Web crawler internals (separate problem)
  • Ad auction and sponsored results
  • Full PageRank iterative computation
  • Image/video search

Functional Requirements

Start by asking your interviewer whether you are designing a full web crawl engine or searching over a fixed corpus. Confirm indexing constraints, ranking depth between BM25 baseline and machine learning reranking, and whether phrase or boolean queries are in scope before enumerating system features.

  • Accept a text query and return a ranked list of relevant web pages.
  • Support keyword matching, phrase matching, and boolean queries.
  • Rank results by relevance combining PageRank, text relevance, freshness, and personalization signals.
  • Display result snippets including title, target URL, and highlighted keyword excerpts.
  • Support autocomplete suggestions and prefix searches as detailed in Typeahead Autocomplete.
  • Support vertical search categories including image, video, and news search.
  • Spell correction, such as transforming "systm desgn" into "Did you mean: system design?".
  • Extract knowledge panels for known entities including people, places, and organizations.
  • Support pagination of ranked search results using shallow page numbers for standard UI browsing and opaque continuation tokens for deep result traversal.

Non-Functional Requirements

Interviewers focus heavily on query latency and index freshness for this problem. The 200ms p99 budget forces sharded inverted indexes and aggressive result caching, so naming that tension early clarifies your architectural choices.

  • Low Latency: Return search results within a strict < 200 ms p99 external SLO, targeting an internal execution budget of ~110 ms with ~90 ms reserved for network jitter, queueing, and tail-latency hedging.
  • High Availability: 99.99% query serving availability.
  • Scalability: 100B indexed pages across tiered storage, where the hot tier serves 1B documents against ~86M queries per day (~1K average QPS, ~10K peak QPS).
  • Freshness: New and updated pages indexed within minutes for news content and within hours for standard web pages.
  • Relevance: Search results must provide high precision and recall, prioritizing quality and authority.
  • Spam Resistance: Robust defense against SEO manipulation, link farms, and web spam.

Capacity Estimations

Establish capacity calculations before selecting a sharding strategy. Corpus size and query QPS determine the number of index servers required, distinguishing between logical index capacity (~1,000 servers at 1 TB each for a 1 PB inverted index) and 3x replicated physical capacity (~3,000 server-equivalents plus operational headroom), while posting list compression dictates whether a petabyte-scale index is feasible.

MetricCalculationValue
Indexed pages (total web)Given tiered index baseline100B
Hot tier (always searched)Given interview baseline (Tier 0)1B
Avg page size (compressed)Given (typical workload assumption)100 KB
Raw index size100B x 100 KB10 PB
Inverted index size (logical)~10% of raw1 PB
Queries / day10K peak x 86400 x ~0.1 duty cycle~86M
Queries / sec (avg)86M ÷ 86400~1K
Queries / sec (peak)Given~10K
Logical index capacity1 PB ÷ 1 TB per server~1,000 servers (1 TB each)
Physical index servers (3x replica)~1,000 logical x 3 replicas (plus headroom)~3,000 server-equivalents
Crawl rateGiven (assumption documented in value)1B pages/day

Architecture Diagram

In an interview setting, confirm whether personalization and ML reranking are in scope, because they add a feature store lookup directly onto the latency-critical hot path.

Walk your interviewer through the system by lifecycle stage. The design splits crawl, indexing, and query serving because each subsystem operates under distinct latency and throughput constraints. Crawled pages land in a distributed document store, while background index builders merge posting lists into sharded inverted indexes. On the query path, popular search queries short-circuit the inverted index through an in-memory query result cache. For cache misses, a stateless coordinator parses the query, broadcasts to Tier-0 document shards in parallel, cascades to Tier 1 and Tier 2 only if candidate count or relevance is insufficient, gathers and merges local candidates, reranks the top 1,000 documents with machine learning models, and generates contextual snippets for the final top 10 results.

Loading...

Component Deep Dives

Walk through each component in the architecture systematically, starting with query understanding on the read path, continuing through index storage structures, and concluding with the offline ingestion pipeline.

Query Service (Query Understanding)

Every search begins here, where tokenization and intent detection run before touching the inverted index.

  • Query parsing: Tokenize input text, convert to lowercase, eliminate stopwords, and stem or lemmatize query terms.
  • Spell correction: Calculate edit distance combined with language model probability, identifying terms like "systm" and proposing the highest-confidence correction.
  • Query expansion: Broaden query context by adding acronym expansions, such as querying "New York City restaurants" when a user searches for "NYC restaurants".
  • Intent detection: Classify whether the query represents navigational intent ("facebook login"), informational intent ("how to cook pasta"), or transactional intent ("buy iPhone 15").
  • Synonym handling: Inject semantic equivalents into the query AST, searching for "automobile" and "vehicle" when users query "car".

Inverted Index: Core Data Structure

The inverted index serves as the foundational data structure for full-text search retrieval.

An inverted index maps terms to a posting list of documents containing that term along with occurrence positions and frequencies. While default boolean AND intersection is used as a straightforward interview model to demonstrate posting list intersection, production search engines employ query-dependent candidate retrieval using minimum-should-match constraints, Block-Max WAND algorithms, phrase position checks, and query expansion:

"system"     -> [doc1:3, doc5:1, doc15:7, doc22:2, ...]   (doc_id:term_frequency)
"design"     -> [doc1:2, doc3:5, doc15:4, ...]
"interview"  -> [doc1:1, doc15:3, doc42:2, ...]

Index Structure per Term (Posting List Record):

JSON
{
  "term": "system",
  "document_frequency": 50000000,
  "posting_list": [
    { "doc_id": 1, "tf": 3, "positions": [15, 42, 78] },
    { "doc_id": 5, "tf": 1, "positions": [201] }
  ]
}

Index Sharding: While term-based sharding partitions data by assigning a subset of terms to each shard, document-based sharding is strongly recommended and represents the industry standard. Under document-based sharding, each shard holds the complete inverted index for a dedicated slice of documents. The query coordinator broadcasts each query across shards within the active index tier simultaneously, each shard computes its local top-K candidates, and the coordinator merges these candidates into a unified global ranking. To mitigate hot-term traversal overhead across document shards, engines layer an optimization: maintaining specialized, pre-warmed in-memory hot-term posting structures or caching posting-list fragments on query coordinators to accelerate viral query lookups.

Index Compression: Posting lists are compressed using Variable-Byte Encoding or PForDelta algorithms. Document IDs are delta-encoded so that monotonically increasing IDs such as [1, 5, 15, 22] are stored as compact gaps [1, 4, 10, 7], because smaller integer deltas compress significantly better. This optimization reduces overall index storage footprint by 5x to 10x.

Ranking Service: The Scoring Pipeline

Ranking represents the deepest stage of search interviews, governed by a multi-stage pipeline combining cheap baseline retrieval with expensive machine learning reranking over candidate subsets.

Stage 1: Initial Retrieval (Coarse Filter): Executes candidate generation across inverted index shards. In an interview setting, default boolean AND matching demonstrates posting list mechanics by intersecting doc_id lists to find documents containing all query terms. In production systems, retrieval balances recall and precision using query-dependent algorithms: OR-style candidate retrieval with minimum-should-match thresholds, Block-Max WAND (Weak AND) algorithms to prune low-scoring postings early, phrase position constraints, and query expansion tokens. The candidate retrieval phase caps the candidate set at 10,000 documents for Stage 2 scoring.

Stage 2: Scoring (BM25 Text Relevance): Evaluates candidate documents using BM25 (Best Matching 25), the industry standard text relevance function that balances term frequency saturation, document length normalization, and inverse document frequency, pruning candidates down to the top 1,000 documents.

Stage 3: PageRank (Link Graph Analysis): Quantifies page authority based on the web graph of incoming hyperlinks. PageRank is computed offline using iterative MapReduce running 20 to 50 iterations until convergence, and scores are stored as precomputed features in the document index. While an illustrative linear formula like final_score = αxBM25_norm + βxPageRank_norm demonstrates the intuition of combining textual relevance with authority, raw BM25 scores (unbounded) and raw PageRank values (power-law distribution) are not directly comparable without non-linear transformation and feature normalization. In production systems, PageRank is ingested as a normalized numerical feature alongside BM25 and freshness into the Stage 4 Machine Learning ranking stack. Scores recalculate daily, while high-velocity or trending pages receive incremental updates from the crawl pipeline as detailed in Web Crawler.

Stage 4: Machine Learning Reranking (Fine-Grained Scoring): Evaluates the top 1,000 candidate documents from Stage 1/2 using a Learning-to-Rank model such as LambdaMART across 200+ ranking signals (BM25 score, PageRank, historical CTR, freshness, domain authority, user location). A second, more computationally intensive neural sub-stage (such as BERT cross-encoders) can optionally score the top 100 to 200 candidates to select the final top 10 results for user display.

Snippet Generator

Snippets are computed exclusively for the final top 10 documents, keeping heavy extraction logic off the retrieval hot path.

  • Generate a relevant summary snippet for each result showing matched query terms in context.
  • Locate the passage within the document body exhibiting the highest query term density and proximity.
  • Wrap matching keywords in bold formatting tags for frontend rendering.
  • Truncate the resulting passage cleanly to approximately 160 characters.

Index Builder (Offline / Incremental)

Index freshness relies on incremental update streams, allowing new or modified documents to become searchable without triggering an expensive full rebuild.

  • Consume index-events from Kafka with at-least-once processing semantics: worker nodes consume event batches, apply index mutations to in-memory segments, and commit partition offsets only after successful persistence.
  • Parser: Tokenize, stem, and construct posting lists partitioned across document shards. Indexers enforce idempotency and reject stale events using monotonic document versions: if incoming.doc_version <= indexed.doc_version, the event is safely ignored.
  • Segment Compaction and Full Rebuild: Background segment compaction continuously merges supplement segments into the primary index to control segment count and query overhead. A separate distributed full rebuild runs weekly for major maintenance, tombstone cleanup, and index regeneration.

Event Bus Design (Kafka)

The event bus decouples crawler ingestion from index generation workers while buffering sudden traffic spikes.

YAML
topic: index-events
partitions: 256
partition_key: doc_id # maintains document update ordering on a single partition
retention: 14 days # enables event replay for complete index rebuilds
replication_factor: 3
min_insync_replicas: 2

producer: crawler and indexer workers (idempotent producer)
event_payload:
  event_id: string
  doc_id: string
  doc_version: int64 # monotonically increasing crawl sequence or generation
  url: string
  action: "add | update | delete"
  content_hash: string
  crawled_at: timestamp

consumer_groups:
  index_builder: merge into inverted index segments in near-real-time
  snippet_cache: precompute static snippets for newly indexed documents
  pagerank_input: stream link graph updates to offline PageRank jobs

guarantees:
  query_path: read-only architecture where index-events feed offline pipelines exclusively
  consumption_semantics: at-least-once processing (consume event -> apply index mutation -> commit Kafka offset)
  idempotency_and_replay: indexer workers are strictly idempotent; crashed nodes replay events without corrupting index state
  version_gating: indexers reject stale events if incoming.doc_version <= indexed.doc_version
  dlq: index-events-dlq after 3 retries with alerting when consumer lag exceeds 1 hour

API Design

Define search query contracts and domain models before exploring internal posting list structures. The API demonstrates standard shallow page-number pagination (?page=1, totalPages) suitable for client-facing result browsing, alongside an opaque continuation token (cursor / nextCursor) for stable deep pagination at scale without expensive deep-offset database or heap traversals. Pagination, spell correction, and knowledge cards follow the core query path.

Search Query Types and Client Contract

TYPESCRIPT
// Search service client request and response domain interfaces
interface SearchRequest {
  query: string;
  page?: number; // illustrative shallow pagination for standard UI results
  cursor?: string; // opaque continuation token for stable deep pagination at scale
  pageSize?: number;
  language?: string;
  country?: string;
}

interface SearchResultItem {
  title: string;
  url: string;
  snippet: string;
  faviconUrl?: string;
  cachedUrl?: string;
  rank: number;
}

interface SearchResponse {
  query: string;
  spellCorrection?: string | null;
  resultsCount: number;
  results: SearchResultItem[];
  relatedSearches: string[];
  knowledgePanel?: Record<string, unknown> | null;
  pagination: {
    page: number;
    totalPages: number;
    nextCursor?: string | null; // serialized continuation token for deep pagination
  };
}

// Internal query coordinator interface
interface SearchCoordinatorService {
  search(request: SearchRequest): Promise<SearchResponse>;
}

Execute Search Query Endpoint

HTTP
GET /api/v1/search?q=system+design+interview&page=1&lang=en&country=US
HTTP/1.1 200 OK
Content-Type: application/json

{
  "query": "system design interview",
  "spell_correction": null,
  "results_count": 125000000,
  "results": [
    {
      "title": "System Design Interview Guide - ByteByteGo",
      "url": "https://bytebytego.com/system-design",
      "snippet": "A comprehensive guide to <b>system design interview</b> preparation...",
      "favicon": "https://bytebytego.com/favicon.ico",
      "cached_url": "https://cache.example.com/system-design",
      "rank": 1
    }
  ],
  "related_searches": ["system design interview questions", "system design primer"],
  "knowledge_panel": null,
  "pagination": {
    "page": 1,
    "total_pages": 100,
    "next_cursor": "eyJxdWVyeSI6InN5c3RlbSBkZXNpZ24iLCJ0aWVyIjowLCJvZmZzZXQiOjEwfQ=="
  }
}

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

The data model follows read-heavy access patterns: posting lists for rapid term lookup, a document store for full text and metadata, and precomputed PageRank scores merged during query ranking.

Inverted Index (SSTable-like format on disk)

YAML
# Inverted Index SSTable-like format on disk
term_dictionary_in_memory:
  "system":
    offset: "0x4A2F"
    doc_freq: 50000000
  "design":
    offset: "0x8B1C"
    doc_freq: 35000000

posting_list_on_disk:
  offset_0x4A2F:
    - doc_id_delta: 1
      tf: 3
      positions: [15, 42, 78]
    - doc_id_delta: 4
      tf: 1
      positions: [201]

Document Store (Bigtable / GFS)

YAML
# Document Store Schema (Bigtable / GFS)
row_key: "doc_id (sha256 hash of URL)"
column_families:
  content:
    title: string
    body_text: string
    meta_description: string
    language: string
  metadata:
    url: string
    domain: string
    crawl_date: timestamp
    content_hash: string
    robots_directives: string
  links:
    outgoing_urls: string[]
    incoming_count: int64
  scores:
    pagerank: float64
    spam_score: float32
    domain_authority: float32

URL Frontier (for Crawler)

YAML
# URL Frontier Priority Queue Entry
priority_queue_entry:
  url: string
  priority: int32
  last_crawled: timestamp
  crawl_interval: int32
  domain: string

PageRank Store

YAML
# PageRank Store Record
key: "doc_id"
value:
  pagerank_score: 0.00042
  last_computed: timestamp

Kafka Topic: index-events

JSON
{
  "event_id": "uuid-v4",
  "doc_id": "hash-of-url",
  "doc_version": 1042,
  "url": "https://example.com/page",
  "action": "add",
  "content_hash": "sha256:...",
  "crawled_at": "2026-03-13T10:00:00Z"
}

Fault Tolerance

Index lag, partial shard failures, and hot query overload represent the most frequent production failure modes.

ConcernSolution
Index server failureEach index shard is replicated 3x across separate racks, allowing the coordinator to route queries to a healthy replica automatically.
Index corruptionChecksum verification detects damaged segments, enabling fast rebuilds from the underlying content store.
Query overloadCircuit breakers trigger graceful degradation, skipping heavy ML re-ranking to serve BM25 text relevance results directly.
Stale indexContinuous incremental index updates stream fresh pages, complemented by a full index rebuild scheduled weekly.
Crawler politenessRespect robots.txt directives and enforce per-domain rate limits as detailed in Web Crawler.

Index Serving Architecture: Tiered Search Execution

  • Tier 0 (Top 1B Pages): Holds the top 1B most authoritative and frequently accessed web documents, evaluated unconditionally in parallel on every search query. It operates as an independent document-sharded, 3x replicated server cluster.
  • Tier 1 (Extends Coverage to Top 10B Pages Total): Expands cumulative coverage to the top 10B documents total by adding approximately the next 9B pages beyond Tier 0, organized as an independent document-sharded cluster. It is queried only if Tier 0 returns an insufficient candidate count or relevance threshold.
  • Tier 2 (Full 100B-Page Corpus): Encompasses the full 100B-page corpus by storing the remaining ~90B long-tail web documents in cold storage across separate sharded clusters, queried exclusively when Tier 0 and Tier 1 cannot satisfy rare or highly specific queries.
  • Execution Cascade: Queries evaluate Tier 0 in parallel. If sufficient high-scoring candidates are found (answering the vast majority of traffic), execution proceeds immediately to ML reranking without querying Tier 1 or Tier 2, protecting both latency SLOs and cluster hardware.

Additional Considerations

Snippets, faceted search, and advertising placement are secondary topics explored once the core query retrieval loop is established.

Index Update Strategy

  • Full Rebuild: A distributed batch job rebuilds the entire index weekly to process tombstones, purge deleted documents, and refresh global authority scores.
  • Incremental Updates: A real-time ingestion pipeline streams newly discovered and updated pages into a fast supplement index.
  • Dual-Index Querying: The query coordinator searches both the primary index and the supplement index in parallel, merging their candidate sets before scoring.
  • Segment Compaction: Background workers periodically merge the supplement index into the primary index to control segment file sprawl.

Anti-Spam and Web Spam Detection

  • Link Spam: Graph algorithms detect artificial link farms that link to each other to manipulate PageRank scores.
  • Content Spam: Machine learning models flag keyword stuffing, hidden zero-pixel text, and domain cloaking patterns.
  • Click Spam: Anomaly detection filters artificial click-through rate inflation on specific search queries.
  • Authority Signals: Trust scores incorporate domain registration age, valid SSL certificates, content uniqueness, and organic link velocity.

Query Result Caching

  • Cache top ranked results for popular search queries in Redis or Memcached clusters.
  • Because roughly 25% of search queries repeat within an hour, caching yields high cache hit ratios and relieves load on index shards.
  • Cache Key Structure: Normalized query text hashed with location, language, and index generation tags.
  • Expiration Policy: A 1-hour TTL serves general queries, while trending news queries expire within 5 minutes.
  • Cache Freshness Trade-off: Result caching trades immediate index freshness for ultra-fast query latency and shard protection. When an incremental index update indexes a new document, that document will not appear in cached popular queries until their TTL expires or a cache generation bump occurs.
  • Cache Generation Invalidation: Ordinary staleness is bounded by TTL. For major index rebuilds, schema migrations, or emergency updates, query coordinators increment an index generation counter in the cache key (e.g. cache_gen = 42), instantly invalidating stale query entries cluster-wide while letting old keys evict naturally under LRU policy.

Semantic Search

  • Extends keyword matching by interpreting query semantic intent rather than literal character overlaps.
  • Transformer Embeddings: Models like BERT encode queries and web documents into dense vector representations.
  • Vector Similarity Search: Approximate Nearest Neighbor (ANN) search engines such as FAISS or ScaNN retrieve semantically adjacent documents.
  • Hybrid Scoring: Combines BM25 keyword matching with vector cosine similarity to produce a balanced relevance score.

Related Problems

Page discovery, link graph extraction, and crawl politeness are detailed in Web Crawler. Learning-to-Rank models, feature store ingestion, and machine learning scoring pipelines are explored in Search Ranking (LTR). Autocomplete prefix caching and query suggestion engines connect to Typeahead Autocomplete. For foundational principles, review Indexing and Query Optimization, Sharding and Partitioning, Caching Patterns and Invalidation, and Kafka Architecture and Guarantees.

Interview Walkthrough

  • 25-minute cut

    Skip PageRank derivation and ML rerank training unless interviewing for staff-level roles.

    • FR/NFR and scope (3 min)
    • Inverted index structure and posting lists (7 min)
    • Document sharding with scatter-gather (8 min)
    • BM25 intuition and two-stage ranking (7 min)
  • Sketch the three-stage pipeline first spanning crawl, index, and query so the interviewer sees you understand the full lifecycle rather than only the query box.
  • Center the deep dive on the inverted index, emphasizing how mapping terms to posting lists containing document IDs, positions, and frequencies forms the core data structure.
  • Shard the index by document ID using Sharding and Partitioning, scattering queries across shards in parallel while enforcing a strict latency budget.
  • Allocate your query budget with clear distinction between the external SLO (< 200 ms p99) and the internal execution target (~110 ms, budgeting 10 ms for query understanding, 5 ms for network fan-out, 25 ms for per-shard evaluation, 5 ms for heap gathering, 50 ms for ML reranking, and 10 ms for snippet extraction, leaving ~90 ms of headroom for tail-latency hedging, queueing, and network jitter).
  • Distinguish between logical index capacity (~1,000 servers at 1 TB each for 1 PB inverted index) and physical capacity (~3,000 server-equivalents under 3x replication plus operational headroom).
  • Apply Caching Patterns and Invalidation to popular queries, because roughly 25% of searches repeat within an hour and make result caching exceptionally high-impact.
  • Describe the two-stage ranking stack, employing BM25 for broad recall to gather the top 1,000 candidates, followed by machine learning models (with an optional neural cross-encoder on the top 100-200) for precision before snippet extraction on the top 10.
  • Highlight incremental indexing through near-real-time supplement segments that merge periodically, ensuring fresh content becomes searchable without full rebuilds.
  • Avoid common pitfalls such as proposing relational database scans with LIKE '%keyword%', because interviewers expect an inverted index with posting lists.

Engineering Trade-offs

Your interviewer will push on the major architectural forks here. Walk through document vs term sharding, then state why document sharding is the industry standard for a Google-scale corpus.

End-to-End Query Execution Flow with Timing Budget

To satisfy an external SLO of < 200 ms p99, execution is structured across parallel stages with an internal target budget of ~110 ms, leaving approximately 90 ms of headroom for network jitter, queueing, cache misses, serialization, speculative hedging, and tail latency:

  • Step 1: Query Understanding (10 ms): Performs tokenization, spell checking, query expansion such as expanding "NYC" into "New York City", and query intent classification.
  • Step 2: Scatter to Index Shards (5 ms network): The stateless Query Coordinator broadcasts the parsed query to 100 document-sharded Tier-0 index servers in parallel, enforcing a 150 ms deadline before skipping non-responsive nodes.
  • Step 3: Per-Shard Processing (25 ms): Each shard intersects compressed posting lists using two-pointer merges over sorted doc ID lists, computes initial BM25 scores, and selects local top 100 candidates into a min-heap.
  • Step 4: Gather and Merge (5 ms): The coordinator collects up to 10,000 local candidates from all shards and runs a K-way max-heap merge to extract the global top 1,000 documents.
  • Step 5: ML Reranking (50 ms): Evaluates the top 1,000 candidates across hundreds of signals including CTR, freshness, domain authority, and PageRank using LambdaMART, supplemented by BERT cross-encoders for the top 100 to 200 candidates to select the final top 10.
  • Step 6: Snippet Generation (10 ms): Extracts high-density matching passages from the document store for the final top 10 ranked documents, highlights query terms, and truncates snippets cleanly.
  • Step 7: Response Assembly (5 ms): Injects related search recommendations and dynamic knowledge cards before serializing the final JSON payload to the client.

Timing Summary: 10 + 5 + 25 + 5 + 50 + 10 + 5 = 110 ms internal execution target, comfortably within the < 200 ms p99 external SLO with ~90 ms reserved for queueing and tail-latency hedging.

Loading...

Scatter-Gather: Index Shard Coordination

Stateless query coordinators scatter queries to all document shards in parallel, presenting a major tail-latency challenge.

Loading...

The Tail Latency (Slow Shard) Problem

  • With 100 index shards, the aggregate query latency is bound by the slowest single shard.
  • At 100 shards, the p99 latency of any single shard is approximately 3 to 5 times its median (p50) latency due to garbage collection pauses, background indexing contention, or physical hardware variance.
  • The likelihood of the entire query being slowed down by at least one laggy shard scales exponentially:
    P(at least one shard is slow) = 1 - (1 - 0.01)¹⁰⁰ = 63%
  • Without software-level mitigations, 63% of all user search queries would suffer from high tail latency.

Technical Solutions

  1. Hedged Requests (Recommended): Send the query to the primary replica of an index shard. If that replica does not respond within 30 ms, dispatch an identical speculative request to a secondary replica. The coordinator accepts whichever replica responds first and cancels the outstanding request. This technique reduces tail latency substantially while incurring only an approximate 5% increase in total throughput load.
  2. Speculative Execution and Partial Returns (Graceful Degradation): If 95 out of 100 shards return results within 100 ms, complete the query using those 95 responses rather than waiting for stragglers. This strategy deliberately trades candidate recall for predictable tail latency. It should only be enabled when the product explicitly accepts approximate retrieval under heavy load, because an omitted shard could theoretically hold a top-ranked result.
  3. Tiered Index Execution: Segregate documents into Tier 0 (the top 1B popular documents), Tier 1 (extending cumulative coverage to the top 10B documents total by adding the next ~9B pages), and Tier 2 (encompassing the full 100B-page corpus by storing the remaining ~90B long-tail pages), with each tier running as an independent document-sharded, replicated cluster. Queries evaluate Tier 0 first; only if Tier 0 yields insufficient candidate volume or relevance does execution cascade to Tier 1 or Tier 2, avoiding massive coordination fan-out for the vast majority of searches.

Shard Routing Trade-off: Document vs. Term Sharding

  • Document Sharding (Google Choice): Each shard stores the complete inverted index dictionary for an assigned subset of documents. Although the query coordinator must broadcast every query across all shards in the active tier (a fixed fan-out of 100 RPCs), per-shard processing is uniform, predictable, and simple to operate. Hot-term load is addressed by pre-warming specialized in-memory hot-term posting structures and caching posting fragments on coordinators rather than restructuring shard ownership.
  • Term Sharding: Shards are partitioned by lexicographical term ranges (for example, Shard A stores terms beginning with A through C). While a search for "system design" contacts only the shards owning 's' and 'd' requiring 2 RPCs, intersecting posting lists across distributed nodes creates heavy network traffic, and popular search terms produce severe hot-shard bottlenecks.

Result Merging Across Document Shards

To accurately merge Top-100 scores returned by independent document shards, we must solve BM25 score comparability and execute sorting efficiently.

Loading...

The Global Score Comparability Problem

  • The BM25 formula depends heavily on Inverse Document Frequency (IDF):
    IDF(q) = log((N - n(q) + 0.5) / (n(q) + 0.5))
    Where N is the total documents in the corpus and n(q) is the count of documents containing the term.
  • If shards computed IDF using only their local document subsets, the same term would yield divergent IDF weights across shards, rendering local BM25 scores globally incomparable.
  • Solution: Global IDF Precomputation during Offline Index Build (MapReduce)
    1. Count total documents N across all index shards.
    2. Calculate document frequency n(q) for each term across all index shards.
    3. Compute the absolute global IDF value for each term.
    4. Persist these global IDF values directly into each shard's term dictionary.
  • At query time, every shard evaluates documents against identical global IDF constants, ensuring BM25 scores remain directly comparable across shards so the K-way merge produces a mathematically correct global top-K list.

K-Way Merge Algorithm

With 100 shards returning their local Top-100 results sorted by score descending, the coordinator must merge 10,000 candidates to select the top 1,000 for ML reranking:

  1. Initialize a Max-Heap of size 100, containing one entry representing the top result from each shard.
  2. Insert the highest scoring candidate from each of the 100 shards into the heap.
  3. Pop the maximum item from the heap, which represents the highest-ranked global result.
  4. Push the next highest scoring result from that same shard into the heap.
  5. Repeat steps 3 and 4 until the target global top 1,000 candidates are fully extracted for the ML reranking stage.

This algorithm executes in O(K x log(num_shards)) time, requiring only approximately 7,000 comparisons (O(1000 x log(100))) to merge 10,000 candidates.

  • Tie-Breaking: When text relevance scores are identical, static PageRank authority scores serve as a deterministic secondary sort key.
  • Near-Duplicate Deduplication: Because document sharding ensures each document resides on exactly one shard, identical document duplicates cannot occur. For near-duplicate content published across different URLs, SimHash fingerprints computed during offline crawling allow ML reranking workers to filter redundant pages.

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