System Design Problem

Design a Web Crawler (Googlebot)

Commonly Asked By:GoogleMicrosoftYahoo

Interview Setup

Interview Prompt

Design a web crawler that discovers and downloads billions of web pages for a search engine index. Respect robots.txt, avoid overloading sites, and deduplicate near-identical content.

Clarifying Questions (ask before designing)

QuestionWhy it matters
What's the crawl rate target and how many domains?Crawling 1B pages/day across 100M domains needs ~11.5K pages/sec with strict per-domain politeness. Crawling 1B pages/day from 10K domains creates an entirely different host bottleneck.
Freshness priority or breadth-first discovery?BFS discovers new sites, whereas the freshness scheduler recrawls known pages using separate queues.
Exact dedup or near-duplicate detection?Exact URL deduplication uses a hash set, while near-duplicate detection across mirrors and templates requires SimHash or MinHash.
JavaScript-rendered pages in scope?Headless browser crawling is roughly 100 times slower, requiring a dedicated rendering pipeline.

Scope

In scope

  • URL frontier with BFS and priority scheduling
  • Politeness: per-domain rate limiting
  • robots.txt fetch, parse, and cache
  • Content deduplication (exact URL + SimHash near-dedup)
  • DNS caching and resolution
  • Distributed crawler worker architecture

Out of scope (state explicitly)

  • HTML parsing and index extraction (downstream indexer)
  • JavaScript rendering (headless Chrome farm)
  • Login/authenticated crawling
  • Image/video binary crawling

Functional Requirements

Start by asking your interviewer about seed URLs, politeness requirements, and deduplication scope. Distributed crawl vs single-machine and JavaScript rendering are the usual scope pivots.

  • Crawl the web starting from a set of seed URLs
  • Discover new URLs by extracting links from crawled pages
  • Download and store web page content for indexing
  • Respect robots.txt directives (politeness)
  • Handle URL deduplication (don't crawl the same page twice)
  • Support recrawling to detect updated content
  • Prioritize important/popular pages for crawling first
  • Handle different content types (HTML, PDF, images, etc.)

Non-Functional Requirements

Your interviewer will care most about politeness and throughput on this problem. Per-domain rate limiting shapes the entire frontier design, so explain this requirement early to justify your two-tier queue architecture.

  • High Throughput: Crawl 1 billion pages per day
  • Politeness: Don't overload any single web server (rate limit per domain)
  • Robustness: Handle spider traps, infinite loops, malformed HTML, and server errors
  • Scalability: Horizontally scalable by adding worker nodes to increase crawl throughput linearly
  • Freshness: Recrawl pages adaptively based on content change frequency
  • Extensible: Modular architecture to easily add content extraction, language detection, and metadata processing
  • Fault Tolerant: Node failures must not cause data loss or interrupt overall crawl progress

Capacity Estimations

Run this math before sizing the worker fleet. Pages per day and average page size determine bandwidth and storage, while the frontier size drives whether an in-memory queue is even possible.

MetricCalculationValue
Pages to crawlGiven (assumption documented in value)15B (entire web, roughly)
Target crawl rateGiven (assumption documented in value)1B pages/day
Pages / secGiven~11,500
Avg page sizeGiven (typical workload assumption)50 KB
Download / day1B x 50 KB~50 TB
Download bandwidth50 TB ÷ 86400~5 Gbps sustained
URL frontier size10B URLs x 200 bytes~2 TB
Content storage / monthRaw payload: ~50 TB/day ≈ ~1.5 PB/month. Planned storage: ~3 PB/month including replication, metadata, retention/versioning, and storage overhead~3 PB planned (~1.5 PB raw payload)
Worker fleet~11,500 pages/sec ÷ ~4.5 pages/sec/worker ≈ ~2,500 concurrent fetch workers across 100 nodes100 nodes (~2,500 worker processes, ~25 per node)

Architecture Diagram

In the room: ask about politeness defaults early because naive parallel crawling gets you blocked and is an instant red flag.

Walk your interviewer through the diagram by data flow. We maintain a URL frontier that schedules what to fetch next while enforcing per-domain rate limits. Fetcher workers download pages, deduplicate content, and store raw HTML in object storage. The parser extracts links back into the frontier and publishes crawl events to Kafka for the downstream indexer. DNS caching and robots.txt checks sit on the hot fetch path so politeness does not add latency at scale.

Loading...

Component Deep Dives

Next we walk through each component in the system architecture. Start with the frontier because it represents the most complex scheduling component, then cover the fetch path, deduplication layers, and index handoff.

URL Frontier: The Most Complex Component

The URL frontier is not a simple queue because it must balance priority scheduling against per-domain politeness constraints. To achieve this, the frontier operates as a two-tier queue system:

1. Priority Queue (Front Queue): Determines which URLs to crawl first based on PageRank, change frequency, domain authority, and freshness requirements. This tier is implemented using multiple priority queues spanning P0 through PN.

2. Politeness Queue (Back Queue): Prevents overwhelming any individual web server by allocating one dedicated FIFO queue per target domain. A rate limiter enforces a default limit of one request per domain per second, or strictly honors the Crawl-delay specified in robots.txt.

Loading...

Fetcher (HTTP Downloader)

Each fetch request must respect robots.txt rules, connection limits, and network timeouts before downloaded content ever reaches persistent storage.

  • Connection pooling: Reuse persistent HTTP connections to the same host
  • Timeouts: Connect timeout of 5 seconds and read timeout of 30 seconds
  • Redirect handling: Follow up to 5 HTTP redirects (handling 301 and 302 status codes)
  • robots.txt compliance: Before crawling any page on a domain, fetch and cache robots.txt. Parse directives including Disallow, Allow, Crawl-delay, and Sitemap, caching parsed rules per domain with a 24-hour TTL
  • Content types: Accept text-based formats including HTML, PDF, and DOC while rejecting raw binary media like video and audio

Content Deduplication

Near-duplicate detection saves index storage because exact URL deduplication alone is insufficient at web scale. Many pages share identical or near-identical content due to mirrors, syndication, and boilerplate templates.

  • Exact deduplication: Compute an MD5 or SHA-256 hash of the page body and verify it against the database of previously seen content hashes.
  • Near-duplicate detection: Execute the SimHash or MinHash algorithm over tokenized shingles. Two pages are classified as near-duplicates when the Hamming distance between their 64-bit SimHashes is 3 or less (Hamming distance ≤ 3).
  • Storage layers: Maintain an in-memory Bloom filter for exact body hashes alongside a locality-sensitive candidate index in Cassandra, calculating exact Hamming distance against retrieved candidates for near-duplicate confirmation.

URL Deduplication (Seen URLs)

URL normalization prevents crawling the same page under multiple distinct URL representations, and should always be applied before checking the seen-URL repository.

  • Seen-URLs Bloom filter: Use a distributed Bloom filter configured for 100 billion entries, consuming approximately 120–125 GB of memory (10 bits per URL = 1 trillion bits) targeting an expected ~1% false-positive rate (with 7 hash functions).
  • Frontier admission check: Query the Bloom filter before adding any candidate URL to the frontier, falling back to an exact database check in Cassandra if a match occurs.
  • URL normalization rules: Distinguish safe, semantics-preserving normalizations (lowercasing scheme and host, removing fragments and default ports, resolving dot-segments) from configurable heuristics (stripping tracking parameters, trailing-slash policy, query parameter sorting/stripping depending on application semantics).

DNS Resolver

Uncached DNS lookups introduce 50 to 200 milliseconds of latency and quickly become a major bottleneck at 11,500 pages per second, requiring aggressive prefetching and multi-tiered caching.

  • Tiered caching: Combine an in-memory LRU cache inside each worker process with a shared Redis cache layer for high hit rates.
  • Batch DNS prefetching: Pre-resolve DNS records asynchronously for upcoming domains queued in the frontier before workers dequeue their URLs.
  • TTL and negative caching: Honor upstream DNS record TTLs while caching negative responses (NXDOMAIN) for 1 hour to prevent repeated lookups for dead domains.

Parser / Link Extractor

Link extraction feeds newly discovered links back into the frontier, providing the underlying engine that drives breadth-first web discovery.

  • Resilient HTML parsing: Parse document structures with a fault-tolerant parser capable of recovering from malformed HTML and unclosed tags.
  • Entity and metadata extraction: Extract hyperlinks, plain text content, OpenGraph metadata, JSON-LD micro-data, and base tags.
  • JavaScript handling strategy: Offload client-rendered JavaScript pages to a specialized headless browser cluster for high-value targets, while processing standard pages through fast server-rendered HTML parsers.

Index Pipeline Handoff

Decoupling crawl throughput from index construction latency ensures that the crawler can operate at peak network speeds while Kafka absorbs downstream backpressure whenever index builds fall behind.

  • After committing raw page content to the object store (S3/GFS), fetch workers emit a page-crawled event to Kafka carrying URL metadata, content hash, and the object storage content reference.
  • The downstream Search Engine indexer consumes these events asynchronously, fetching the stored document by storage key to update inverted indices and posting lists.
  • Decoupling crawl throughput and raw document persistence from index build latency allows Kafka message retention to absorb downstream pipeline backpressure without pausing active fetchers.

Event Bus Design (Kafka)

The event bus decouples crawl workers from downstream search indexing consumers and buffers discovered URLs for frontier ingestion, absorbing traffic spikes and enabling independent pipeline scaling.

YAML
# Kafka Event Bus Topology for Crawler Coordination and Search Index Ingestion

topics:
  page-crawled:
    purpose: "Notifies search index pipeline of successfully crawled pages already committed to S3/GFS content store"
    event_schema:
      event_id: "string"
      url: "string"
      domain: "string"
      content_hash: "string (lookup key in S3/GFS content store)"
      storage_key: "string (s3://crawler-content/{content_hash})"
      status_code: "number"
      crawled_at: "timestamp"
      simhash: "string"

  discovered-urls:
    purpose: "New URLs extracted by parser routed to URL Frontier, keyed by domain"
    partitions: 1024
    partition_key: "domain (ensures per-domain ordering and partition affinity within Kafka transport)"
    retention: "3 days"
    replication_factor: 3
    min_insync_replicas: 2
    event_schema:
      event_id: "string"
      url: "string"
      domain: "string"
      priority: "number"
      depth: "number"
      referrer_url: "string"

  crawl-failures:
    purpose: "HTTP errors, connection resets, and timeouts routed to retry and deprioritization queues"

producers:
  crawler_fetch_workers:
    type: "Idempotent producer"
    behavior: "Crawl workers persist raw HTML to S3/GFS, then publish page-crawled and discovered-urls events asynchronously so the fetch loop never blocks on downstream indexing"

consumer_groups:
  frontier-ingest:
    responsibility: "Consumes domain-keyed URL partitions from Kafka, verifies against Bloom filter, consults the consistent-hash ring to route URLs to the authoritative domain owner, and enqueues into local frontiers"
    dead_letter_queue: "discovered-urls-dlq after 3 retries"
    alerting: "Trigger alert when consumer group lag exceeds 5 minutes"
  index-pipeline:
    responsibility: "Consumes page-crawled events and uses the stored content reference to trigger downstream Search Engine indexing"
    dead_letter_queue: "page-crawled-dlq after 3 retries"
    alerting: "Trigger alert when consumer lag exceeds 15 minutes"

API Design

These are internal operational APIs, so present seed submission and cluster status endpoints to demonstrate how operators manage and inspect the crawl fleet.

Crawler Service API Interface

TYPESCRIPT
// Internal service contracts for crawl coordination and administrative operations
interface AddSeedsRequest {
  urls: string[];
  priority?: "critical" | "high" | "normal" | "low";
  crawlDelayOverrideMs?: number;
}

interface AddSeedsResponse {
  accepted: number;
  duplicateCount: number;
  status: "queued" | "rejected";
}

interface CrawlStatusResponse {
  pagesCrawledToday: number;
  pagesInFrontier: number;
  crawlRatePerSec: number;
  activeNodes: number;
  activeWorkerProcesses: number;
  errorsToday: number;
}

interface BlockDomainRequest {
  domain: string;
  reason: string;
  expiresAt?: string;
}

interface CrawlerAdminService {
  addSeeds(request: AddSeedsRequest): Promise<AddSeedsResponse>;
  getCrawlStatus(): Promise<CrawlStatusResponse>;
  blockDomain(request: BlockDomainRequest): Promise<{ success: boolean; blockedAt: string }>;
}

Add Seed URLs

HTTP
POST /api/v1/crawler/seeds
Content-Type: application/json

{
  "urls": ["https://example.com", "https://news.ycombinator.com"],
  "priority": "high"
}

HTTP/1.1 202 Accepted
Content-Type: application/json

{
  "accepted": 2,
  "duplicateCount": 0,
  "status": "queued"
}

Get Crawl Status

HTTP
GET /api/v1/crawler/status

HTTP/1.1 200 OK
Content-Type: application/json

{
  "pages_crawled_today": 892345678,
  "pages_in_frontier": 5234567890,
  "crawl_rate_per_sec": 11500,
  "active_nodes": 98,
  "active_worker_processes": 2450,
  "errors_today": 12345
}

Block Domain

HTTP
POST /api/v1/crawler/block
Content-Type: application/json

{
  "domain": "spam-site.com",
  "reason": "spam"
}

HTTP/1.1 200 OK
Content-Type: application/json

{
  "success": true,
  "blockedAt": "2026-03-13T10:00:00Z"
}

Common Error Responses

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
202 Accepted: asynchronous job queued successfully, poll GET /jobs/{id} for completion status
408 Request Timeout: background job is still executing, continue polling status endpoint

Data Model

The data model mirrors the crawl lifecycle by maintaining structured queue entries for frontier scheduling, Bloom filter states for fast deduplication, and immutable blob records for page content storage.

URL Frontier Entry

JSON
{
  "url": "https://example.com/page/123",
  "normalized_url": "https://example.com/page/123",
  "domain": "example.com",
  "priority": 2,
  "discovered_at": "2026-03-13T10:00:00Z",
  "last_crawled_at": "2026-03-12T08:00:00Z",
  "recrawl_interval": 86400,
  "retries": 0,
  "depth": 3,
  "referrer_url": "https://example.com/"
}

Content Store (S3 / GFS)

YAML
bucket: "crawler-content"
key: "{content_hash}"
value:
  url: "https://example.com/page/123"
  content_type: "text/html"
  status_code: 200
  headers:
    content-type: "text/html; charset=UTF-8"
    server: "nginx"
  body: "<html>...</html>"
  crawled_at: "2026-03-13T10:00:00Z"
  content_hash: "sha256:7f83b1657ff1fc53b92dc18148a1d65dfc2d4b1fa3d677284addd200126d9069"
  simhash: "0xA3F2B1C4D5E6F7A8"

Robots.txt Cache (Redis)

YAML
key: "robots:{domain}"
ttl_seconds: 86400
value:
  rules:
    - user_agent: "*"
      disallow:
        - "/admin"
        - "/private"
    - user_agent: "Googlebot"
      allow:
        - "/api"
  crawl_delay: 2
  sitemaps:
    - "https://example.com/sitemap.xml"
  fetched_at: "2026-03-13T00:00:00Z"

Bloom Filter: Seen URLs

YAML
type: "Distributed Bloom Filter (Redis-backed or custom off-heap)"
capacity: "100 billion URLs"
false_positive_rate: "~1% (expected with 7 hash functions)"
memory_size: "~120–125 GB (10 bits per element x 100B URLs = ~1 trillion bits)"
hash_functions: 7

Kafka Topics

YAML
topics:
  page-crawled:
    payload: "(url, content_hash, storage_key, simhash, crawled_at)"
    destination: "Search Engine Index Pipeline (indexer reads raw document from S3/GFS by storage_key)"
  discovered-urls:
    partitions: 1024
    partition_key: "domain"
    payload: "(new URLs extracted from parsed pages)"
    destination: "URL Frontier (Kafka buffers domain-keyed URLs, consumer consults consistent-hash ring to forward to authoritative domain owner)"
  crawl-failures:
    payload: "(HTTP errors, timeouts, connection resets)"
    destination: "Retry and backoff scheduler"

Fault Tolerance

Operational resilience requires mitigating worker failures, resolver outages, and adversarial crawler traps before they consume cluster resources.

ConcernSolution
Worker crashURL remains in the frontier until crawl completion is confirmed, then reassigned to another active worker
DNS failureRetry with exponential backoff, then fall back to secondary DNS resolvers
HTTP timeoutRetry up to 3 times with backoff, then mark the URL as failed and deprioritize it
Content store failurePersist to S3 with 11 nines of durability and enable cross-region replication
Frontier data lossCheckpoint the frontier to disk periodically, and rebuild from crawled URLs stored in the content store if lost
Bloom filter lossRebuild the filter from frontier state and stored document URLs during recovery

Spider Traps and Infinite Loops

Adversarial and poorly designed websites can generate infinite URL loops using dynamic calendars, recursive directory structures, or session-tracking tokens. The crawler enforces five distinct defensive layers to isolate and mitigate traps:

  • Maximum crawl depth: Terminate link traversal once a branch reaches depth 15 from its originating seed URL.
  • Layered per-domain page caps: Require a verified sitemap once a domain surpasses 100,000 discovered URLs, trigger an automated operational alert and quarantine throttling at 500,000 URLs, and enforce a hard stop at 1 million pages per crawl cycle.
  • URL pattern detection: Halt expansion when more than 1,000 URLs from a single host match an identical path prefix or template structure.
  • Maximum page size: Discard any document exceeding 10 MB to protect worker memory and avoid decompression bombs.
  • Content hash deduplication: Compare SHA-256 checksums of downloaded content so that identical bodies served under distinct URLs are dropped immediately without link expansion.

Additional Considerations

JavaScript rendering, sitemap parsing, and crawl budget allocation are senior follow-ups.

Freshness Scheduler

Discovery and freshness operations use separate frontiers. The freshness scheduler maintains a per-URL next_crawl_at timestamp derived from PageRank, sitemap changefreq metadata, and observed historical change rates. A dedicated worker pool drains the freshness queue so that it never competes with breadth-first discovery workers for politeness slots on the same domain. When the index pipeline (Search Engine) signals backpressure, the scheduler pauses low-priority recrawls first while preserving the 15-minute recrawl cadence for high-priority news URLs.

Recrawl Strategy

  • News sites: Recrawled every 15 minutes to capture breaking updates.
  • E-commerce product pages: Recrawled every few hours to track price and inventory fluctuations.
  • Reference and documentation sites: Recrawled every few days.
  • Static content pages: Recrawled on a weekly or monthly cadence.
  • Adaptive recrawling algorithm: Track how often a page changes historically and adjust the recrawl interval dynamically. If content changed since the last fetch, halve the interval down to the minimum threshold. If the content remained unchanged, double the interval up to the maximum threshold.

Distributed Architecture

  • Master-Worker coordination: A central master node assigns URL batches to worker nodes, which execute fetches and report status back to the coordinator.
  • Masterless coordination with consistent hashing: Each worker node manages an independent subset of domains assigned through Consistent Hashing on domain names. This design eliminates any single point of failure because workers coordinate asynchronously through Kafka partitions.

Legal and Ethical Considerations

  • robots.txt compliance: Always fetch and strictly respect robots.txt directives for every visited domain.
  • Crawl-delay adherence: Honor requested delays between successive HTTP requests to avoid degrading host performance.
  • noindex directives: Forward noindex robots metadata to the indexer to ensure excluded pages are never searchable.
  • Terms of service: Respect published terms of service where domains explicitly disallow automated scraping beyond robots.txt.
  • Privacy regulations: Redact and discard personally identifiable information (PII) during extraction to comply with GDPR and privacy mandates.

Sitemap Processing

  • Parse sitemap.xml documents to obtain curated manifests of authoritative site URLs.
  • Extract priority and changefreq XML hints to initialize priority queues and recrawl timers.
  • Support nested sitemap index hierarchies that point to multiple compressed child sitemap files.

Related Problems

Crawl output feeds the indexing pipeline in Search Engine. Politeness and frontier scheduling patterns overlap with Price Comparison Engine for tiered crawling and Top-K Rankings for velocity scoring of priority URLs. You can also explore foundational architecture primitives in Message Queues Fundamentals, Consistent Hashing, Distributed File Systems (GFS/HDFS), and System Design Interview Patterns.

Interview Walkthrough

  • 25-minute cut

    Skip arch50/arch75 depth unless staff.

    • FR/NFR and politeness scope (3 min)
    • Two-tier URL frontier design (8 min)
    • Per-domain rate limiting and robots.txt (7 min)
    • Bloom filter dedup and crawl throughput math (7 min)
  • Start with the URL frontier as the central scheduler that decides what to crawl next while enforcing politeness constraints.
  • Explain the two-tier frontier where front priority queues feed per-domain back queues for politeness rate limiting, as this architecture is a primary interview probe.
  • Use a Bloom filter for visited-URL deduplication to keep the frontier memory-efficient at billions of URLs.
  • Describe worker distribution: master-worker for simplicity or consistent-hashing workers by domain for fault tolerance.
  • Cover recrawl scheduling using adaptive intervals calculated from page change frequency rather than a static global timer.
  • Emphasize robots.txt and crawl-delay compliance because ignoring politeness gets your crawler blocked and creates an immediate red flag.
  • Quantify throughput using Back-of-the-Envelope Estimation where pages per second multiplied by average page size determines network bandwidth and storage needs.
  • Highlight a common pitfall: relying on a naive global BFS queue without per-domain rate limiting, which floods individual hosts and triggers IP bans.

Engineering Trade-offs

Your interviewer will push on frontier design and coordination. Walk through priority vs politeness queues, then state what you'd pick for a billion-pages-per-day crawler.

URL Frontier: Priority + Politeness Architecture

The URL Frontier acts as the central scheduler and is built as a two-tier system: Front Queues (Priority) and Back Queues (Politeness). Together they balance finding the most important pages early while respecting domain rate limits.

Loading...

1. Front Queues (Priority Allocation):

  • URLs are classified into multiple priority levels from critical to low:
    • Queue P1 (Critical): e.g., [google.com/news, bbc.com/latest]
    • Queue P2 (High): e.g., [wikipedia.org/..., nytimes.com/]
    • Queue P3 (Normal): e.g., [example.com/about, blog.io/post]
    • Queue P4 (Low): e.g., [random-site.xyz/page42]
  • Priority Assignment Signals:
    • PageRank: Pages on domains with higher PageRank get prioritised first.
    • Freshness Need: Frequently updated directories (e.g. news sites) are scheduled on higher-priority tiers.
    • Sitemap Hints: Explicit <priority> hints in sitemaps (e.g. 0.8 high vs 0.2 low).
    • Historical Change Frequency: If historical scans show a page mutates daily, its priority is increased.
  • Selection Process: The prioritizer selects URLs using a weighted random probability distribution: P1: 40% chance, P2: 30% chance, P3: 20% chance, and P4: 10% chance.

2. Back Queues (Politeness Rate Limiting):

  • To avoid crashing web servers (Denial-of-Service), the crawler maintains exactly one FIFO queue per target domain.
  • Each domain queue tracks performance states:
    • Queue [google.com]: Tracks the last fetch timestamp and active crawl_delay (such as 2 seconds) so that the next crawl window is calculated dynamically.
    • Queue [wikipedia.org]: With a longer crawl delay (such as 5 seconds), the next fetch is scheduled 5 seconds after the last successful download.
  • Worker Loop: Worker threads scan the back queues to find domains where next_fetch_time ≤ now(), dequeue one URL, issue the HTTP fetch, and set the next available window.

Crawl Execution Flow: One URL End-to-End

Processing a single target URL (e.g., https://example.com/products/shoes) follows a strictly ordered synchronous timeline to maximize efficiency while validating constraints:

  1. Step 1: Dequeue from Frontier (0.1 ms)
    Checks eligibility of the back queue (e.g., [example.com]). If last_fetch was 3s ago and crawl_delay is 2s, the queue is eligible. Pops the URL from the queue.
  2. Step 2: DNS Resolution (1 ms: Cached)
    Checks the local DNS cache for the host IP (such as resolving example.com to 93.184.216.34). On a cache miss, queries recursive DNS resolvers and updates the cache with a standard TTL (such as 300 seconds).
    Optimization: Custom batch DNS prefetching executes in the background for upcoming back-queue items.
  3. Step 3: Robots.txt Constraint Validation (0.1 ms: Cached)
    Fetches cached robots rules from Redis (e.g., key robots:example.com). If the URL is matched against a block pattern (e.g. Disallow: /products/shoes) for our user agent, the crawler drops the URL instantly, logs a rejection, and moves to the next candidate.
  4. Step 4: HTTP Fetch (200 ms: Network)
    Leases a connection from the HTTP client pool for the domain. Transmits an HTTP GET request declaring a distinct User-Agent:
    User-Agent: MyCrawler/1.0 (+https://mycrawler.com/about)
    Follows HTTP redirects (301, 302, 307) up to 5 levels max. Imposes a 30-second absolute timeout and caps response sizes at 10 MB to prevent crawler bloat.
  5. Step 5: Content Deduplication (1 ms)
    Computes the 64-bit SimHash on the downloaded body. Queries the LSH candidate index using 16-bit sub-bands to retrieve candidate fingerprints sharing at least one band, then calculates the exact bitwise Hamming distance against each candidate. If the Hamming distance is ≤ 3, the page is confirmed as a near-duplicate and skipped. Otherwise, the fingerprint is written to the index and processing continues.
  6. Step 6: Parsing and Link Extraction (5 ms)
    Runs robust HTML parsers to pull structured headers, meta tags, and all child anchor links (<a href="...">) along with JSON-LD micro-data and language parameters.
  7. Step 7: Distributing Storage & Index Handoff (10 ms)
    Persists the raw HTML, headers, URL, and time metadata into GFS/S3 under a unique key generated by the page content hash:
    Write to S3: { url, content, headers, crawled_at, content_hash } (Deduplication at storage level using content_hash key).
    Once successfully written to object storage, emits a page-crawled event to Kafka containing the URL metadata and the content_hash storage reference, allowing downstream Search Engine indexers to ingest the document asynchronously without blocking the crawler.
  8. Step 8: Discovered URLs Ingestion (1 ms)
    Each extracted child URL is normalized, checked against the global seen-URLs Bloom filter, and added to the frontier with a calculated priority score if it is brand new.
    Trap Protection: Increments crawled:{domain} in Redis. Employs layered thresholds: requires a verified sitemap after 100,000 URLs, triggers an operational alert and quarantine throttling at 500,000 URLs, and halts crawl operations for that domain at 1 million pages to prevent infinite crawler traps.
  9. Step 9: Release & Back Queue Cooling
    Sets last_fetch_time = now() for example.com, locking the back queue from worker dequeues for the duration of the crawl delay.

Timeline Summary: The total time per URL is approximately 220 ms (dominated by network I/O). A single thread achieves roughly 4.5 URLs per second. To scale to 1 billion pages per day, the crawler scales up to ~2,500 active worker processes deployed across a fleet of 100 worker nodes (averaging ~25 concurrent fetch processes per node, with each node handling ~115 pages/sec).

Distributed Coordination: Domain-Sharded Architecture

In a large-scale system with ~2,500 concurrent fetch workers across 100 nodes, coordinating which worker crawls which domain is critical to prevent politeness rate violations and high locking overheads.

The Shared Frontier Concurrency Problem:

  • If multiple workers query a single centralized queue of URLs, they might fetch from the same domain at the exact same moment.
  • This issues sudden spike loads to web hosts, violating politeness agreements and generating massive distributed lock contention across workers.

The Consistent Hashing Solution:

  • We partition domains deterministically across workers using a Consistent Hashing ring with virtual nodes:
    node = ring.find_successor(hash(domain))
  • How the ring assigns and distributes domains:
    • Hash Ring Mapping: A 64-bit cryptographic hash ring places both worker nodes (with 100 to 200 virtual node tokens per physical machine) and target domain names on the same keyspace.
    • Deterministic Ownership: For any domain, the worker node located first clockwise from the domain's hash coordinates crawl operations for that domain.
    • Smooth Rebalancing: When a worker node crashes or scales out, only a fraction (1/N = ~1%) of domains migrate to adjacent nodes, avoiding the catastrophic global reshuffling of naive modulo hashing.
  • Each worker node maintains its own local URL frontier and Back Queues for its assigned subset of domains.
  • Zero-Lock Politeness for Normal Domains: Because each standard domain is owned by a single worker node at any given moment, the node manages politeness using an in-memory local token bucket with zero distributed lock overhead on the hot fetch path.

Data Ingestion, Failover, and Hot-Domain Politeness via Kafka:

  • Ingestion Pipeline: Extracted child URLs are published to the Kafka topic discovered-urls, configured with 1,024 partitions and keyed by domain hash (key = domain). Kafka provides durable, ordered buffering and decouples URL ingestion from crawl execution. Keying by domain guarantees that all URLs for a single domain arrive on the same partition in FIFO order.
  • Ingestion Routing and Hash Ring Decoupling: Worker nodes participate in a unified Kafka consumer group (frontier-ingest). While Kafka's consumer group protocol dynamically balances the 1,024 partitions across consumer nodes (~10 partitions per node) for ingestion I/O, Kafka consumer assignment does not determine domain crawl ownership. Instead, the consumer node reads newly discovered URLs, verifies them against the seen-URLs Bloom filter, and then consults the consistent-hashing ring (ring.find_successor(hash(domain))) to route deduplicated URLs to the authoritative worker node managing that domain's local frontier.
  • Local Frontier Management: The authoritative worker node receives candidate URLs for its assigned domains and enqueues them into its local frontier back queues, enforcing rate limits with zero distributed lock overhead.
  • Failover and Rebalancing: If a worker node crashes, the consistent-hashing ring smoothly shifts ownership of its domains to the next clockwise node in the ring. The new owner node retrieves crawl checkpoints and timestamps from Redis, resuming frontier execution without duplicate fetches or lost state. Separately, the Kafka consumer group detects heartbeat loss within 10 to 30 seconds and automatically rebalances partition ingestion across surviving nodes.
  • Hot-Domain Sharding with Shared Host Politeness: If an exceptionally massive domain (such as wikipedia.org) overloads a single node's queue memory or network bandwidth, the crawler splits the domain across multiple worker nodes by path prefix (e.g. wikipedia.org/wiki/A-M vs wikipedia.org/wiki/N-Z). Crucially, to prevent overloading the target server, all worker nodes crawling that host coordinate against a shared host-level rate-limit budget (persisted in Redis as a token-leasing bucket). Path sharding scales queue throughput while preserving domain-level politeness.

Advantages over Centralized Frontier:

  • ✓ No Single Point of Failure: Masterless ring architecture ensures there is no centralized coordinator or master node bottleneck to fail.
  • ✓ Linear Scalability: Adding more worker nodes dynamically expands queue memory and domain capacity across the fleet.
  • ✓ Zero-Lock Politeness for Normal Domains: Single-owner worker nodes enforce rate limits via local in-memory token buckets without distributed lock contention.
  • ✗ Hot-Domain Coordination Overhead: When massive domains are split across nodes by path prefix, workers must coordinate against a shared Redis rate-limit lease pool to protect target hosts.

BFS vs DFS vs Priority-Based Crawling

Selecting the path traversal model changes which pages are crawled first and impacts memory consumption:

  • Breadth-First Search (BFS):
    • Crawls all root seeds, then all depth-1 pages, followed by depth-2 pages.
    • ✓ Pros: Discovers crucial landing pages early, traversing from root homepages into major navigation sections.
    • ✗ Cons: Tends to waste valuable network bandwidth on low-value pages at the same depth level, and cannot prioritize important content.
  • Depth-First Search (DFS):
    • Follows links straight down a site structure before backtracking.
    • ✓ Pros: Low memory consumption since only the active page traversal path is stored.
    • ✗ Cons: Easily trapped in deep recursive page paths, and takes a long time to discover major sections of other domains.
  • Priority-Based Crawling (⭐ Recommended):
    • Scores each URL based on relative importance and crawls the highest-scored URLs first.
    • Scoring Equation:
      Score = α x PageRank(domain) + β x depth_penalty + γ x freshness_need
    • ✓ Pros: Helps prioritize high-value pages for earlier discovery and caching.
    • ✓ Adaptive Recrawling: Can reprioritize dynamic schedulers based on real-time change discovery.
    • ✗ Cons: Requires managing complex distributed priority queues.

Practical Hybrid Implementation: Start with BFS from seeds to quickly map a website's core architecture, switch to a priority-based scheduler for daily operations (focusing on important pages), and use DFS locally to walk sitemap hierarchies efficiently.

URL Normalization: Why It Matters

Websites often serve the exact same page content across multiple distinct URL formats. Normalization prevents wasting network and storage bandwidth on redundant crawls.

The Redundancy Problem:

Without normalization, the following 7 URLs would be treated as completely different entities, causing the crawler to fetch the same page 7 separate times (wasted bandwidth):

  • https://example.com/page (Canonical template)
  • https://Example.COM/page (normalizes host casing to lowercase)
  • https://example.com/page/ (strips trailing slash)
  • https://example.com/page?a=1&b=2 (sorts query parameters alphabetically)
  • https://example.com/page?b=2&a=1 (normalizes query parameters to canonical ordering)
  • http://example.com/page (normalizes protocol scheme)
  • https://example.com/./page/../page (resolves relative path segments)

Normalization Hierarchy: Safe Rules vs Configurable Heuristics:

The normalization pipeline distinguishes universally semantics-preserving rules from policy-dependent canonicalization heuristics:

  • 1. Universally Safe Normalizations (Strict Semantics Preservation):
    • Scheme and Host Lowercasing: Protocol schemes and DNS hostnames are strictly case-insensitive per RFC 3986 (e.g. HTTP://Example.COM becomes http://example.com).
    • Default Port Removal: Eliminate standard protocol ports (e.g. :80 for HTTP and :443 for HTTPS).
    • Fragment Identifier Removal: Strip URI fragments (#section) because fragments reference client-side DOM anchors and are never transmitted in HTTP requests.
    • Path Segment Resolution: Resolve dot-segments by resolving and removing inline /. and /.. path references.
    • Unreserved Percent-Encoding Decoding: Decode percent-encoded uppercase hex sequences for unreserved characters (such as converting %41 to A).
  • 2. Configurable Canonicalization Heuristics (Policy-Dependent):
    • Trailing Slash Handling: Stripping trailing slashes is safe on root hosts (e.g. example.com/), but must be configured per site on path hierarchies where servers differentiate directory indexes from file resources.
    • Query Parameter Sorting: Reordering query parameters alphabetically helps match identical parameter sets, but is policy-dependent because dynamic applications and APIs may treat parameter order as semantically significant.
    • Tracking Parameter Removal: Strip recognized marketing and analytics query keys (e.g. utm_source, utm_campaign, fbclid, gclid, ref), but preserve functional parameters required for page rendering.
    • Canonical Tag Adherence: Extract and honor HTML canonical link tags (<link rel="canonical" href="...">) when verified by the downstream parser to collapse duplicate parameter variants into a single canonical target.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...