SysDesignPrep.com
Study guide 72 of 183

Web search engine architecture

How a Google-style web search engine fits together: crawling, document processing, building inverted indexes, document-partitioned versus term-partitioned indexes, tiered indexes, query serving with scatter-gather, ranking signals including links, freshness, caching and the numbers involved.

Reading is half of it. See this used in a real interview: walk through Design a Web Crawler →

Web search is the archetypal large-scale system: hundreds of billions of pages, an index of petabytes, tens of thousands of queries per second, and answers in a few hundred milliseconds. Few interviews ask you to design all of Google, but crawler and search questions borrow its architecture, and knowing the overall shape helps you place every component.

The pipeline

  1. Crawl: fetch pages politely and prioritise by importance and change rate. See web crawling at scale.
  2. Process: parse HTML, extract text, links, titles and metadata; detect language; remove duplicates and near-duplicates; classify spam. See hashing and encoding.
  3. Analyse links: build the link graph; compute link-based importance (PageRank) and anchor text for target pages. See large-scale graph processing.
  4. Index: build inverted indexes and per-document data (forward index) for ranking and snippets.
  5. Serve: answer queries by retrieving and ranking candidates.

Steps 2 to 4 run as large batch and streaming jobs. See MapReduce and Spark.

The inverted index at web scale

For each term, a posting list of documents containing it, with positions and per-document features, compressed heavily (delta encoding of document ids plus variable-length integers). Postings are sorted by document id (for fast intersection) or by a static quality score (so the best documents come first and queries can stop early). See how Elasticsearch works.

Partitioning the index

SchemeHowProsCons
Document-partitionedeach shard indexes a subset of documents, all termseach shard answers independently; easy to add documents; balancedevery query hits every shard
Term-partitionedeach shard holds all postings for some termsa query touches only shards for its termshuge posting lists for common terms; hard to balance; multi-term queries move lots of data

Web search engines use document partitioning: a query fans out to thousands of shards in parallel, each returns its top results, and aggregators merge them. Each shard is replicated for throughput and availability. See sharding and partitioning.

Serving a query

  1. Query understanding: spelling correction, synonyms, language, intent (navigational, informational, local), entities. See search ranking and relevance.
  2. Cache check: popular queries are answered from a result cache.
  3. Scatter: send the query to all index shards (through a tree of aggregators to limit fan-in).
  4. Per-shard retrieval and scoring: find matching documents and score them with fast signals; return the top few.
  5. Gather and re-rank: merge, then apply heavier ranking (machine-learned models) to the top hundreds.
  6. Snippets: fetch document text from a document store and generate query-dependent snippets for the top results.

With thousands of shards, the slowest shard dominates latency: use replicas, hedged requests and partial results after a deadline. See tail latency.

Tiered indexes

Not all pages deserve equal resources. A small top tier of the most important and fresh documents lives in memory on many replicas; larger lower tiers sit on SSDs and are consulted only when the top tier does not yield enough good results. Most queries are satisfied by the top tier, which keeps cost manageable.

Freshness

News and fast-changing pages must appear within minutes, while rebuilding the full index takes days. Use a separate real-time index fed by a streaming pipeline for recent documents, merged with the main index at query time (an LSM-like approach), and periodically fold it into the main index. See Lambda vs Kappa architecture.

Ranking signals

Hundreds of signals: text relevance (BM25-style), anchor text, link-based authority, freshness, page quality and spam scores, user engagement, location and language, and personalisation, combined by learned models. Evaluated with human relevance ratings and online experiments. See feature flags and A/B testing.

The numbers

Hundreds of billions of documents, compressed indexes of hundreds of petabytes across tiers, tens of thousands of queries per second at peak, and a few hundred milliseconds end to end. Even a single query touches thousands of machines, which is why caching popular queries and tiering the index matter so much. See back-of-envelope estimation.

Query suggestions are built from query logs into prefix structures. See tries and autocomplete and Design Typeahead.

In the interview

For a search-heavy extension of Design a Web Crawler: "Crawled pages are deduplicated and processed in batch; link analysis computes authority; the inverted index is document-partitioned into shards replicated across machines, with an in-memory top tier and SSD-based lower tiers, plus a real-time index for fresh pages. Queries go through a result cache, then scatter-gather across shards with deadlines, then ML re-ranking and snippet generation."

Checklist

  • Crawl, process, link analysis, index, serve as separate stages.
  • Compressed posting lists ordered for fast intersection or early termination.
  • Document-partitioned, replicated shards with aggregator trees.
  • Result caching, tiered indexes and a real-time index for freshness.
  • Deadline-bound scatter-gather with hedging.
  • Multi-signal learned ranking evaluated offline and online.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.