Google Search is the most studied large system in computing and one of the least documented in detail. Google has never published its current architecture, index size or ranking weights. It has, however, published a trail of papers and documentation spanning 25 years: the original 1998 design, a 2003 description of how a query is served across thousands of machines, the 2010 move to continuous indexing, the techniques it uses to tame tail latency, and public guidance on crawling, rendering and its named ranking systems.
This article reconstructs the architecture from those sources, separating what Google has said from what is inference. It walks the three pipelines, crawling, indexing and serving, with code for the core algorithms, works a sizing example for a smaller web index, and ends with the lessons that transfer to any search system you build. Where a claim comes from a specific paper, the paper is named.
The architecture in one picture
What is public, and what is not
Four primary sources anchor the picture. Brin and Page's 1998 paper, The Anatomy of a Large-Scale Hypertextual Web Search Engine, describes the original crawler, the forward and inverted indexes split into barrels, and PageRank. Barroso, Dean and Holzle's 2003 IEEE Micro paper, Web Search for a Planet, describes query serving: the index divided into shards by document, each shard served by a pool of replicated machines, and separate document servers that produce titles and snippets. Peng and Dabek's OSDI 2010 paper on Percolator describes the incremental indexing system behind Caffeine, which Google announced in June 2010. Dean and Barroso's 2013 The Tail at Scale describes how latency-critical fan-out services cope with slow machines.
Google Search Central documentation adds the current, public view: Search works in three stages, crawling, indexing and serving; Googlebot renders pages with a recent version of Chromium; and ranking is done by many named systems described in Google's guide to its ranking systems. Everything else, including the number of shards, the hardware and how signals are combined, is not public, and you should be suspicious of any article that states it with confidence.
Crawling: frontier, budget and politeness
Crawling starts from a frontier of known URLs, seeded from previous crawls, sitemaps and links discovered in fetched pages. The hard problems are scheduling and politeness, not fetching. Google's documentation on crawl budget describes two inputs: a crawl capacity limit, how fast Googlebot can fetch from a site without hurting it, which rises when the server responds quickly and falls on errors and slow responses; and crawl demand, how much Google wants to crawl, driven by popularity and staleness. A site gets the crawl that both allow.
A frontier that respects both looks like a priority queue of URLs plus a per-host timer:
import heapq, time
from urllib.parse import urlsplit
class Frontier:
def __init__(self):
self.queues = {} # host -> heap of (-priority, url)
self.ready = [] # heap of (next_allowed_time, host)
self.delay = {} # host -> seconds between fetches
def add(self, url, priority):
host = urlsplit(url).netloc
if host not in self.queues:
self.queues[host] = []
self.delay[host] = 1.0
heapq.heappush(self.ready, (time.time(), host))
heapq.heappush(self.queues[host], (-priority, url))
def next(self):
when, host = heapq.heappop(self.ready)
time.sleep(max(0.0, when - time.time()))
_, url = heapq.heappop(self.queues[host])
return host, url
def done(self, host, latency_s, ok):
# capacity limit: back off on errors or slow responses, speed up on fast ones
d = self.delay[host]
d = d * 2 if (not ok or latency_s > 2.0) else max(0.1, d * 0.9)
self.delay[host] = d
if self.queues[host]:
heapq.heappush(self.ready, (time.time() + d, host))Priority is where crawl demand lives: a page's estimated importance and its predicted change rate decide how soon it is refetched. robots.txt is fetched per host and cached, and HTTP status codes carry meaning: a 5xx or 429 slows the crawl, while a 404 eventually removes the URL.
Rendering as a separate phase
Since 2019 Googlebot has rendered pages with an evergreen headless Chromium, so content inserted by JavaScript can be indexed. Rendering is far more expensive than fetching HTML, and Google's JavaScript SEO documentation describes it as a separate phase: pages are queued for rendering after the initial crawl, and links and content found in the rendered page then feed back into crawling and indexing. The architectural lesson is to decouple a cheap first pass from an expensive enrichment pass, and to schedule the expensive one by value.
Indexing: from batch to Percolator
Indexing turns fetched pages into postings. Each page is parsed into tokens with positions and fields such as title, anchor text from inbound links and body; near-duplicate pages are clustered and one canonical URL is chosen per cluster, using rel=canonical as a strong hint but not a command; and links are extracted for both the frontier and link analysis.
Until 2010 Google rebuilt its index in large batches with MapReduce, so a newly crawled page waited for the next batch. Percolator replaced that with incremental processing on Bigtable: cross-row, cross-table transactions with snapshot isolation, and observers, code that runs when a watched column changes. When a document is fetched, an observer processes it, updates the cluster of duplicates it belongs to, and triggers further observers downstream. The paper reports that the change reduced the average age of documents in search results by 50 percent. The cost was higher resources per document than the batch system; Google accepted that for freshness.
Postings are partitioned by document, not by term: each shard holds a complete inverted index for a subset of pages. Web Search for a Planet explains why. Every shard can answer any query on its own subset, so a query fans out to all shards in parallel and the results are merged, and adding pages means adding shards. Term partitioning would send heavy terms to a few overloaded machines and require shipping long posting lists between machines to intersect them.
Serving: fan-out, merge and the tail
A query arrives at a front end, which handles spelling, language and query understanding, then goes to a root that fans it out to one replica of every index shard. Each shard intersects posting lists for the query terms, scores matches with cheap signals and returns its local top results. The root merges these, applies more expensive ranking to the shortlist, and asks document servers for titles and snippets of the final results. The shape is easy to express:
import heapq, math
from concurrent.futures import ThreadPoolExecutor, wait
def shard_search(shard, terms, k):
# shard.postings: term -> {doc_id: term_frequency}; shard.prior: doc_id -> static score
lists = [shard.postings.get(t, {}) for t in terms]
if not all(lists):
return []
lists.sort(key=len) # intersect from the rarest term
docs = set(lists[0]).intersection(*lists[1:])
def score(d):
tfidf = sum(math.log1p(l[d]) * math.log(shard.n_docs / len(l)) for l in lists)
return tfidf + 2.0 * shard.prior.get(d, 0.0) # prior such as link-based authority
return heapq.nlargest(k, ((score(d), d) for d in docs))
def search(shards, terms, k=10, deadline_s=0.2):
with ThreadPoolExecutor(len(shards)) as pool:
futures = [pool.submit(shard_search, s, terms, k) for s in shards]
finished, _ = wait(futures, timeout=deadline_s) # partial results beat late results
partial = [f.result() for f in finished]
return heapq.nlargest(k, (hit for hits in partial for hit in hits))The deadline matters as much as the scoring. With fan-out, the slowest shard sets query latency. The Tail at Scale gives the techniques: hedged requests, sending a second copy of a request to another replica if the first has not answered within roughly the 95th-percentile latency; tied requests, sending to two replicas that cancel each other when one starts work; micro-partitioning, many more shards than machines so load can be moved in small pieces; and returning good-enough results when a small fraction of shards are slow. The paper reports a Bigtable benchmark in which deferring a second request by 10 ms cut 99.9th-percentile latency from 1,800 ms to 74 ms while adding only about 2 percent more requests.
Ranking: PageRank and the cascade
The 1998 paper's best-known contribution is PageRank: a page is important if important pages link to it. It is the stationary distribution of a random surfer who follows a random outlink with probability d, typically 0.85, and jumps to a random page otherwise, computed by power iteration:
def pagerank(links, d=0.85, iters=50):
# links: page -> list of pages it links to
pages = set(links) | {q for outs in links.values() for q in outs}
n = len(pages)
pr = {p: 1.0 / n for p in pages}
for _ in range(iters):
nxt = {p: (1 - d) / n for p in pages}
for p in pages:
outs = links.get(p, [])
if outs:
share = d * pr[p] / len(outs)
for q in outs:
nxt[q] += share
else: # dangling page: spread evenly
for q in pages:
nxt[q] += d * pr[p] / n
pr = nxt
return prModern ranking is far richer, and Google says PageRank remains one of many signals. Its public guide to ranking systems names, among others, RankBrain, launched in 2015 to relate words to concepts; neural matching; BERT, applied to Search in 2019 to understand how word combinations express meaning; and systems for freshness, reliable information, spam and deduplication. Architecturally, these sit in a cascade: cheap signals score millions of candidates inside the shards, and progressively more expensive models, up to large neural networks, rescore only the short list that reaches the root. That cascade is the only way neural models can run inside a query's latency budget.
Worked example: a 1-billion-page index
Scale the design down to something you might build: a vertical search engine over 1 billion pages. Assume 10 KB of extracted text per page, about 1,500 tokens, and a compressed positional index of around 30 percent of text size, so about 3 TB of postings. If each serving machine holds 64 GB of index in memory to keep latency low, you need about 47 shards; round to 64 for headroom and growth. At 500 queries per second, with each shard replica handling 100 queries per second, each shard needs 5 serving replicas plus one for zone loss, so 64 x 6 = 384 index machines.
Tail latency now dominates. If each shard answers within 50 ms 99 percent of the time, the chance that all 64 shards answer within 50 ms is 0.99^64, about 53 percent; nearly half of all queries wait on at least one slow shard. Hedging after the 95th percentile, a deadline with partial results, and replicas spread across failure domains turn that from an outage-shaped problem into a few percent of extra load. This sizing logic is developed further in Designing Search at Scale, and the components of a general search stack are in the search architecture walkthrough.
Failure modes
- Crawler as attacker. An aggressive crawler overloads small sites. Adaptive per-host delays and honouring 429 and 5xx are mandatory.
- Duplicate floods. Session IDs, tracking parameters and faceted URLs create infinite URL spaces. Canonicalise, cap URLs per host and detect near-duplicates before indexing.
- Stale index. Batch rebuilds leave fresh content invisible for days. Incremental indexing, as Percolator showed, trades resources for freshness.
- Fan-out tail. One slow shard delays every query. Hedge, tie, micro-partition and accept partial results.
- Expensive ranking on too many candidates. Running a large model per candidate blows the budget. Keep the cascade and measure each stage's candidate count.
- Adversarial content. Link farms and keyword stuffing target any public signal. Assume every signal will be gamed and keep spam systems separate from relevance.
Trade-offs
| Decision | Google's published choice | Alternative and its cost |
|---|---|---|
| Index partitioning | By document, fan-out to all shards | By term: fewer machines per query, hot terms and big intersections |
| Index freshness | Incremental (Percolator) | Batch: cheaper per document, stale results |
| Tail latency | Hedging and partial results | Wait for all shards: complete but slow |
| Ranking | Cascade from cheap to neural | One model on all candidates: too slow |
| Rendering | Separate, deferred phase | Render on fetch: slower crawl everywhere |
To go deeper on pieces of this stack: link analysis at scale is in PageRank on Spark, the storage system under Percolator is in Bigtable architecture, and how lexical and neural retrieval combine is in hybrid search with BM25, dense retrieval and reranking.
What to do next
- Read the primary sources in order: the 1998 Anatomy paper, Web Search for a Planet, the Percolator paper and The Tail at Scale. Each is short and each changed practice.
- Build a toy engine from the code here over a few hundred thousand pages, then measure query latency with 8, 32 and 64 shards to see the tail grow.
- Add a deadline and hedged requests to your fan-out and measure p99 before and after, along with the extra request rate.
- Implement a two-stage ranker: a cheap first stage over all matches and a model over the top 100, and record the latency of each stage.
- If you run a website, read Google's crawl budget and JavaScript SEO documentation and check your server's error rate and response times as Googlebot sees them in Search Console.
- For any system you design, decide document versus term partitioning, batch versus incremental indexing and the latency budget per stage before writing code.