Web crawling at scale
How large crawlers fetch billions of pages: the URL frontier, politeness per host, prioritisation and recrawl scheduling, DNS and fetching, URL normalisation and deduplication, content fingerprints, robots.txt, traps and distributing the crawl.
Reading is half of it. See this used in a real interview: walk through Design a Web Crawler →
A web crawler looks like a simple loop: take a URL, fetch it, extract links, add them to the list. At the scale of billions of pages, every part of that loop needs design: which URL to fetch next, how to avoid hammering one website, how to skip the same page reached through a thousand different URLs, and how to spread the work across hundreds of machines. "Design a web crawler" is a classic interview question that rewards knowing these parts.
The loop
- Frontier chooses the next URL to fetch.
- DNS resolver turns the host into an address (cached aggressively).
- Fetcher downloads the page, respecting robots.txt and rate limits.
- Store writes raw content to object storage.
- Parser extracts text, metadata and links.
- URL filter and deduplication normalise links and drop ones already seen.
- New URLs go back into the frontier.
Each stage is a pool of workers connected by queues, so stages scale independently. See background jobs.
The URL frontier
The frontier balances two goals: priority (fetch important and changed pages first) and politeness (never overload one host). A common design, from the Mercator crawler:
- Front queues by priority: URLs are scored (page importance, how often it changes, freshness need) and placed in priority queues.
- Back queues by host: each back queue holds URLs for one host. A heap records the earliest time each host may be fetched again.
- A fetcher thread takes the host whose next allowed time has arrived, fetches one URL, and reschedules that host after a delay (for example a few seconds, or proportional to its last response time).
This keeps many hosts busy in parallel without sending any single host more than about one request at a time.
Politeness and robots.txt
- Fetch and cache each host's robots.txt; obey disallowed paths and any crawl-delay.
- Limit concurrent connections and request rate per host and per IP (many sites share hosting).
- Back off on errors and slow responses; stop on repeated 429 or 503.
- Identify the crawler with a clear user agent and contact page.
Deduplication
The same content is reachable through many URLs:
- Normalise URLs before checking: lowercase the host, remove default ports and fragments, sort or strip tracking parameters, resolve relative paths.
- Seen-URL set: billions of URLs do not fit comfortably in memory, so use a Bloom filter in front of a disk-backed store, accepting that a tiny fraction of new URLs are skipped. See probabilistic data structures.
- Content fingerprints: hash page content to catch exact duplicates, and use SimHash or MinHash to detect near-duplicates (the same article with a different ad or date). See hashing and encoding.
Recrawl scheduling
The web changes constantly. Estimate how often each page changes from its history and importance, and schedule recrawls accordingly: a news homepage every few minutes, a rarely changing page every few months. Use conditional requests (If-Modified-Since, ETags) so unchanged pages cost little. Sitemaps and feeds hint at what changed.
Traps and junk
- Crawler traps: infinite calendars, session ids in URLs, endlessly deep paths. Limit URL length and depth per site and cap pages per host.
- Spam and low-quality sites: score and deprioritise.
- Soft 404s: pages that return 200 but say "not found"; detect by content.
- Rendering: many pages need JavaScript; render only the important ones in headless browsers, since rendering is many times more expensive than fetching.
Distributing the crawl
- Partition by host (hash of the host name) so each crawler node owns a set of hosts. Politeness then stays local to one node, and per-host state does not need coordination.
- Nodes exchange discovered URLs that belong to other partitions in batches.
- Use consistent hashing so adding nodes moves few hosts.
- Place crawler nodes in several regions to fetch from closer locations.
Estimates
For 1 billion pages a month: about 400 pages per second on average. At 100 KB per page, that is roughly 100 TB per month of raw content before compression, and around 40 MB per second of bandwidth. DNS lookups and connection setup often dominate fetch latency, so cache DNS and reuse connections. See back-of-envelope estimation.
In the interview
For Design a Web Crawler: a pipeline of fetch, parse, filter and store; a frontier with priority front queues and per-host back queues for politeness; URL normalisation, a Bloom-filter-backed seen set and content fingerprints for deduplication; recrawl by change rate; partitioning by host; and trap protection.
Checklist
- Frontier with priority and per-host politeness.
- robots.txt cached and obeyed; per-host and per-IP rate limits.
- URL normalisation, seen-URL Bloom filter, content and near-duplicate fingerprints.
- Change-rate-based recrawl with conditional requests.
- Trap limits on depth, URL length and pages per host.
- Partitioning by host; DNS caching and connection reuse.