Why this problem matters
A web crawler (also called a spider or bot) is a program that downloads web pages, finds the links inside them, and then downloads those linked pages too, over and over. Search engines use crawlers to discover what to put in their index. Archives, price-comparison sites, SEO tools and researchers building datasets all run crawlers too.
The basic loop fits in ten lines: take a URL from a list, download it, extract links, add new links to the list. The interview is about everything that loop ignores at scale:
- the list of URLs to visit (the frontier) grows to billions of entries and must decide what to fetch next,
- you must be polite: never hammer one website, and obey its
robots.txtrules, - the same page appears under many URLs, and the same content appears on many pages, so you need deduplication,
- some sites generate infinite URLs (crawler traps),
- pages change, so you must recrawl them on a sensible schedule,
- the work must be spread across many machines without them stepping on each other.
Interviewers commonly probe the frontier design, politeness, dedup with Bloom filters, and how to partition the crawl. A short outline exists in Design workshops; this lesson is the full design. It reuses ideas from Message queues and Storage, search and geo.
Step 1: Problem and scope
Starting from a set of seed URLs, continuously discover and download public web pages, store their content, and hand it to a search indexer, while being polite to websites and avoiding duplicate work.
In scope: seed URLs, the URL frontier with priority and politeness, robots.txt, DNS resolution, fetching, parsing and link extraction, URL normalization and dedup, content dedup, page storage, recrawl scheduling, trap avoidance, distributing the crawl, feeding an index.
Out of scope unless asked: ranking search results, rendering JavaScript-heavy pages in a full browser (we will mention it), images and video, the search query side.
Step 2: Clarifying questions
| Question | Assumed answer | Why it matters |
|---|---|---|
| What is it for? | Feeding a general web search index | Decides what to store and how fresh |
| Pages per month? | 1 billion fetches (new pages plus recrawls) | Throughput |
| Content types? | HTML only | Parser and storage size |
| Average page size? | 100 KB of HTML | Bandwidth and storage |
| Keep old versions? | Latest version only, plus a content hash | Storage |
| How fresh? | Popular pages daily, most pages weekly to monthly | Recrawl scheduling |
Respect robots.txt? | Always | Politeness and legal safety |
| JavaScript rendering? | Not for most pages | Fetcher cost |
| How long to keep data? | 5 years for capacity planning | Total storage |
| Number of known URLs? | Around 10 billion | Size of the "seen" set |
Step 3: Requirements
Functional requirements
- Accept seed URLs and crawl outward by following links.
- Obey
robots.txtand per-site crawl delays. - Download HTML, extract text and links, and store the page.
- Skip URLs already seen and pages whose content is a duplicate.
- Recrawl pages based on how often they change and how important they are.
- Publish new and changed pages to the indexing pipeline.
Non-functional requirements
- Politeness: at most one request at a time to a host, with a delay between requests. A crawler that overloads a small website is effectively a denial-of-service attack.
- Scalability: hundreds to thousands of pages per second, spread over many machines.
- Robustness: the web is full of malformed HTML, huge files, slow servers, redirect loops and malicious traps. One bad page must not crash a worker.
- Extensibility: easy to add new content types (PDFs, images) or new processing steps later.
- Freshness versus coverage: balance recrawling known pages against discovering new ones.
Step 4: Back-of-the-envelope estimates
A 30-day month has 2,592,000 seconds.
Throughput
- 1 billion pages per month ÷ 2,592,000 ≈ 386 pages per second on average.
- Peak at 2 times average: about 772 pages per second.
Bandwidth
- 386 pages per second × 100 KB ≈ 38.6 MB/s, which is about 309 Mbps. Large, but a modest fleet of fetchers with ordinary data-center links handles it.
Storage
- Raw HTML: 1 billion × 100 KB = 100 TB per month.
- HTML compresses well. Assuming 4 to 1 compression: 25 TB per month.
- If you kept every monthly crawl for 5 years: 25 TB × 60 = 1.5 PB. Keeping only the latest copy of each page (as we assumed) is far less, because recrawls overwrite.
Links
- If each page has about 50 links, parsers extract 386 × 50 ≈ 19,300 links per second. Most are already known, so the "have I seen this URL?" check runs about 19,300 times per second. That check must be fast, which leads to the Bloom filter.
The "seen URLs" set
- 10 billion known URLs. Storing an exact 16-byte hash of each needs 10 billion × 16 = 160 GB, which is too large for one machine's memory but fine on disk or sharded.
- A Bloom filter with a 1% false-positive rate needs about 9.6 bits per URL: about 12 GB for 10 billion URLs, with 7 hash functions. At 0.1% it needs about 18 GB with 10 hash functions. (The formulas are in the deep dive.)
Concurrency
- Fetches are slow: say 1 second on average including connection setup and download. To sustain 386 pages per second you need about 386 fetches in flight at once, more at peak.
- Politeness says one connection per host at a time with a delay. If you wait 2 seconds between requests to the same host, each host contributes at most 0.5 pages per second, so you need at least 772 different hosts active at the same moment. A crawl that only has a few hosts in its queue cannot go fast, no matter how many machines you add.
What the numbers told us
Throughput is moderate (hundreds of pages per second), but storage is large (tens of TB per month) and the dedup check is hot (about 19,000 per second). Politeness, not hardware, limits speed: you need many hosts in flight at once.
Step 5: API design
A crawler has no public API, but it has internal interfaces between components. Describing them shows you know where the boundaries are.
Frontier
add(urls[], priority, discoveredFrom)
next(workerId) -> url, host, notBefore (respects politeness)
done(url, result: ok | not_modified | error, retryAfter?)
Admin / control plane
POST /seeds { "urls": [...], "priority": "high" }
POST /blocklist { "host": "example.com" }
GET /stats pages/s, errors, queue sizes per host
Output event (to the indexer, via a queue)
page_fetched { url, canonicalUrl, fetchedAt, httpStatus,
contentHash, simhash, storageKey, outLinks[] }
Step 6: Data model and storage choice
url_state (key-value store, sharded by host)
url_hash -> url, host, status, last_fetched, next_fetch,
etag, last_modified, content_hash,
change_rate, priority, fail_count
host_state (in the frontier's memory, persisted)
host -> ip, ip_expires, robots_rules, robots_expires,
crawl_delay, next_allowed_time, error_rate
page_store (object storage or a distributed file system)
storage_key -> compressed HTML, response headers
content_fingerprints (key-value store)
content_hash -> canonical url (exact duplicates)
simhash bands -> urls (near duplicates)
Why these stores:
- URL state is billions of small records accessed by key, so a horizontally partitioned key-value or wide-column store (Cassandra, HBase, Bigtable-style) fits. Partitioning by host keeps a site's URLs together, which matches how the frontier works.
- Page content is large, written once and read sequentially by the indexer, so object storage (or HDFS, the Hadoop distributed file system) is cheap and durable. Storing many small pages together in large files reduces per-object overhead; the WARC format, a standard archive format for web crawls, does exactly that.
- Host state is small (millions of hosts) and read constantly, so it lives in memory on the frontier nodes, with periodic snapshots.
Step 7: High-level design
seed URLs
|
v
+---------------------------------+
| URL FRONTIER |
| prioritizer -> per-host queues |
| -> politeness scheduler |
+---------------+-----------------+
| next URL (host is allowed now)
v
+---------------+-----------------+ +--------------+
| FETCHERS |<---->| DNS resolver |
| robots.txt check, HTTP GET | | + cache |
+---------------+-----------------+ +--------------+
| raw HTML
v
+---------------+-----------------+ +--------------+
| PARSER / CONTENT PROCESSOR |----->| Page store |
| extract text + links, | | (object |
| content hash + simhash dedup | | storage) |
+-------+-----------------+-------+ +--------------+
| links | page_fetched event
v v
+-------+--------+ +----+-----------+
| URL filter + | | Indexing queue |--> search indexer
| normalizer + | +----------------+
| "seen?" check |
| (Bloom filter) |
+-------+--------+
| new URLs only
+-------------> back to URL FRONTIER
Recrawl scheduler: reads url_state, re-adds URLs
whose next_fetch time has arrived
The design is a loop with a queue at its center. Each stage scales independently: fetchers are network-bound, parsers are CPU-bound, and the frontier is memory- and coordination-bound.
Step 8: Request flows
Crawl path for one URL
- The frontier picks
https://example.com/blog/post-1, because its host's politeness timer has expired and it is the highest-priority URL for that host. - The fetcher checks the host's cached
robots.txtrules. If they are missing or older than about a day, it fetcheshttps://example.com/robots.txtfirst. If the path is disallowed, it marks the URL as blocked and stops. - It resolves
example.comto an IP address through the local DNS cache. - It sends
GETwith a clear user-agent string (identifying the crawler and a contact URL), and, for a recrawl, the conditional headersIf-None-Match: <etag>andIf-Modified-Since: <date>. - Outcomes:
304 Not Modified: the page has not changed. Updatelast_fetched, reschedule, done.200 OK: continue. Enforce a size limit (say 10 MB) and a timeout (say 30 seconds).301/302: treat the target as a new URL (through normalization and dedup); limit redirect chains to about 5 hops.429or503: the site is overloaded. Back off on the whole host and honorRetry-After.404/410: mark the URL dead; recheck rarely.
- It tells the frontier the host is free; the frontier sets the host's next allowed time to now plus its crawl delay.
Processing path
- The parser decodes the HTML (handling character encodings), extracts visible text, the title, the
<link rel="canonical">hint,noindex/nofollowmeta tags, and allhreflinks. - It computes a hash of the text (exact duplicate check) and a simhash (near-duplicate check). If the content duplicates a page already stored, it records the URL as an alias of the canonical one and does not store or index the content again.
- Otherwise it stores the compressed HTML in the page store and updates
url_state. - It publishes a
page_fetchedevent to the indexing queue. - Each extracted link is normalized, filtered (wrong scheme, blocked host, too long, looks like a trap), checked against the seen set, and new ones are added to the frontier with a priority.
Step 9: Deep dives
Deep dive 1: The URL frontier (priority and politeness)
The frontier answers one question millions of times a day: which URL should be fetched next? It has two goals that pull in different directions:
- Priority: fetch important pages first (popular sites, pages linked from many places, pages that change often).
- Politeness: never send more than one request at a time to a host, and wait between requests.
A single first-in, first-out queue fails both. A page with 500 links to its own site would put 500 URLs for one host in a row, and fetchers would either hammer that host or sit idle waiting.
The classic solution (described in the Mercator crawler design and used in many textbooks) uses two layers of queues:
new URLs
|
v
+----------------+
| Prioritizer | assigns priority 1..F
+--+----+----+---+
| | |
v v v
[F1] [F2] [F3] front queues: one per priority level
| | |
+----+----+
| biased selector: picks high-priority
v queues more often
+----------------+
| Router | host -> back queue (via table)
+--+----+----+---+
| | |
v v v
[B1] [B2] [B3] ... back queues: each holds URLs of ONE host
| | |
+----+----+
|
+----------------+
| Min-heap of | (next_allowed_time, back queue)
| host timers |
+-------+--------+
| pop the host whose time has come
v
fetcher
- Front queues handle priority. The prioritizer gives each URL a score and places it in one of F queues. A selector picks from high-priority queues more often than low ones, but not exclusively, so low-priority URLs are not starved forever.
- Back queues handle politeness. Each back queue holds URLs for exactly one host. A mapping table says which back queue serves which host.
- A min-heap (a priority queue that always returns the smallest item) holds one entry per back queue: the earliest time that host may be contacted again. A fetcher pops the top entry, waits if the time has not come, takes one URL from that host's queue, and after the fetch pushes the host back with
now + delay.
The delay can be fixed (say 1 to 2 seconds), taken from robots.txt's Crawl-delay directive if present (not all crawlers honor it, but a polite one should), or adaptive: some crawlers wait a multiple of how long the last fetch took, so slow servers automatically get gentler treatment.
Here is a simplified, single-process version of the back-queue scheduler:
import heapq
from collections import deque
class Frontier:
"""Per-host FIFO queues plus a heap ordered by when each host may be
fetched next. Simplified: one process, no priorities."""
def __init__(self, delay_seconds: float = 1.0):
self.delay = delay_seconds
self.queues = {} # host -> deque of URLs
self.ready = [] # heap of (next_allowed_time, host)
def add(self, host: str, url: str) -> None:
if host not in self.queues:
self.queues[host] = deque()
heapq.heappush(self.ready, (0.0, host))
self.queues[host].append(url)
def next_url(self, now: float):
while self.ready and self.ready[0][0] <= now:
_, host = heapq.heappop(self.ready)
queue = self.queues[host]
if queue:
url = queue.popleft()
heapq.heappush(self.ready, (now + self.delay, host))
return url
del self.queues[host] # host drained; re-added on next add()
return None
f = Frontier(delay_seconds=1.0)
for u in ["a.com/1", "a.com/2", "a.com/3"]:
f.add("a.com", u)
f.add("b.org", "b.org/1")
print([f.next_url(now=0.0) for _ in range(3)]) # ['a.com/1', 'b.org/1', None]
print(f.next_url(now=1.0), f.next_url(now=1.0)) # a.com/2 None
Walk through the output: at time 0, the first call returns a.com/1 and schedules a.com for time 1; the second returns b.org/1; the third finds no host allowed yet, so it returns None, even though a.com still has two URLs waiting. At time 1, a.com/2 is released, and a second call at the same instant gets nothing. That is politeness in action.
Frontier size. With billions of URLs, the queues do not fit in memory. Keep the head and tail of each queue in memory and the middle on disk, reading and writing in large batches.
Priority signals. Common inputs are: how many known pages link to the URL (a rough importance estimate, related to PageRank), the site's overall quality, how close the URL is to a seed (link depth), how often the page has changed in the past, and whether it is new (never crawled).
Interview tip
Draw the two layers explicitly and say what each solves: "Front queues are for priority, back queues are for politeness, one host per back queue, and a heap of host timers tells fetchers which host is next." That one sentence is what interviewers listen for.
Deep dive 2: robots.txt and DNS caching
robots.txt is a plain text file at the root of a site (https://example.com/robots.txt) that tells crawlers which paths they may fetch. The rules are standardized in RFC 9309 (the Robots Exclusion Protocol).
User-agent: *
Disallow: /admin/
Disallow: /search
Allow: /search/help
User-agent: BadBot
Disallow: /
Sitemap: https://example.com/sitemap.xml
How a polite crawler uses it:
- Fetch it before crawling a host, cache the parsed rules per host, and refresh them about daily (RFC 9309 suggests not caching for more than 24 hours in general).
- Find the group that matches your user-agent (or
*), then apply the most specific (longest) matching rule. Here/search/helpis allowed even though/searchis disallowed. - If
robots.txtreturns 404, treat everything as allowed. If it returns a server error (5xx), treat the whole site as disallowed for now, because the site may be in trouble. - Use the
Sitemaplines to discover URLs the site wants crawled. - Inside pages, respect
<meta name="robots" content="noindex, nofollow">: do not index the page, do not follow its links.
DNS caching. Every fetch needs the host's IP address. A DNS lookup can take tens to hundreds of milliseconds, and at hundreds of fetches per second, the default resolver on each machine becomes a bottleneck. Production crawlers:
- run their own caching resolver on each fetcher node (or a shared cluster),
- cache results per host, respecting the DNS TTL (time-to-live, how long the answer may be reused), with a minimum and maximum bound,
- resolve asynchronously (many lookups in flight) rather than blocking a thread per lookup,
- prefetch: when a new host enters the frontier, resolve it before its first URL is due.
Caching the IP also helps politeness: several hostnames (a.example.com, b.example.com) may live on one server, and some crawlers apply politeness per IP address as well as per host.
Deep dive 3: URL normalization and deduplication
The web has two kinds of duplicates: many URLs for one page, and one page's content at many URLs. You handle them at different stages.
URL normalization
Before checking "have I seen this URL?", turn every URL into a single canonical form:
- resolve relative links against the page's URL (
../shoponhttps://example.com/blog/becomeshttps://example.com/shop), - lowercase the scheme and host (paths are case-sensitive, so leave them),
- remove default ports (
:80for http,:443for https), - remove the fragment (
#section), which never reaches the server, - sort query parameters and remove known tracking parameters such as
utm_source, - optionally apply site-specific rules learned over time (such as removing session id parameters).
from urllib.parse import urlsplit, urlunsplit, parse_qsl, urlencode, urljoin
TRACKING = {"utm_source", "utm_medium", "utm_campaign", "fbclid", "gclid"}
def normalize(base: str, href: str) -> str:
url = urljoin(base, href) # resolve relative links
parts = urlsplit(url)
scheme = parts.scheme.lower()
host = (parts.hostname or "").lower()
port = parts.port
if (scheme, port) in (("http", 80), ("https", 443)):
port = None # drop default ports
netloc = host if port is None else f"{host}:{port}"
path = parts.path or "/"
query = sorted((k, v) for k, v in parse_qsl(parts.query)
if k not in TRACKING)
return urlunsplit((scheme, netloc, path, urlencode(query), ""))
print(normalize("https://Example.com/blog/", "../shop?b=2&a=1&utm_source=x#top"))
# https://example.com/shop?a=1&b=2
print(normalize("http://EXAMPLE.com:80/", "/index.html"))
# http://example.com/index.html
Also honor the page's <link rel="canonical" href="..."> hint, which tells you the site's preferred URL for this content.
The seen-URL check with a Bloom filter
About 19,300 links per second need a "seen before?" check against 10 billion URLs. An exact set of hashes is 160 GB. A Bloom filter is a compact probabilistic set: a bit array of m bits and k hash functions.
- Add a URL: compute k hashes, set those k bits to 1.
- Check a URL: compute the k hashes; if any of those bits is 0, the URL is definitely new. If all are 1, it is probably seen.
It never gives a false negative (it never says "new" for a URL it has seen), but it can give a false positive (say "seen" for a new URL), because other URLs may have set those bits.
Sizing formulas for n items and false-positive rate p:
m = -n * ln(p) / (ln 2)^2 bits
k = (m / n) * ln 2 hash functions
n = 10,000,000,000, p = 0.01:
m = 10^10 * 4.605 / 0.4805 = 9.585 * 10^10 bits = 11.98 GB
k = 9.585 * 0.693 = 6.64 -> use 7
n = 10,000,000,000, p = 0.001:
m = 1.438 * 10^11 bits = 17.97 GB, k = 9.97 -> use 10
So a 1% filter takes about 12 GB, which can be sharded across frontier nodes by host. What does the 1% error cost? One in a hundred genuinely new URLs is wrongly skipped. For a search crawler, that is acceptable: important pages are linked from many places, so they will be found through another link or a sitemap. If it is not acceptable, use the Bloom filter as a fast first check, and confirm "probably seen" answers against the exact store on disk.
Common mistake
A Bloom filter cannot delete items (clearing a bit could remove other URLs too), and its false-positive rate climbs as you add more items than it was sized for. Size it for expected growth, rebuild it periodically, or use variants such as counting Bloom filters or cuckoo filters if you need deletes.
Content dedup: exact and near duplicates
Different URLs often serve the same content: mirrors, http and https versions, printer-friendly pages, syndicated articles. Studies of web crawls have long found that a substantial share of pages are duplicates or near duplicates.
- Exact duplicates: hash the extracted text (for example with SHA-256 or a fast 64-bit hash) and look it up in
content_fingerprints. A match means "same content as URL X"; store a pointer, not the content. - Near duplicates: two pages differ only by a date, an ad or a "you may also like" box. Their exact hashes differ completely. Simhash (by Moses Charikar) produces a 64-bit fingerprint in which similar documents get fingerprints that differ in only a few bits.
How simhash works, step by step:
- Split the text into features (words, or short word sequences called shingles).
- Hash each feature to 64 bits.
- Keep 64 counters. For each feature, for each bit position, add 1 if the feature's hash has a 1 there, otherwise subtract 1. (Features can be weighted, for example by frequency.)
- The final fingerprint has a 1 in each position where the counter is positive.
- Compare two pages by Hamming distance: the number of bit positions where their fingerprints differ.
import hashlib
def simhash(text: str, bits: int = 64) -> int:
weights = [0] * bits
for word in text.lower().split():
h = int.from_bytes(hashlib.md5(word.encode()).digest()[:8], "big")
for i in range(bits):
weights[i] += 1 if (h >> i) & 1 else -1
return sum(1 << i for i in range(bits) if weights[i] > 0)
def hamming(a: int, b: int) -> int:
return bin(a ^ b).count("1")
page_a = "cheap flights to goa book now best prices on flights to goa " * 3
page_b = page_a + "updated 10 october"
page_c = "python tutorial for beginners learn lists loops and functions " * 3
print(hamming(simhash(page_a), simhash(page_b))) # small (a few bits)
print(hamming(simhash(page_a), simhash(page_c))) # large (around 32)
Unrelated texts differ in about half of the 64 bits; near copies differ in only a few. A common rule treats fingerprints within 3 bits as near duplicates. To find candidates without comparing against every stored fingerprint, split the 64 bits into blocks and index each block: two fingerprints within 3 bits must agree exactly on at least one of 4 blocks of 16 bits, so you look up only pages sharing a block.
Deep dive 4: Recrawl scheduling, traps and distribution
Recrawl scheduling
Pages change at very different rates: a news homepage changes every few minutes, a 2015 blog post almost never. Recrawling everything at the same rate wastes most fetches on unchanged pages and still misses fast-changing ones.
- Track change history: on each fetch, record whether the content hash changed. Estimate a change rate per page.
- Adapt the interval: if the page changed since last time, shorten its interval (for example halve it, with a minimum such as 1 hour). If it did not, lengthen it (for example multiply by 1.5, with a maximum such as 30 days).
- Weight by importance: an important page that changes weekly deserves more frequent checks than an obscure one.
- Use cheap checks: conditional requests with
ETagandLast-Modifiedgive a304 Not Modifiedwith no body when nothing changed, saving bandwidth for both sides. - Use hints: sitemaps with
lastmoddates, RSS and Atom feeds for news sites.
The recrawl scheduler periodically scans url_state for URLs whose next_fetch has passed and pushes them into the frontier with their priority, competing with newly discovered URLs. The share of capacity reserved for recrawls versus discovery is a tunable policy.
Crawler traps
A crawler trap is a site structure that produces endless URLs, accidentally or on purpose:
- calendars with "next month" links forever (
/calendar?month=2099-12, then 2100-01, ...), - relative-link loops that grow the path (
/a/b/a/b/a/b/...), - session ids in every URL, so every visit looks new,
- faceted search producing endless filter combinations (
?color=red&size=9&sort=price...), - deliberately generated infinite pages.
Defenses, layered:
- maximum URL length (for example 2,000 characters) and maximum path depth,
- detect repeating path segments (
/a/b/a/b), - per-host page budget: no host gets more than a set number of pages per crawl cycle unless it earns more through importance,
- strip known session parameters during normalization,
- near-duplicate content detection: many URLs producing near-identical pages is a strong trap signal,
- monitoring: alert when one host's queue grows unusually fast, and let operators blocklist it.
Distributing the crawl
One machine cannot do all this. Spread the work across N crawler nodes, each running a frontier, fetchers and parsers. The key decision is how to split URLs among nodes.
Partition by host (hash of the hostname, or better the registered domain, modulo N, with consistent hashing so nodes can be added):
- all URLs of one host live on one node, so politeness is enforced locally, with no cross-node coordination,
- robots rules and DNS results for a host are cached on one node,
- links to other hosts are forwarded to the owning node, in batches to reduce messaging.
node 0 owns hash(host) in range A node 1 owns range B
+---------------------+ +---------------------+
| frontier (A hosts) | links to B | frontier (B hosts) |
| fetchers, parsers |------------->| fetchers, parsers |
| Bloom shard A |<-------------| Bloom shard B |
+---------------------+ links to A +---------------------+
Partitioning by URL hash instead would spread one host's URLs across all nodes, forcing them to coordinate politeness for every request, which is why host partitioning is the standard answer. The downside is skew: a giant site lands on one node. Mitigate by splitting very large hosts across a few nodes with a shared rate limit, or by giving them a higher per-host budget on a stronger node.
Geography. Crawling from data centers near the sites (Asia, Europe, the Americas) cuts latency. Partitioning can include region so that, for example, Indian sites are fetched from an Indian region.
Feeding a search index
The crawler's output is the page_fetched event stream. The indexer consumes it, tokenizes the text, and updates an inverted index (word to list of documents). It also uses the extracted link graph to compute importance scores, which flow back into the crawler's priorities: important pages are crawled more often, and fresh crawls make the index better. The query side of search is covered in Design a search query cache, and the index itself in Design a Twitter timeline.
Step 10: Scaling and bottlenecks
- Fetchers are network-bound. Use asynchronous I/O (one process handling thousands of connections) rather than one thread per connection. Add nodes to add throughput, as long as the frontier has enough distinct hosts ready.
- Parsers are CPU-bound. Run them as a separate pool fed by a queue so slow parsing never blocks fetching.
- Frontier memory is the hardest to scale. Keep queue heads in memory and the rest on disk; shard by host.
- Seen-URL set: Bloom filter shards in memory, exact store on disk.
- DNS: local caching resolvers with asynchronous lookups.
- Page store: object storage scales almost without limit; batch small pages into large files.
- JavaScript-heavy pages: rendering with a headless browser is 10 to 100 times more expensive than a plain fetch, so do it only for important pages that need it, in a separate pool.
Step 11: Failure handling
| Failure | Effect | Mitigation |
|---|---|---|
| Crawler node crashes | Its hosts stop being crawled | Frontier checkpointed to disk; consistent hashing reassigns its hosts |
| Fetch timeout, slow server | Worker stuck | Strict connect and read timeouts; size limits |
| Site returns 429/503 | We are overloading it | Exponential backoff for the whole host; honor Retry-After |
| Malformed or huge HTML | Parser crash or memory blowup | Lenient parser, size cap, isolate parsing per page |
| Crawler trap | Frontier fills with junk | URL length, depth and per-host budgets; trap detection |
| Bloom filter lost | Re-crawl of known URLs | Rebuild from url_state; dedup again on content hash |
| Indexer falls behind | Index is stale | Queue buffers events; crawl continues |
| Duplicate processing after retry | Page stored twice | Idempotent writes keyed by URL hash and fetch time |
Step 12: Trade-offs and alternatives
| Decision | Choice here | Alternative | When the alternative wins |
|---|---|---|---|
| Traversal order | Priority-based | Pure BFS (breadth-first) | Small focused crawl of one site |
| Partitioning | By host | By URL hash | No politeness concerns (crawling your own sites) |
| Seen-URL check | Bloom filter + exact store | Exact set only | Small crawl that fits in memory |
| Near-dup detection | Simhash | MinHash with shingles | You need Jaccard similarity estimates |
| Recrawl | Adaptive per page | Fixed interval | Small, uniform set of pages |
| Rendering | Plain HTML fetch | Headless browser for all | Sites that only work with JavaScript |
| Storage | Latest copy only | Every version | Web archive, change analysis |
What interviewers probe
"How do you make sure you never hit one site too hard?" One back queue per host, one fetch at a time per host, a delay between fetches tracked in a heap of host timers, backoff on 429 and 503, and respect for Crawl-delay. Partitioning by host keeps this local to one node.
"What does a Bloom filter false positive cost here?" A new URL is skipped. Important pages have many inbound links and sitemaps, so they are found anyway. If unacceptable, confirm positives against an exact store.
"BFS or DFS?" BFS from seeds finds important, well-linked pages early; DFS dives deep into a single site and gets stuck in traps. Real crawlers use priority queues, which generalize BFS.
"How do you detect that a page changed?" Conditional GET with ETag or Last-Modified, then compare the content hash. Use simhash to ignore trivial changes like timestamps.
"What if a node dies?" Its host range moves to other nodes through consistent hashing; they load the checkpointed frontier from durable storage. Some URLs may be fetched twice, which is harmless.
"How would you crawl only Indian news sites?" A focused crawler: seeds from known news sites, a classifier that scores pages and links for relevance, and priorities that favor relevant links while dropping off-topic hosts.
Interview questions
Q1. What is a URL frontier?
The data structure that holds URLs waiting to be fetched and decides the order. It combines priority (important pages first) and politeness (one request at a time per host, with delays). At scale it is sharded by host and partly stored on disk.
Q2. Why have both front queues and back queues?
Front queues order URLs by priority. Back queues each hold URLs for a single host so the crawler can enforce a delay per host. A selector moves URLs from front to back queues, and a heap of host timers tells fetchers which host may be contacted next.
Q3. How does a crawler use robots.txt?
It fetches the file before crawling a host, caches the parsed rules (refreshing about daily), and applies the most specific matching rule for its user-agent before every fetch. A 404 means everything is allowed; a server error means pause the site. It also reads sitemap links from the file.
Q4. Why cache DNS lookups?
Every fetch needs an IP address, and lookups are slow relative to the crawler's rate. A local caching, asynchronous resolver removes this bottleneck and reduces load on DNS servers. Results are cached for the record's TTL.
Q5. What is URL normalization and why is it necessary?
Turning equivalent URLs into one canonical form: resolving relative links, lowercasing scheme and host, removing default ports, fragments and tracking parameters, and sorting query parameters. Without it, the same page would be fetched under many URLs.
Q6. How does a Bloom filter work, and how big is one for 10 billion URLs?
It is a bit array plus k hash functions; adding sets k bits, checking tests them. It can say "probably seen" for a new item but never "new" for a seen one. At a 1% false-positive rate it needs about 9.6 bits per item, about 12 GB for 10 billion URLs with 7 hash functions.
Q7. How do you detect duplicate content across different URLs?
Exact duplicates by hashing the extracted text and looking it up. Near duplicates with simhash: similar pages produce 64-bit fingerprints that differ in only a few bits, so a small Hamming distance (for example 3 or less) marks a near duplicate.
Q8. How do you decide when to recrawl a page?
Estimate each page's change rate from past fetches and adapt its interval: shorter when it changed, longer when it did not, within bounds, weighted by importance. Conditional requests make unchanged checks cheap, and sitemaps and feeds give hints.
Q9. What is a crawler trap and how do you avoid one?
A site structure that generates endless URLs, such as infinite calendars or repeating paths. Limit URL length and depth, detect repeating segments, strip session parameters, cap pages per host, and watch for many URLs with near-identical content.
Q10. How do you distribute a crawler across machines?
Partition by host (or registered domain) with consistent hashing. Each node owns its hosts' frontier, politeness timers, robots rules and dedup shard, and forwards links for other hosts to their owners in batches. This avoids cross-node coordination on every fetch.
Q11. What is the crawler's throughput limited by?
Often politeness, not hardware. With a 2-second delay per host, each host gives at most 0.5 pages per second, so 772 pages per second needs at least about 772 hosts active at once. Bandwidth, DNS and parsing CPU are the next limits.
Q12. How does the crawler feed a search index?
It publishes an event for each new or changed page with a pointer to the stored content. The indexer tokenizes the text and updates the inverted index, and link-graph analysis produces importance scores that feed back into crawl priorities.
Q13. How do you handle a site that responds slowly or with errors?
Use connection and read timeouts, back off exponentially on that host for 429 and 5xx responses, honor Retry-After, and lower the host's priority if errors persist. One bad host must never tie up the fetcher pool.
Key takeaways
- The crawl loop is simple; the design is about the frontier, politeness, dedup, traps and distribution.
- Under our assumptions: about 386 pages per second, 309 Mbps, 25 TB of compressed HTML per month and about 19,000 seen-URL checks per second.
- Front queues give priority, back queues (one per host) plus a heap of host timers give politeness.
- Always obey robots.txt and back off on 429 and 503; politeness limits throughput more than hardware does.
- Normalize URLs before dedup; use a Bloom filter (about 12 GB for 10 billion URLs at 1%) for the fast seen check.
- Use content hashes for exact duplicates and simhash with Hamming distance for near duplicates.
- Recrawl adaptively by change rate and importance, using conditional GETs.
- Partition by host with consistent hashing so politeness stays local to one node.
Next lesson
Continue with Design a search query cache.

