Skip to main content

Web Crawler (Googlebot Scale)

Sharpened prompt. Design a crawler fetching 10B pages/month across 200M domains, never sending more than one request per second to any single host, deduplicating URLs against a set of 500B already-seen URLs, escaping infinite calendar and faceted-navigation traps, recrawling news sites hourly and static pages monthly, and rendering JavaScript where it matters.

The crawler's difficulty is not fetching — HTTP is easy. It is scheduling under a per-domain constraint, deduplication at a scale where the index does not fit in memory, and not being destroyed by adversarial or accidental infinity.

1. Problem framing​

Functional requirements​

  • Discover URLs from seeds, links, and sitemaps.
  • Fetch politely, respecting robots.txt and crawl-delay.
  • Deduplicate URLs and detect near-duplicate content.
  • Prioritise: crawl important and fresh pages more often.
  • Render JavaScript for pages that require it.

Non-functional requirements​

PropertyTargetConsequence
Throughput10B pages/month ≈ 3,900 pages/secThousands of concurrent fetches across many machines
Politeness≤ 1 req/sec/host, alwaysA per-host serialisation point, distributed
URL dedup500B URLs, O(1) membershipBloom filter in front of an on-disk store
Trap resistanceNo crawler stuck in an infinite spaceHeuristic circuit breakers per domain
FreshnessNews in minutes, archives in monthsPriority is a function of change rate and importance

Back-of-the-envelope​

Rate: 10B/month = 3,900 pages/sec sustained
Bandwidth: 3,900 × ~100 KB = 390 MB/s = 3.1 Gbps sustained inbound
Storage: 10B × 100 KB = 1 PB/month raw; ~200 TB/month compressed (text only ~50 TB)
URL set: 500B URLs × 16 B fingerprint = 8 TB -> too big for RAM
Bloom filter @1% FP: 500B × 9.6 bits = 600 GB -> still too big for one node
-> shard by domain hash: 100 nodes × 6 GB Bloom each. Now it fits.
Politeness: 200M domains, ≤1 req/s each. But traffic is Zipfian — a handful of domains
could absorb most of your capacity if you let them.
Frontier: ~50B pending URLs × 100 B = 5 TB -> disk-backed priority queues

The Bloom-filter sharding arithmetic is the number that decides the architecture: you cannot hold the seen-set on one machine, so the crawler must be partitioned by domain — which is convenient, because politeness wants exactly the same partitioning.

2. High-level architecture​

3. Component inventory​

ComponentConcrete choiceWhy this one
FrontierMercator two-tier queues, disk-backedSeparates priority from politeness — the key structural idea
URL dedupSharded Bloom filter → RocksDBBloom answers "definitely new" cheaply; RocksDB is the authority on the rest
FetcherAsync I/O (Go/Rust/Python asyncio), thousands of sockets per nodeCrawling is I/O-bound; threads are the wrong model
DNSLocal caching resolver with a long TTL overrideDNS is a hidden bottleneck at 3,900 fetches/sec
RendererHeadless Chrome pool, used selectively~100× the cost of a plain fetch; must be rationed
Content dedupSimHash with Hamming-distance bucketsNear-duplicate detection at index time
CoordinationPartition by hash(registrable_domain)One owner per domain gives politeness for free

4. The toughest parts​

4.1 Politeness: the constraint that shapes everything​

Why it's hard. You must never overwhelm a small site — that is both an ethical obligation and, practically, how you get permanently blocked or sued. But politeness is a per-host serialisation constraint sitting inside a massively parallel system: with 3,900 fetches/sec across thousands of workers, any worker might pick a URL from a host another worker just hit a millisecond ago. A global lock per host would serialise the whole crawler.

Solution — the Mercator two-tier frontier, which decouples "what to crawl" from "when we may crawl it."

class Frontier:
"""Front queues encode PRIORITY. Back queues encode POLITENESS.
Exactly one back queue per host, so a host is serialised by construction."""

def __init__(self, n_front=10, n_back=3000):
self.front = [DiskQueue() for _ in range(n_front)] # priority 0 = highest
self.back = {} # host -> DiskQueue
self.host_of_back = {} # back queue -> host
self.heap = [] # (next_ok_time, back_id)

def add(self, url, priority):
self.front[priority].push(url) # discovery is cheap

def _refill_back(self, back_id):
"""Pull from a front queue (biased toward high priority) into an empty
back queue, assigning it to whatever host that URL belongs to."""
url = self._pop_biased()
host = registrable_domain(url)
self.back[back_id] = DiskQueue([url])
self.host_of_back[back_id] = host

def next_url(self):
# The heap gives the back queue whose host is eligible soonest.
t, back_id = heapq.heappop(self.heap)
if (wait := t - time.time()) > 0:
time.sleep(wait) # or yield in async
url = self.back[back_id].pop()
host = self.host_of_back[back_id]

# Reschedule this host by its own crawl delay: robots.txt, or 10x the
# observed response time — slow servers get MORE space, not less.
delay = max(robots.crawl_delay(host) or 0, 10 * stats.avg_latency(host), 1.0)
if self.back[back_id].empty():
self._refill_back(back_id)
heapq.heappush(self.heap, (time.time() + delay, back_id))
return url

Three properties fall out of this structure. One back queue per host means the politeness constraint is enforced by the data structure rather than by a lock. The heap makes "which host may I hit next?" an O(log n) operation instead of a scan. And adaptive delay — 10× the observed response time — means a struggling server automatically gets crawled less, which is the single most important courtesy a crawler can offer and is much better than a fixed constant.

Two additional rules worth stating: rate-limit by IP address as well as hostname, because thousands of virtual hosts often share one server, and cache robots.txt per host with a 24-hour TTL, treating a fetch failure as "disallow" rather than "allow."

4.2 Deduplicating URLs when the set does not fit in memory​

Why it's hard. Every page yields ~100 links, so 10B pages produce a trillion link observations against a set of 500B distinct URLs. Checking "have I seen this?" must be sub-millisecond, and the set is 8 TB of fingerprints — far beyond RAM on any node. A naive database lookup per link means a trillion random reads.

Solution — canonicalise first, then a sharded Bloom filter in front of an LSM-tree store.

def canonicalise(url: str) -> str:
"""Most 'distinct' URLs are the same page. Canonicalisation is the cheapest
possible dedup and typically collapses the URL space by 30-50%."""
u = urlsplit(url.strip())
scheme = "https" if u.scheme in ("http", "https") else u.scheme
host = u.hostname.lower().removeprefix("www.")
port = "" if u.port in (None, 80, 443) else f":{u.port}"
path = posixpath.normpath(unquote(u.path) or "/").replace("//", "/")

# Strip tracking parameters; sort the rest so order doesn't create duplicates.
q = [(k, v) for k, v in parse_qsl(u.query, keep_blank_values=True)
if not k.lower().startswith(("utm_", "fbclid", "gclid", "msclkid", "ref"))]
query = urlencode(sorted(q))

return urlunsplit((scheme, host + port, path, query, "")) # fragment always dropped

def seen(url: str) -> bool:
fp = xxhash.xxh128(canonicalise(url)).digest()
shard = fp[0] % N_SHARDS # partition by fingerprint
if not bloom[shard].might_contain(fp):
bloom[shard].add(fp)
rocksdb[shard].put(fp, b"") # authoritative record
return False # definitely new
return rocksdb[shard].contains(fp) # Bloom said maybe: verify on disk

The two-stage structure is what makes it work. The Bloom filter has no false negatives, so "might_contain == False" is a definitive "new" and needs no disk read — and since the overwhelming majority of discovered links are already known, most lookups are answered by the second branch, with the Bloom filter's false-positive rate (1%) determining how many disk reads you actually perform. Sizing it at 9.6 bits per element gives 1% FP; going to 14.4 bits gives 0.1% and trades RAM for I/O.

Shard by fingerprint so each node owns 1/100th of the space and holds a 6 GB filter. Route link discovery to the owning shard — a small network hop, batched, rather than a global structure.

The trade-off to acknowledge: a Bloom false positive means a genuinely new URL is treated as seen and never crawled. That is why RocksDB verifies rather than the filter deciding alone. If you were willing to accept the loss, you could skip the disk check entirely and crawl 1% fewer pages — a legitimate design point worth naming, and one that some crawlers actually take.

4.3 Spider traps: infinity generated by accident​

Why it's hard. Enormous numbers of URLs are auto-generated and effectively infinite. A calendar widget offers "next month" forever. Faceted navigation on an e-commerce site produces a combinatorial explosion of filter permutations. A misconfigured server maps /a/b/a/b/a/b/... to the same page. A session ID in the query string makes every URL unique. None of these are malicious; all of them will consume your entire crawl budget on one domain if you let them.

Solution — layered heuristics, each cheap, plus a per-domain budget as the backstop.

def trap_score(url: str, domain_stats: DomainStats) -> float:
u, score = urlsplit(url), 0.0
segs = [s for s in u.path.split("/") if s]

# 1. Depth: legitimate content is rarely more than ~10 levels deep.
if len(segs) > 10: score += 0.3 * (len(segs) - 10)

# 2. Repeating segments: /shop/shop/shop/ is a server misconfiguration.
if len(segs) != len(set(segs)) and max(Counter(segs).values()) >= 3: score += 1.0

# 3. Parameter explosion: faceted navigation.
params = parse_qsl(u.query)
if len(params) > 6: score += 0.2 * (len(params) - 6)

# 4. Session-like parameters make every URL unique.
if any(k.lower() in {"sid", "sessionid", "phpsessid", "jsessionid"} for k, _ in params):
score += 1.0

# 5. Date-like paths far in the future: calendar widgets.
if (d := extract_date(u.path)) and d > date.today() + timedelta(days=365):
score += 1.0

# 6. The strongest signal: this domain yields lots of URLs and no new CONTENT.
if domain_stats.pages_crawled > 1000 and domain_stats.unique_content_ratio < 0.1:
score += 2.0

return score

Signal 6 is the one that generalises, and it is the one to emphasise: a trap is defined by producing URLs without producing information. Rules 1–5 encode specific known patterns and will always be incomplete; measuring the ratio of unique content fingerprints to pages crawled catches novel traps you never anticipated.

Back it with a hard per-domain budget proportional to the domain's importance (PageRank-like score, or observed traffic): a major news site gets millions of pages, an unknown domain gets a thousand until it proves it deserves more. The budget is what guarantees a trap can waste at most a bounded amount of crawl capacity, no matter how clever it is.

4.4 Deciding what to crawl next, and how often to come back​

Why it's hard. The frontier holds 50B URLs and you can fetch 3,900/sec. Crawl order determines the quality of the entire index. Meanwhile pages change at wildly different rates — a news homepage every few minutes, a 2011 blog post never — and crawling both on the same schedule wastes bandwidth on one and misses the other entirely.

Solution — a priority score for discovery, and an adaptive recrawl interval driven by observed change.

def discovery_priority(url, referrer) -> int:
s = 0.0
s += 3.0 * domain_authority(registrable_domain(url)) # PageRank-ish prior
s += 2.0 * (1.0 / (1 + path_depth(url))) # shallow pages matter more
s += 1.5 * inlink_count(url) # many pages point here
s += 2.0 if is_sitemap_listed(url) else 0.0 # the site told us it matters
s += 1.0 if referrer.is_recently_changed else 0.0
return bucket(s, n_buckets=10) # -> a front queue index

def next_recrawl(page) -> timedelta:
"""Adaptive: converge on each page's actual change rate rather than guessing."""
if page.changed_since_last_crawl:
page.interval = max(MIN_INTERVAL, page.interval * 0.5) # halve: check sooner
else:
page.interval = min(MAX_INTERVAL, page.interval * 1.5) # grow: back off
# Importance compresses the interval: important pages are checked more often
# even when they change rarely, because staleness costs more.
return page.interval / (1 + math.log1p(page.importance))

The multiplicative-decrease / multiplicative-increase loop converges quickly on each page's real change rate without needing a model, and it self-corrects when a dormant page becomes active. Bound it at both ends (say 5 minutes to 90 days) so nothing oscillates or is forgotten.

Two cheap wins worth naming. Conditional requests: send If-Modified-Since and If-None-Match; a 304 Not Modified costs a few hundred bytes instead of 100 KB, so most recrawls become nearly free. And sitemaps with lastmod: sites that publish them are telling you exactly what changed, which is strictly better information than any inference — always prefer explicit signals over heuristics when a site offers them.

4.5 The same content at a thousand URLs​

Why it's hard. URL dedup catches identical URLs. It does nothing about the same article at example.com/article, example.com/article?print=1, m.example.com/article, and example.com/amp/article, or about syndicated content republished across hundreds of sites. Indexing all of them wastes storage and produces duplicate search results.

Solution — content fingerprinting with SimHash, plus honouring explicit canonical signals.

def simhash(text: str, bits: int = 64) -> int:
"""Locality-sensitive: similar documents produce similar fingerprints, so
near-duplicates are found by Hamming distance rather than exact match."""
v = [0] * bits
for shingle, weight in shingles(text, n=4): # 4-word shingles
h = xxhash.xxh64(shingle).intdigest()
for i in range(bits):
v[i] += weight if (h >> i) & 1 else -weight
return sum(1 << i for i in range(bits) if v[i] > 0)

def near_duplicate(fp: int) -> int | None:
"""Find any stored fingerprint within Hamming distance 3.
Trick: split the 64 bits into 4 blocks of 16. Two fingerprints within
distance 3 must match exactly on at least one block (pigeonhole), so
four exact-match index lookups replace a full scan."""
for block in range(4):
key = (fp >> (16 * block)) & 0xFFFF
for candidate in index[block].get(key, ()):
if popcount(fp ^ candidate) <= 3:
return candidate
return None

The pigeonhole trick is the detail worth showing: without it, near-duplicate detection is a scan over billions of fingerprints; with it, it is four hash lookups. That is the difference between feasible and not.

Before reaching for fingerprints, honour the explicit signals: <link rel="canonical">, og:url, and 301 redirect targets. A site telling you which URL is canonical is more reliable than any inference, and using it costs nothing.

When duplicates are found, pick a canonical by a stable rule (highest authority domain, shortest URL, earliest first-seen) and store the others as pointers. Also feed this back into crawl scheduling: a domain producing overwhelmingly duplicate content gets its budget cut — which is signal 6 from 4.3, closing the loop between content analysis and politeness.

4.6 JavaScript, and why you cannot render everything​

Why it's hard. A large fraction of the modern web renders content client-side, so a plain HTTP fetch returns an empty shell. But headless Chrome costs roughly 100× a simple fetch in CPU and memory and takes seconds rather than milliseconds. Rendering all 3,900 pages/sec would need an enormous fleet, and most pages do not need it.

Solution — a two-pass architecture: fetch cheaply, detect emptiness, and render selectively.

async def crawl(url):
raw = await http_fetch(url) # ~50 ms, ~1 MB of RAM

if not needs_rendering(raw):
return parse(raw) # ~95% of pages take this path

# Queue for rendering rather than blocking a fetcher on a 3-second render.
await render_queue.put(url, priority=render_priority(url))
return parse(raw) # index what we have meanwhile

def needs_rendering(raw) -> bool:
text_ratio = len(visible_text(raw)) / max(1, len(raw.body))
return (text_ratio < 0.05 # almost no text
or has_spa_root(raw) # <div id="root"></div>
or (script_bytes(raw) > 10 * text_bytes(raw))) # JS-dominant

Points worth making. Detect rather than assume — the emptiness heuristic (text-to-markup ratio, SPA root elements) is cheap and accurate enough. Learn per domain: once a domain has needed rendering ten times, render it by default and skip the first pass. Ration by importance, since the render queue will always be shorter than the fetch queue. And cap render time at a few seconds with a hard kill; a page that has not settled by then is not worth the fleet capacity.

Where a site offers server-side rendering, an API, or structured data (JSON-LD, microdata), prefer it — it is cheaper for both parties and more reliable than scraping a rendered DOM.

4.7 Coordinating a thousand crawler nodes​

Why it's hard. The crawler runs on many machines. Two nodes must never fetch the same URL, must never both hit the same host, and must survive nodes joining and leaving without losing frontier state or violating politeness during the transition.

Solution — partition by registrable domain, with consistent hashing and durable per-partition state.

# One owner per domain. Politeness, dedup shard, and frontier partition all align
# on the same key, so there is nothing to coordinate at crawl time.
def owner_of(url: str) -> NodeId:
return ring.get(registrable_domain(url)) # consistent hash, vnodes

# Discovered links are routed to their owner in batches — a small network cost
# that buys complete independence between nodes.
async def emit_links(links):
by_owner = defaultdict(list)
for l in links:
by_owner[owner_of(l)].append(l)
await asyncio.gather(*(rpc[node].add_urls(batch) for node, batch in by_owner.items()))

Aligning three concerns on one partition key is the design insight worth stating: politeness needs one owner per host, dedup needs a sharded set, and the frontier needs partitioning — so make them the same partition. Then no distributed locks are needed anywhere in the hot path.

Handle rebalancing carefully. When a node joins, some domains move; the new owner must not immediately hammer a host the old owner just fetched. Persist per-host last-fetch timestamps in a shared store (Redis or the frontier's own durable state) so politeness survives ownership changes. And keep the frontier itself on durable local storage with replication, since losing 50B pending URLs means re-discovering the web.

5. What breaks first​

EventFirst failureMitigation
A huge site (millions of pages)One node's frontier and budget consumedPer-domain budget proportional to importance; shard very large domains by path prefix
Spider trapCrawl capacity wasted on infinityLayered heuristics plus the unique-content-ratio signal and a hard budget
DNS resolver saturationFetch latency spikes fleet-wideLocal caching resolver, long TTL override, pre-resolve popular hosts
Bloom filter fills upFalse-positive rate climbs, new URLs skippedScalable Bloom filters (add layers) or periodic rebuild from RocksDB
Node lossIts domains stop being crawledConsistent hashing reassigns; shared last-fetch timestamps preserve politeness
Render queue backlogJS-heavy sites go stalePriority rationing; index the raw HTML meanwhile; scale renderers independently
Site blocks the crawlerLoss of an important sourceAdaptive delay on 429/503, honour Retry-After, publish a crawler identity and contact

6. Cheat sheet​

  • Frontier: Mercator two tiers — front queues carry priority, back queues (one per host) carry politeness, with a heap over next-eligible times.
  • Politeness: delay = max(robots crawl-delay, 10× observed latency, 1s). Rate-limit by IP as well as host. Failed robots.txt means disallow.
  • URL dedup: canonicalise (strip tracking params, sort query, drop fragment) → sharded Bloom filter → RocksDB verification. 9.6 bits/URL at 1% FP.
  • Traps: depth, repeated segments, parameter explosion, session IDs, future dates — and above all the unique-content-to-pages-crawled ratio. Hard per-domain budget as the backstop.
  • Recrawl: multiplicative decrease on change, multiplicative increase on no-change, compressed by importance. Conditional requests make most recrawls nearly free.
  • Content dedup: SimHash with the 4-block pigeonhole trick for Hamming-distance-3 lookups; honour rel=canonical first.
  • JS: detect emptiness, render selectively, learn per domain, ration by importance.
  • Coordination: partition by registrable domain so politeness, dedup, and frontier all share one key — no distributed locks anywhere.
  • The one-liner: "A politeness-constrained scheduler wrapped around a set-membership problem too big for RAM — the Mercator two-tier frontier solves the first, a sharded Bloom-plus-LSM solves the second, and everything else is defending crawl budget against infinity."