Web Crawler
Worked solution

Web Crawler — a worked solution

Politely traverse the web at scale. Don't crawl yourself in circles.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the Web Crawler workspace

The problem

Build a production-grade web crawler that traverses the open web at hundreds of millions of pages per day, respects every host's politeness budget and robots.txt, never crawls itself in circles, and feeds a downstream search index without becoming the bug that pages someone else's on-call.

The hard problems are not throughput. The hard problems are: politeness (one bad config rolls out, you DDoS Wikipedia and your IP range is blocked within an hour), de-duplication (without a seen-set you crawl forever), traps (calendar widgets and faceted search will eat 30% of your fetch budget if you let them), and making the multi-store page commit atomic enough that a parser crash doesn't leave the system with a URL marked seen but never enqueued.

This canonical sizes for the Mid tier: 10 B URLs/month over 100 M hosts (~3.86 K sustained / 7.7 K peak fetch QPS, ~250 TB/month compressed WARC ingest). Numbers are derived in the Capacity section, not assumed.

The reference architecture

Reference architecture for Web Crawler: 16 components — Operator / Admin, CDN, API Gateway, Scheduler / Politeness, Frontier (host subqueues), Seen-set (Bloom + Cassandra), Fetcher Fleet, robots.txt Cache, Renderer (Headless Chromium), Target Web (origins), Fetched Pages Stream, Parser / Extractor, WARC Archive, Search Index, Pages DLQ, DNS Upstream — connected by 20 flows.Operator / AdminclientCDNCloudflareAPI GatewayEnvoy + WAFScheduler / PolitenessMercator-style back-que…Frontier (host subque…Cassandra (24-node clus…Seen-set (Bloom + Cas…Bloom bitarray sidecar …Fetcher FleetOkHttp 5 + co-located U…robots.txt CacheRedis 7 cluster + googl…Renderer (Headless Ch…Headless Chromium pool …Target Web (origins)externalFetched Pages StreamKafka (512 partitions, …Parser / ExtractorJVM workers (Tika + jso…WARC ArchiveS3 (Standard + Glacier …Search IndexOpenSearch (32 primarie…Pages DLQKafka compacted DLQ (RF…DNS UpstreamTiered: Quad9 → Cloudfl…
16 components, 20 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Operator / Admin

Internal users (SREs, ops, legal) submit seed URLs, takedown requests, and read freshness/coverage dashboards. Off the crawl path; only ~100s of QPS even at peak.

Why it exists. Crawl pipelines aren't autonomous — they need a control surface. Seeds enter the system here, takedowns enter here, and the audit trail of operator actions is the legal-compliance record. Without an explicit operator surface every emergency action becomes a deploy.

When it fails. Operator surface down → no new seeds, no takedowns, no dashboards. Crawl pipeline keeps running on its existing seed list — degrades gracefully. Legal-mandated takedowns become the load-bearing concern: runbook says page legal + ops directly and push the takedown via the emergency control plane.

CDNCloudflare

Fronts the operator dashboard + the public WARC index mirror. Caches static dashboard assets (24 h) and the public crawler-identity endpoint (60 s). The crawl loop itself does NOT traverse the CDN.

Why it exists. The public surfaces (dashboard, /.well-known/crawler-ip-ranges) are vulnerable to volumetric attacks. Anycast absorbs DDoS at the edge; the origin sees only L7 misses. Without it, every operator dashboard refresh hits the gateway directly.

When it fails. Edge degrades → operators see stale dashboards (acceptable). Edge fails completely → public crawler-identity endpoint goes dark, which means target-site SREs trying to verify our bot identity see 5xx and may blanket-ban us. Mitigation: identity surface is also reachable via raw-IP allowlist published out-of-band.

API GatewayEnvoy + WAF

Authenticates operators (mTLS upstream, bearer + ticket-system reference inbound), terminates TLS, enforces per-operator rate limits (10 RPS), routes /api/v1/* to the scheduler and the takedown propagation channel to the robots cache. Hosts the public crawler-identity endpoint.

Why it exists. Single-tenant policy enforcement point: WAF rules, auth, rate limits, audit logging all live here. Without a gateway every backing service grows its own auth — an audit nightmare and a security hole.

When it fails. Gateway down → no new seeds, no takedowns, no dashboard. Crawl pipeline keeps running but legal-mandated actions can't propagate. Multi-AZ + blue/green deploy minimizes blast radius; an emergency control plane bypasses the gateway for catastrophic legal escalations.

Scheduler / PolitenessMercator-style back-queue + multi-tier token bucket

Runs the Mercator-style back-queue dispatch loop: pop next ready host from the frontier, check politeness windows (host, /24, AS), assign URL to a fetcher pod via grpc lease. Sharded by hash(host); each shard is leader-elected via etcd. Holds politeness state in a centralized Redis cluster shared across all 12 shards so cross-shard AS budgets coordinate correctly.

Why it exists. Without an explicit scheduler, fetchers either pull randomly (politeness violations) or push-from-frontier (creates a SPOF). The scheduler is the single component that sees the global picture of host readiness + politeness budgets and makes dispatch decisions deterministically.

When it fails. Scheduler shard loses leadership → its hosts stop dispatching until re-election (~5 s). Whole tier dies → fetcher fleet idles within seconds, frontier backs up. etcd quorum loss → all 12 shards lose leadership, dispatch halts globally; runbook is restore-etcd-then-redo-elections.

Frontier (host subqueues)Cassandra (24-node cluster, RF=3)

Persistent per-host FIFO subqueues + back-queue index, partitioned by hash(host). Stores URLs awaiting fetch, ordered by priority+enqueued_at. Read by the scheduler on dispatch; written by the parser on URL discovery. Carries depth + host_quota for trap detection and per-host fairness.

Why it exists. Without a durable frontier, a process restart means losing all in-flight discovery and replaying it from the parser's stream — wastes politeness budget against every host. Persistence is the load-bearing property: re-crawl on loss is a community-trust catastrophe.

When it fails. Frontier corruption → restore-from-12 h-old-backup → 12 h of crawled URLs become 'unseen' again → re-crawl burst against every host → politeness storm → community ban. Mitigation: 15-min Cassandra incrementals + controlled re-warm runbook (10 % rate ramp over 4 h).

Seen-set (Bloom + Cassandra)Bloom bitarray sidecar + Cassandra journal (RF=3)

Two-tier dedup: in-memory Bloom bitarray (180 GB at 0.1 % FPR for 100 B known URLs, sharded 32-way) sits in front of a leaderless Cassandra journal for false-positive resolution and durable membership. Parser writes here on every commit; scheduler reads as a control-plane verification when popping URLs.

Why it exists. Without the seen-set a crawler re-fetches every URL forever — the literal 'crawl yourself in circles' failure mode. Bloom is the speed layer; the journal is the truth layer; together they answer 'have we seen this URL?' at p99 < 5 ms while remaining recoverable.

When it fails. Bloom shard loss → enters warming mode (bypass Bloom, hit Cassandra journal directly, 5–10× slower but correct). 5-min Bloom S3 snapshots restore in ~10 min; without them, replay-from-journal is hours-to-days during which FPR spikes and we re-enqueue billions of URLs.

Fetcher FleetOkHttp 5 + co-located Unbound DNS (c6a.4xlarge)

On lease-from-scheduler: (1) check robots cache, (2) DNS resolve via co-located Unbound (95 % local hit) or upstream, (3) GET the page from the target origin with 30 s hard timeout and 15 MB byte cap, (4) hand JS-required pages to the renderer pool, (5) publish raw fetched bytes to pages-stream with acks=all. Each pod runs 256 ToeThreads with exactly-1 concurrent fetch per host.

Why it exists. The hot path of the system. Every other component supports the fetcher's politeness contract. Sized for 32 % headroom over peak demand because tail events (slow-host clusters) compound — sized-for-the-median is how production crawlers stall.

When it fails. Fetcher fleet outage → no fetches happen → freshness SLO breached. Single pod loss = 0.13 % capacity loss (negligible). Correlated tail event (one ISP slowing) → main pool slots saturate → long-tail subpool is the safety valve. Bad deploy with politeness regression → progressive canary with politeness assertion in CI catches it before fleet-wide rollout.

robots.txt CacheRedis 7 cluster + google/robotstxt parser

RFC 9309-parsed entries keyed by host, 24 h TTL. Cluster-sharded by consistent-hash(host). Hot-reloads on operator takedown command (synthetic Disallow injection). Parse-failure entries are flagged fail-CLOSED so a 200-OK-with-HTML-body doesn't accidentally permit crawling for 24 h.

Why it exists. Robots compliance is the legal contract with the public web. A miss is a paged on-call incident. 1.16 K robots fetches/sec peak (15 % of fetch QPS) is non-trivial — caching is mandatory, and the takedown propagation channel runs through here.

When it fails. Robots cache cluster down → fetchers fail-closed (refuse to fetch) until restored. Stale cache on AOF restore (RPO 60 s) means up to 60 s of takedowns may not be visible — defense-in-depth is the scheduler-side suppression. Misparse-as-200 silently misclassifies hosts; alert is parse-failure-rate spike across hosts.

Renderer (Headless Chromium)Headless Chromium pool + per-tab memory cap

Renders SPA pages by handing the URL to a pool of headless Chromium tabs, returning the rendered DOM to the fetcher. 400 c6a.4xlarge × 24 renders/instance = 9 600 concurrent capacity. Per-tab 1 GB hard memory cap; HEAD pre-check skips render when Content-Length > 50 MB. Bounded queue (depth 5 K) with overflow=drop-to-raw-html so a render saturation degrades the index gracefully.

Why it exists. 20 % of pages need JS to surface meaningful content. Skipping renders means losing them from the index. The renderer is opt-in (not in the hot path of every fetch) because it's the dominant compute cost — $179 K/month — and we'd rather have a separate budget knob than couple it to fetcher fleet sizing.

When it fails. Renderer outage → SPAs drop out of the index until restored (raw HTML still flows). OOM cascade → kill-tab-not-pod isolates the bad page; queue-depth alert auto-scales before the cascade saturates the fleet. 200 MB-of-JS attack page → HEAD precheck skip + content-length gate.

Target Web (origins)

The actual public web — millions of hosts, every imaginable misconfiguration, anti-bot defenses (Cloudflare, Datadome, PerimeterX), expired certs, redirect loops, slow tail. Out of our control. We respect their robots.txt and politeness; we're guests.

Why it exists. The thing we're crawling. The crawler exists to ingest content from here without becoming the bug that pages someone else's on-call.

When it fails. Target hosts will fail individually constantly — that's the baseline. The crawler must tolerate it. Catastrophic mode: an entire CDN provider (Cloudflare) blanket-bans our IP range → 100 % 403s on tens of millions of hosts. Mitigation: 403-rate-per-CDN-signature alert + verified-bot programs.

Fetched Pages StreamKafka (512 partitions, RF=3)

Carries raw fetched bytes from fetcher to parser. 512 partitions × RF=3 with two-level keying (host + url-bucket) so a hot host like Wikipedia spreads across N partitions instead of pinning one. 7-day retention for parser-bug reprocessing without re-fetching the web. acks=all + min.insync=2 guarantees RPO=0 on broker loss.

Why it exists. Decouples fetch from parse so a parser deploy or stall doesn't backpressure the fetch loop. The 7-day retention is what lets us fix a parser bug and replay without re-fetching the web — a strict cost-and-politeness saver.

When it fails. Broker loss → ISR shrink alert + acks=all means the producer awaits ack; fetcher's local outbox absorbs the gap. min.insync=2 violation → page on-call (lost RPO=0 guarantee). Consumer lag spike → parser backpressure into pages-stream → the dreaded 'fetcher blocks on publish, slot exhausts' cascade — mitigated by the local fetcher outbox.

Parser / ExtractorJVM workers (Tika + jsoup + SimHash)

Consumes pages-stream (8 partitions per worker × 64 workers = 512 partitions covered). For each page: encoding detect → link extract → URL normalize → SimHash for near-dup → outbox-batch commit (seen + frontier in one Cassandra batch) → durable PUT to WARC → async batch publish to index. Idempotent on (url, fetch-ts) so a partition rebalance replay doesn't double-write.

Why it exists. All discovery + dedup + indexing logic lives here. Separated from the fetcher so a parser deploy doesn't touch the fetch fleet (different release cadences for different concerns).

When it fails. Parser deploy regression → poison messages on a partition → DLQ catches them after 3 retries (without the DLQ, one bad page stalls 1/512 of the crawl indefinitely). Outbox lag spike → page on-call before downstream divergence becomes visible. Cassandra batch failure → seen+frontier rolls back, URL re-fetched on next discovery; WARC may have orphaned bytes (acceptable, lifecycle-cleaned).

WARC ArchiveS3 (Standard + Glacier IA tiering)

ISO 28500 WARC.gz files, ~1 GB rotated, written by the parser as the durable archive of fetched bytes. ~250 TB/month ingest, 1.5 PB/year compressed. Revisit records (WARC-Refers-To) shrink dedup payloads by referring back to the original fetch instead of restoring the body. Object-lock + SSE-S3 + lifecycle to IA at 30 d → Glacier Deep at 365 d.

Why it exists. The archival layer + the rebuild substrate for the index. Anything that fails downstream of WARC is rebuildable from WARC. Without it, an OpenSearch corruption means re-crawling the web.

When it fails. S3 region outage → WARC writes fail with retries → outbox-driver eventually sidebands to a warc-failed table for later replay. We accept losing *that page's* archive in the worst case (the seen+frontier commit already succeeded, so the URL won't re-fetch on its own; manual replay from the warc-failed table is the recovery path).

Search IndexOpenSearch (32 primaries × 2 replicas each = 96 shards)

32 primary shards × 2 replicas (96 shards total), keyed by url-hash for even distribution. Bulk-ingest API batched 5 K docs / 10 s flush. 30 s refresh interval — freshness is per-document via _bulk timing, not refresh-interval. Hot tier holds last 90 d; older docs roll to a frozen tier on cheaper storage.

Why it exists. The downstream consumer of the crawl — the actual searchable corpus that powers the product (search engine, RAG corpus, AI training data). Decoupled from the parser via async batch ingest so an index outage doesn't block the parser.

When it fails. Index outage → bulk ingests fail → parser's index DLQ catches them for replay (no silent drops). Index corruption → reindex from WARC (the rebuild substrate). Cluster split-brain → OpenSearch's quorum elects a primary, ingests block briefly; indexing pipeline lag spikes.

Pages DLQKafka compacted DLQ (RF=3)

Receives messages that exceeded the parser's max-3-retries (Tika OOM on a malicious PDF, UTF-8 explosion, oversize body, header injection). 32 partitions × RF=3, 30 d retention so on-call has time to triage. Replay tooling re-injects fixed messages back to pages-stream after parser hotfix.

Why it exists. Without a DLQ, one poison message on partition 47 stalls 1/512 of the entire crawl indefinitely until on-call manually skips offsets — not a runbook. The DLQ is the difference between 'parser bug → degraded freshness' and 'parser bug → 0.2 % of crawl frozen for hours.'

When it fails. DLQ backlog spike alert → on-call inspects pattern, ships parser hotfix, replays. DLQ broker loss → pages-stream still flows; new poisons accumulate locally on parser pods until DLQ recovers.

DNS UpstreamTiered: Quad9 → Cloudflare 1.1.1.1 → Google 8.8.8.8

On co-located Unbound miss (~5 % of fetches), forwards to a tiered upstream: Quad9 → Cloudflare 1.1.1.1 → Google 8.8.8.8 in failover order. Resolver short-circuits to next provider on >100 ms p99 sustained. Negative caching at 5 min keeps typo-trap traffic from amplifying.

Why it exists. DNS resolver overload is failure mode #4 — a crawler DDoSing its own resolver. The tiered upstream is the load-bearing fix; making it a separate node makes DNS p99 a first-class observable, not buried inside the fetcher pod.

When it fails. All three upstream providers down → Unbound's local cache covers ~95 % of fetches; the remaining 5 % fail with circuit-breaker; freshness SLO degrades but no cascade. Single provider outage → invisible (next in chain). Slow upstream → resolver short-circuit kicks in within seconds.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Typical clarifications to surface:

  • Scope: Whole-web (Common Crawl shape) or vertical (e.g. e-commerce only)? Bounds the seed list and the abuse-tolerance.
  • Freshness vs. coverage: Is the goal "discover every URL once" (archival) or "keep the top 1 M sites fresh" (search-engine-shape)? Re-crawl scheduling looks completely different.
  • JS rendering: Index SPAs? If yes, the renderer fleet dominates compute cost.
  • AI-training corpus or live-search? Determines retention, dedup tolerance, and whether you publish the WARC archive.
  • Operator constraints: Take-down SLA? Verifiable identity (reverse-DNS + published IP allowlist) required to avoid Cloudflare AI Audit blocking?
  • Multi-region: Single-region keeps the frontier simple; multi-region requires partitioned-by-host frontiers + cross-region gossip for discovery.

Assumptions to state:

  • 10 B URLs/month = ~3.86 K sustained fetch QPS, ~7.7 K peak (2× burst).
  • 100 M hosts in scope; politeness floor 1 RPS/host, average sustained 0.05 RPS/host (long-tail dominated).
  • 20 % of pages need JS rendering.
  • Single AWS region (us-east-1); colocate everything with S3.
  • 5-year WARC retention.
  • Public, identified crawler — verifiable UA + reverse-DNS, IP allowlist published at /.well-known/crawler-ip-ranges.

02Functional reqs

What must this system actually do?

  • Submit seed URLs / sitemap pulls (operator API, low QPS).
  • Discover new URLs from parsed pages and enqueue if not seen.
  • Fetch URLs respecting per-host / per-IP / per-AS politeness windows + robots.txt rules.
  • Render JS-required pages via headless Chromium pool.
  • Archive fetched bytes to WARC in S3.
  • Index parsed text/metadata into the search index.
  • Take down a host or URL pattern within < 1 hour of operator command.
  • Re-crawl by adaptive schedule (change-rate driven, not size driven).
  • Publish identity (UA, IP allowlist, public robots.txt for our agent).

03Non-functional

What must it promise about speed, uptime and correctness?

  • Crawl freshness (top 1 M sites) p99: < 24 h since last successful fetch.
  • Robots.txt compliance: 99.99 % — every miss is a paged on-call incident with legal escalation.
  • Politeness violation rate: < 0.01 % of fetches exceed configured per-host RPS.
  • Fetcher availability: 99.9 % — drops are acceptable (the freshness SLO is what matters), but cascading freezes are not.
  • Frontier durability: 99.99 % within-region (~52 min annual budget) — single-region multi-AZ Cassandra at QUORUM with 5-min Bloom snapshots and 15-min Cassandra incrementals. Five-nines is reserved for cross-region crawlers that the next inflection-point upgrade will deliver; we explicitly do not claim 99.999 % on a single-region design with a 12 h restore RTO.
  • Indexing pipeline lag p99: < 5 min fresh, < 1 h backfill.
  • Take-down SLA: < 1 h global propagation, audit-logged.
  • Cost ceiling: Render fleet is the dominant variable cost; cap it at 30 % of total infra spend or scale the render fraction down.

04Capacity estimation

How much load and data does this have to hold?

Mid-tier walk-through. All numbers derived; alter the inputs in the simulator to re-derive.

Fetch QPS sustained = 10B / (30 × 86400)        = 3.86K  qps
Fetch QPS peak      = 3.86K × 2                  = 7.72K  qps

Concurrent hosts in flight. Politeness caps each host at 1 RPS, but the average polite host sustains ~0.05 RPS (long-tail wordpress sites). So:

Concurrent hosts = 7.72K / 0.05 = 154K hosts in flight at peak

That number is the design constraint. Heritrix tops out around 250 ToeThreads per process, but on a 2025 c6a.4xlarge (32 vCPU/64 GB) with HTTP/2 multiplexing we run 256 ToeThreads/pod conservatively. We size for 30 % headroom over the median demand because tail events (slow-host clusters) compound: 800 fetcher pods × 256 ToeThreads = 204 K slot capacity vs 154 K demand = 32 % cushion. 10 % of slots reserved for a "long-tail" subpool that handles hosts with historical p99 > 10 s, so slow hosts don't saturate the main fleet. The crawler is bound by host count, not raw fetch capacity.

Bloom filter for the seen-set. 100 B known URLs at 0.1 % FPR:

bits/URL = -ln(0.001) / (ln 2)^2 ≈ 14.4 bits ≈ 1.8 B
total    = 100B × 1.8 B          = 180 GB sharded
shards   = 32 (per-shard 5.6 GB, fits in RAM in front of the Cassandra journal)

robots.txt fetches. 100 M hosts / 24 h = 1.16 K req/s — 15 % of page-fetch QPS. Non-trivial; needs its own cache and dedicated robots-fetch headroom.

WARC ingest & storage.

ingest_MB/s = 3.86K × 25 KB        = 96 MB/s
daily       = 8.3 TB
monthly     = 250 TB
5y archive  = ~15 PB compressed

Render fleet. 20 % render-required × 7.72 K peak = 1.54 K renders/sec. At 5 s avg duration: 7.7 K concurrent renders. Chromium mean RSS is ~600 MB but p99 is 1.5–2 GB; the 1 GB per-tab cap is a kill threshold, not the mean. Sized for p99 burst: 400 c6a.4xlarge × 24 renders = 9 600 concurrent capacity vs 7 700 demand = 25 % headroom. ≈ $179 K/month (dominant compute cost). Bounded queue depth 5 K with overflow-policy = drop-to-raw-html so a render saturation degrades the index gracefully instead of cascading.

Kafka pages-stream. 7.72 K msg/s peak × 32 KB envelope (gzipped page + headers + URL + envelope) = 245 MB/s through the broker. 512 partitions sized for two-level keying (host + url-bucket) so a hot host like Wikipedia doesn't pin a single partition — the power-law p99/p50 of 5–10 × in flat hash-by-host needs the extra partitions plus the bucket dimension to stay inside per-partition 60 MB/s ceilings. RF=3 + min.insync=2 + acks=all gives RPO = 0 across an AZ loss.

Why 16 frontier shards / 32 seen shards. Frontier writes are dominated by parser discovery (~5–10× the fetch QPS in URLs/s but tiny payloads); 16 shards × 1 K writes/sec/shard fits comfortably in Cassandra at quorum. Seen-set is much larger (180 GB working set) and reads are RYW for parser self-checks → 32 shards isolates the hot edge of the keyspace.

Why these timeouts. Fetcher → target-web at 30 s = the hard ceiling (industry-standard, matches Googlebot). Fetcher → renderer 35 s > renderer → target-web 25 s so the cascading-retry rule (callee_total_budget < caller_per_attempt_timeout) holds and the renderer is allowed to fail-fast inside the fetcher's window. Fetcher → pages-stream at 5 s with acks=all+fixed-retries = the producer DOES await ack; on Kafka unavailability the fetcher writes to a local disk-backed outbox (60 s capacity) so the slot is freed and the page is not dropped. Scheduler → fetcher at 100 ms = the dispatch lease ack budget; the fetch itself is async after the ack. Cascade chain: scheduler 100 ms > fetcher-ack ~5 ms ✓; fetcher 30 s > publish 5 s × 3 retries = 15 s ✓ but the local outbox absorbs publish failures so the slot SLA holds at fetch+15 s = 45 s p99 worst-case; scheduler dispatch loop tolerates this because slots are not slot-pinned-per-URL but slot-pinned-per-fetch.

05API design

What does the outside world call, and what comes back?

POST /api/v1/seeds
Authorization: Bearer <op-token>
{ "urls": ["https://example.com/", "https://example.org/sitemap.xml"], "priority": "normal" }
202 Accepted        # gw → scheduler (e3); scheduler normalizes + enqueues into the frontier (e15b)

POST /api/v1/takedown
{ "scope": "host", "value": "example.com", "reason": "operator-request", "ticketId": "TKT-1234" }
202 Accepted        # propagates to scheduler + robots cache within 1h

GET /api/v1/freshness?host=example.com
200 { "host": "example.com", "lastFetchedAt": "...", "knownUrls": 1834221, "fetchedThisMonth": 18234 }

# Public crawler-identity surface (unauthenticated, CDN-cached):
GET /.well-known/crawler-ip-ranges
200 { "ranges": ["203.0.113.0/24", ...], "userAgents": ["ArqlyBot/1.0"], "verifyDoc": "https://arqly.example/bots" }

06Data model

What gets stored, and what is it looked up by?

Frontier (Cassandra, 16-shard, RF=3):

fieldtypenotes
hosttextpartition key (hash routes to shard)
prioritysmallintclustering key 1
enqueued_attimestampclustering key 2
urltext
depthsmallinttrap-detection input

Per-host quota (Cassandra counter table — a counter column can't share a table with non-counter columns, so the per-host cap lives beside the frontier, not inside it):

fieldtypenotes
hosttextpartition key
enqueued_ctcounterrunning enqueue count this month; capped at 5 % of frontier budget per host

Seen-set (Cassandra journal beneath an in-memory Bloom):

fieldtypenotes
url_fpbigintpartition key — 64-bit fingerprint
seen_attimestamp
simhashbigintfor near-duplicate detection at indexing time
disposesmallint0=keep, 1=trap-suspect, 2=takedown

robots.txt cache (Redis hash):

fieldnotes
hostkey
etagfor conditional GET on refresh
parsed_rulesRFC 9309 byte-encoded, compact
parse_failedbool — fail-CLOSED if true
fetched_attimestamp; TTL 24 h

Database choice — recommended. Cassandra (or ScyllaDB) for both frontier and seen-set. Why: leaderless replication survives an AZ loss without a write outage; quorum tunability lets the seen-set serve RYW reads from the parser without paying linearizable cost on the frontier dequeue. Discord's trillion-message ScyllaDB migration shows the pattern at write-heavy scale. Postgres alone falls over at 100 B-row dedup tables; DynamoDB works but locks you to AWS and the per-row pricing on a 100 B-row table is brutal. Redis is the right cache for robots.txt (24 h TTL, 200 GB working set) but not the system of record. S3 for WARC is uncontroversial — it's what Common Crawl and Internet Archive use, with 11×9 s durability built in.

07High-level design

Which components handle a request, and in what order?

Architecture summary — the three loops:

  1. Crawl loop (hot path). Scheduler pops a host whose politeness window has elapsed → leases the URL to a fetcher pod → fetcher checks robots cache → fetches the page (or hands off to renderer for SPAs) → publishes the raw fetch to pages-stream. The fetcher never blocks on a downstream commit — the parser side runs independently.
  2. Parse-and-commit loop. Parser consumes pages-stream → extracts links → SimHash for near-dup detection → outbox-batch commit (seen + frontier in one Cassandra batch keyed by host) → durable PUT to WARC → async batch publish to the search index. The outbox is the load-bearing atomicity primitive: seen + frontier must never diverge, or you re-crawl forever (or never at all).
  3. Operator/admin loop. Operator submits seeds and takedowns through the dashboard → CDN → Gateway → Scheduler. Takedown propagates to the scheduler's politeness controller AND the robots cache as a synthetic Disallow within < 1 h. Seed injection: the scheduler validates + normalizes the seed and enqueues it into the per-host frontier subqueue (e15b, scheduler → frontier, op: write) — the same normalize+enqueue terminus the parser uses on discovery (e15). A seed routed gw → scheduler (e3) therefore has a real frontier terminus: it becomes a frontier row exactly like a discovered URL. The seed enqueue is a single best-effort write (no page-commit txn); only the parser's discovery write (e15) is part of the outbox group, since only there must seen + frontier commit atomically.

Why the Mercator back-queue. The classic 1999 design — per-host FIFO subqueues with a back-queue selector picking the next ready host — is still the right primitive. It enforces the per-host concurrency invariant as a data-structure property, not a runtime check. Heritrix's BdbFrontier, StormCrawler's ES-status-index, and Googlebot's "Jack" all descend from this. The scheduler shards by hash(host) so politeness state stays local and a 12-pod scheduler scales horizontally.

Why Bloom + Cassandra journal for the seen-set. A naive all-Cassandra seen-set requires 100 B point reads per discovery batch (10× fan-out of the fetch rate) — that's 80 K reads/sec on the parser side, untenable at quorum latency. The Bloom filter cuts 99.9 % of those before they hit the journal, and the Cassandra journal serves the remaining false-positive lookups in single-digit ms at quorum. The 0.1 % FPR is the right knob: a false positive means we don't enqueue a URL we haven't seen, which is mild data loss in a fundamentally best-effort archival pipeline. False negatives (saying we haven't seen something we have) are impossible — that's why Bloom over hash-set is correct here, not the other way around.

Why fail-CLOSED on robots.txt parse error. RFC 9309 is strict; mis-served HTML-as-robots is a real failure mode. Treating a parse failure as "Allow: /" trades 24 h of compliance violation for a few minutes of operator pain. Treating it as "Disallow: /" risks dropping a few thousand crawlable hosts but bounds legal exposure. Production crawlers (Google, Bing, Common Crawl) all fail-closed; we follow.

Why outbox, not 2PC, on the page commit. The parser writes seen + frontier in a single Cassandra batch (single shard, both keyed by host) PLUS an outbox row in the same batch — that is the only atomic step (mode: outbox, participants: dedup + enqueue). The outbox row is then drained by the parser to: (a) WARC PUT (sync, exp-backoff with max 5 retries before sideband to a warc-failed table), and (b) bulk-publish to the index ingester. WARC is not a participant in the page-commit txn — it is downstream of the outbox drain, so a WARC PUT failure cannot un-commit seen+frontier and a successful seen+frontier commit cannot block on WARC. This matches the JSON: only e14 (parser→seen) and e15 (parser→frontier) carry txn.group=page-commit; e13 (parser→warc) and e16 (parser→index) are downstream propagators. 2PC across Cassandra and S3 is impossible in practice (S3 has no participant API); saga is overkill since the compensating action is "the next discovery re-enqueues the URL and we re-fetch," which is not a true compensation but the natural retry semantic. The accepted trade-off: a crash between the Cassandra batch commit and the WARC PUT loses that page's archive (not the index entry, which is rebuilt on next re-crawl). Outbox-lag dashboard catches systemic versions.

Why per-host AND per-IP AND per-AS politeness. A naive per-host limiter explodes against shared hosting (100 K vhosts on one /24). Cloudflare's published per-AS observations show the long tail of crawler abuse: most polite-per-host crawlers are accidentally rude per-AS. Three-tier token bucket in the scheduler is the load-bearing fix. Coordination across the 12 scheduler shards is centralized in a separate Redis cluster (one logical cluster, shared by all scheduler shards) keyed by (tier, key) where tier ∈ {host, /24, AS}. Per-shard local cache absorbs reads with 100 ms TTL; writes (decrement) go to the central Redis with optimistic concurrency. Without this central coordination, scheduler shard 5 and shard 7 would each grant 1 RPS to AS-12345 → 2× the per-AS budget. The Redis cluster is RDB-persisted with 5-min snapshots so a restart doesn't reset every bucket to "full" (which would permit a politeness-regression burst).

Takedown control plane (< 1 h SLA). The takedown propagation is engineered, not asserted. (1) Operator submits to gw with bearer + ticket-system reference. (2) gw writes synchronously to a Kafka compacted topic takedown-bus (RF=3, retain forever) AND issues control RPCs to the scheduler (immediate suppression on dispatch) and to the robots cache (synthetic Disallow: / injection for the targeted host). (3) Fetchers re-pull from the robots cache on every fetch (no fetcher-resident takedown table) so robots-cache-injection alone covers the fleet within one robots-TTL refresh. (4) Audit log writes to takedown-bus with timestamps for every step; the legal team gets a runbook of "submission → propagation → confirmation". Worst-case propagation budget: gw→scheduler RPC (1 s) + gw→robots RPC (1 s) + 600 fetcher pods discovering the takedown on next robots-cache-refresh (worst case = robots TTL = 24 h, but typical = next fetch within seconds). To meet the < 1 h SLA the robots cache injection is the load-bearing step (immediate); the scheduler suppression is a defense-in-depth.

Multi-region posture. Single-region (us-east-1, colocated with S3) for this canonical. Cross-region adds: partitioned-by-host frontier with cross-region gossip for discovery (each region owns a host-shard and discovers URLs locally; a small cross-region delta replicates "newly discovered" deltas every 30 s). For an SRE proposing global active-active, the burden is showing that the cross-region egress + politeness coordination is worth it; for most use-cases, single-region with multi-AZ for HA is the right call.

Primary sources

  • Mercator: A Scalable, Extensible Web Crawler — Heydon & Najork (Compaq SRC, 1999)
  • Detecting Near-Duplicates for Web Crawling — Manku, Jain, Das Sarma (Google, WWW 2007)
  • Inside Googlebot — Google Search Central Blog (Mar 2026)
  • Heritrix Politeness Parameters — Internet Archive Wiki
  • Common Crawl monthly statistics & WARC file format
  • Cloudflare AI Audit — blocking AI crawlers with one click
  • Cloudflare on Perplexity stealth crawlers (Aug 2025)
  • RFC 9309 — Robots Exclusion Protocol
  • Google open-source robots.txt parser (google/robotstxt)
  • Apache StormCrawler vs Nutch (DZone)

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work Web Crawler yourself