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
- Crawl: fetch pages politely and prioritise by importance and change rate. See web crawling at scale.
- Process: parse HTML, extract text, links, titles and metadata; detect language; remove duplicates and near-duplicates; classify spam. See hashing and encoding.
- Analyse links: build the link graph; compute link-based importance (PageRank) and anchor text for target pages. See large-scale graph processing.
- Index: build inverted indexes and per-document data (forward index) for ranking and snippets.
- 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
| Scheme | How | Pros | Cons |
|---|---|---|---|
| Document-partitioned | each shard indexes a subset of documents, all terms | each shard answers independently; easy to add documents; balanced | every query hits every shard |
| Term-partitioned | each shard holds all postings for some terms | a query touches only shards for its terms | huge 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
- Query understanding: spelling correction, synonyms, language, intent (navigational, informational, local), entities. See search ranking and relevance.
- Cache check: popular queries are answered from a result cache.
- Scatter: send the query to all index shards (through a tree of aggregators to limit fan-in).
- Per-shard retrieval and scoring: find matching documents and score them with fast signals; return the top few.
- Gather and re-rank: merge, then apply heavier ranking (machine-learned models) to the top hundreds.
- 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.
Autocomplete and related features
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.