SYSTEM DESIGN #11 · INTERVIEW GUIDE
Distributed Web Crawler
The web is 5.5 billion pages and growing — and your crawler needs to visit the right ones first, respect rate limits, survive machine failures mid-crawl, and never visit the same URL twice. This Distributed Web Crawler covers every hard problem in the space: a priority-sorted URL frontier in Redis Sorted Sets, per-domain politeness queues backed by token buckets, Bloom filters for near-O(1) URL deduplication at 1B+ URL scale, and a Kubernetes fetcher fleet that scales horizontally without coordination overhead. The system handles JavaScript-rendered pages with headless Chromium, respects robots.txt and Crawl-Delay directives, and stores extracted content in S3 Parquet for downstream indexing pipelines. You will leave knowing how Google, CommonCrawl, and Ahrefs crawl at billions of pages per day without getting banned or crashing targets.
Redis ZADDBloom FilterKafkaPlaywrightBFS
💡
The Gist — What Problem Are We Solving?
A librarian exploring an infinite library — forever
A web crawler starts from a list of known URLs (seeds), fetches each page, extracts all links found inside, adds them to a queue, and repeats — indefinitely and at massive scale. The challenges: never hammering the same website twice in quick succession (politeness), never crawling the same page twice (deduplication), and handling infinite URL traps that malicious sites generate to confuse crawlers.
💬Think of it as a robot librarian that reads every book in an infinite library, notes all cross-references, and files them — without visiting the same shelf twice and without bothering any one library more than once per second.
These are the capabilities the system must deliver — what users and operators can actually do with it.
🌱
Seed Management
🌱Accept seed URLs via API; prioritise by PageRank estimate and freshness
⬇️
Fetching
⬇️Respect robots.txt; crawl delay; user-agent identification; handle redirects
🔍
Content Parsing
🔍Extract text, links, metadata; detect encoding and MIME type; parse sitemaps
🔁
Deduplication
🔁Exact dedup via Bloom filter; near-dup via SimHash; URL normalisation
💾
Storage
💾Raw HTML to S3; parsed content and metadata to distributed KV store
⚡
Non-Functional Requirements
These define how well the system must perform — the quality attributes that separate a toy from a production system.
⚡ Throughput
⚡1B pages/day (~11K pages/sec)
⏱️ Politeness
⏱️Max 1 req/sec per domain; respect Crawl-Delay from robots.txt
🔁 Dedup Window
🔁24h exact dedup; near-dup across full corpus
🛡️ Robustness
🛡️Handle spider traps, redirect loops, malformed HTML
📅 Recrawl
📅Adaptive recrawl based on observed change frequency
📊
Key Metrics — The Numbers That Define This System
The headline numbers to know cold — and be ready to explain how each one is achieved.
🏗️
System Architecture Diagram
Full data flow from source to serving. Each layer scales independently.
Ingestion
Distributed Web Crawler Flow
Seed URLs
→
URL Frontier
Redis sorted set per domain
→
→
↓
↓
Storage
Content Store
S3) + Link Extractor
→
🗺️
End-to-End User Journey
Trace a single request end-to-end — the story interviewers want you to tell fluently.
1
Seed URL submitted
— API accepts seed URL; DNS resolved; robots.txt fetched and cached; URL added to priority queue
2
Fetcher picks URL
— Worker dequeues URL; checks per-domain politeness token bucket; waits if needed
3
Page fetched
— HTTP GET with crawler user-agent; follows up to 5 redirects; 30s timeout
4
Content processed
— HTML parsed; links extracted and normalised; Bloom filter check for exact dedup; SimHash for near-dup
5
New links enqueued
— Novel URLs added to frontier with priority score; raw HTML stored to S3; parsed content to KV store
🔭
High-Level Design — Component Breakdown
Core components — each with a single, well-defined responsibility. The key architectural insight: each layer scales independently, and failure in one component is isolated from the rest.
1 — Seed API
Handles responsibilities for the Seed API layer. Designed for independent horizontal scaling — additional instances added without architectural changes. Communicates asynchronously with adjacent components to maximise throughput and fault isolation.
2 — URL Frontier
Handles responsibilities for the URL Frontier layer. Designed for independent horizontal scaling — additional instances added without architectural changes. Communicates asynchronously with adjacent components to maximise throughput and fault isolation.
3 — Fetcher Fleet
Handles responsibilities for the Fetcher Fleet layer. Designed for independent horizontal scaling — additional instances added without architectural changes. Communicates asynchronously with adjacent components to maximise throughput and fault isolation.
4 — Parser
Handles responsibilities for the Parser layer. Designed for independent horizontal scaling — additional instances added without architectural changes. Communicates asynchronously with adjacent components to maximise throughput and fault isolation.
5 — Bloom Filter
Handles responsibilities for the Bloom Filter layer. Designed for independent horizontal scaling — additional instances added without architectural changes. Communicates asynchronously with adjacent components to maximise throughput and fault isolation.
6 — S3 + KV
Object storage for raw events (Iceberg/Parquet), media assets, ML model artefacts, and snapshots. Lifecycle rules tier data to Glacier after 90 days. Versioning disabled on ephemeral buckets (snaps, stories) to ensure hard deletes.
🔬
Low-Level Design — Deep Dives
Deep dives worth explaining in detail in any senior engineering interview. For each: know the data structure, the algorithm, the why, and the trade-off you made.
1 — URL Frontier
Redis ZADD · Priority Queue
URL frontier implemented as Redis Sorted Set: ZADD frontier score url_hash, where score = -priority (lowest score = highest priority). ZPOPMIN(100) dequeues 100 highest-priority URLs per worker per second. Priority score = inlink_count × freshness_multiplier × (1/age_days). New URLs from seed API given priority=0.5 (medium). Known high-value domains (top 1M Tranco list) given priority=1.0. Domain throttling enforced via per-domain queue with Crawl-Delay respected.
# Enqueue with priority
redis.zadd(‘frontier’, {url_hash: -priority})
# Dequeue batch
urls = redis.zpopmin(‘frontier’, 100)
# Politeness: per-domain last-crawl timestamp
if redis.get(f’crawled:{domain}’) > time.time() – crawl_delay:
redis.zadd(‘frontier’, {url_hash: -(priority-0.1)})
2 — Bloom Filter Deduplicator
10B bits · 0.008% FP
Counting Bloom filter with 10B bits (1.25GB RAM), 7 hash functions, <0.008% false positive rate at 1B URLs. Filter persisted to S3 every 15min via memory-mapped file serialisation. On worker restart, filter loaded from latest S3 snapshot. Distributed crawl: per-shard Bloom filters merged daily — cross-shard sync via S3. False positives (URL not-yet-crawled seen as crawled) acceptable — tiny fraction of pages missed. No false negatives possible by construction.
from pybloom_live import ScalableBloomFilter
filter = ScalableBloomFilter(mode=ScalableBloomFilter.LARGE_SET_GROWTH)
def is_duplicate(url_hash):
if url_hash in filter:
return True
filter.add(url_hash)
return False
3 — Politeness Queue
Per-domain Rate Limit
Each domain gets an independent token bucket controlling crawl rate. robots.txt Crawl-Delay directive sets the minimum inter-request interval. Default: 1 request/10 seconds per domain. Top-100 domains with explicit crawl permission: 1 req/sec. Token bucket state per domain stored in Redis (key: politeness:{domain}). Crawler worker checks token before fetching; if no token available, URL re-enqueued with a future score in the priority queue.
def can_crawl(domain):
key = f’politeness:{domain}’
delay = get_crawl_delay(domain) # from robots.txt cache
last = float(redis.get(key) or 0)
if time.time() – last >= delay:
redis.set(key, time.time(), ex=86400)
return True
return False
4 — Parser & Link Extractor
BeautifulSoup · Link Graph
HTML parser extracts: outgoing links (href, src, action attributes), page metadata (title, meta description, canonical URL, OG tags), structured data (JSON-LD, microdata for content classification). Links normalised (relative → absolute, query string canonicalisation, fragment removal). Extracted content stored in S3 as NDJSON (one file per domain per hour). Link graph edges written to Kafka link_graph topic for downstream PageRank computation.
def extract_links(html, base_url):
soup = BeautifulSoup(html, ‘lxml’)
links = set()
for tag in soup.find_all([‘a’,’link’], href=True):
url = urljoin(base_url, tag[‘href’])
if is_valid_url(url):
links.add(canonicalise(url))
return links
⚖️
Trade-offs & Decision Log
Every senior interview comes down to these decisions. Know the exact trade-off, the reasoning, and the specific numbers that justify each choice.
⚖️ BFS vs Priority-Based URL Frontier
✓
Priority Queue (PageRank score) ✅ Chosen
- High-value pages crawled first — better index freshness
- Redis ZADD gives O(log N) priority insertion and O(1) pop
- Crawl budget spent on pages that matter most
- Priority score requires initial seed data to bootstrap
→
BFS (FIFO queue)
- Simple — first discovered is first crawled
- No priority scoring overhead
- Wastes crawl budget on low-value pages (error pages, archives)
- Deep graph traversal hits unimportant content quickly
💡Decision: Priority queue scored by inlink count + freshness signal; FIFO fallback for cold-start domains with no score history
⚖️ Centralised Frontier vs Distributed Frontier
✓
Centralised Redis Frontier ✅ Chosen
- Single source of truth — no duplicate URL assignment
- Redis Sorted Sets give sub-ms enqueue/dequeue at 10M URLs
- Politeness enforced globally — no cross-node coordination needed
- Redis becomes bottleneck above ~50K URL ops/sec
→
Distributed Per-Worker Frontier
- No single bottleneck — scales linearly with workers
- Risk of duplicate crawls without bloom filter coordination
- Politeness enforcement requires cross-worker communication
- Complex consistent hash ring management
💡Decision: Centralised Redis frontier up to 50K URLs/sec; partition frontier by domain hash for larger crawls
🎯Interview Questions — Answered
The exact questions interviewers ask — with production-grade answers
Q1
How does the crawler respect robots.txt at scale?
robots.txt is fetched and cached per domain at crawl start and refreshed every 24 hours. The politeness queue uses the robots.txt Crawl-delay directive as the minimum inter-request delay for that domain. A custom robots.txt parser handles non-standard extensions (Allow: patterns, Sitemap: directives, crawl-delay with decimal values). Disallowed paths are checked against the URL before it enters the frontier — never fetched. User-Agent in robots.txt is matched to the crawler’s declared agent (‘CloudWizardBot/1.0’). If robots.txt fetch fails (404 or timeout), the crawler defaults to a conservative 10-second crawl delay.
Q2
How does the Bloom filter handle URL deduplication at 1B+ URL scale?
A counting Bloom filter with 10 billion bits (1.25GB RAM) and 7 hash functions gives a false positive rate of 0.008% at 1B URLs. False positives (seen as ‘already crawled’) are acceptable — a tiny fraction of pages revisited slightly less often than optimal. False negatives are impossible by construction (Bloom filters never say ‘not seen’ for URLs that have been seen). The filter is persisted to S3 every 15 minutes and loaded on worker restart. For distributed crawls, a Bloom filter per-shard with periodic cross-shard sync via S3 ensures global deduplication with eventual consistency.
Q3
How is crawl freshness balanced against crawl budget?
Adaptive crawl scheduling: page re-crawl interval = base_interval × (1 + content_change_rate)^(-1). Pages that change frequently (news, inventory) get shorter intervals. Pages with stable content (about pages, static docs) get longer intervals. Change detection: SHA-256 of HTML fingerprint compared to last crawl. If fingerprint unchanged, interval is doubled (up to max 30 days). If changed, interval is halved (down to min 1 hour). Budget allocation: 60% to high-priority (new URLs), 30% to frequently-changing pages, 10% to re-validation of stable pages. Total budget managed as a token bucket: N pages/hour, refilled each hour.
System Design Series · Every Tuesday & Thursday
Level up your system design interviews
Each post covers Gist, Functional & Non-Functional Requirements, Key Metrics, System Diagram, User Journey, HLD, LLD, and Trade-offs & FAQs.
Subscribe to never miss a post →
Previous Articles
Categories: System Design
Tags: bfs, bloom filter, distributed systems, interview prep, politeness, system design, web crawler
Leave a Reply