DEV Community

Cover image for Scraping at Scale: Architecture Patterns for Millions of Pages a Day
PromptCloud
PromptCloud

Posted on

Scraping at Scale: Architecture Patterns for Millions of Pages a Day

A scraper that works for ten thousand pages does not work for ten million by running a thousand copies of itself. Somewhere on the way up, scale stops being a bigger version of the same problem and becomes a different problem, a distributed-systems problem, dominated by three forces a small scraper never feels: decoupling, politeness versus throughput, and cost. The patterns below are the ones that serve those three, and they are what separate a script that scales to a system that does.

The naive way to scale a scraper is to take the one that works and run more of it. It fails in specific, predictable ways. A thousand parallel copies with no coordination hammer the same domains and get everyone blocked. They re-fetch the same URLs because nothing tracks what has been seen. One misbehaving source backs up the whole fleet because there is no isolation. And the bill explodes, because at ten million pages a day every wasted request and every unnecessary browser render is multiplied by ten million. Scaling a scraper is not a throughput exercise, it is an architecture exercise, and here are the patterns that matter.

Pattern 1: decouple everything with a queue

The foundational move is to stop thinking of a scrape as a linear program and start thinking of it as stages connected by queues. A URL frontier holds work to be done; a pool of fetch workers pulls from it; fetched pages flow to a render stage if they need it, then to parse, validate, and store, each stage its own pool of workers connected by its own queue.

discovery -> [url frontier] -> fetch workers -> [fetched queue] ->
    render workers (only if needed) -> [parsed queue] ->
        extract + validate -> [results queue] -> store / deliver
Enter fullscreen mode Exit fullscreen mode

This decoupling buys you three things at once. Each stage scales independently, so when rendering is your bottleneck you add render workers without touching the fetch pool. Bursts are absorbed by the queues instead of overwhelming a stage. And workers can be stateless, holding no shared state locally, which means you can autoscale them horizontally on queue depth and lose one at any time without losing work. The frontier, the seen-set, and the results all live in external stores; the workers are cattle. This one pattern is the difference between a system you can grow by adding machines and a monolith you can only grow by making it bigger.

Pattern 2: be massively parallel globally, gentle per domain

Here is the tension unique to scraping at scale: you need enormous global throughput, and you need to be polite to every individual site, and those pull in opposite directions. Ten million pages a day is a firehose in aggregate, but pointed at any single domain it must be a trickle, or you take that site down and get blocked. So global concurrency and per-domain concurrency are two different controls, and you need both.

The pattern is per-domain scheduling on top of global parallelism: shard the frontier by domain, and gate each domain with its own rate limit or token bucket and a concurrency cap, so a worker may only pull the next URL for a domain when that domain's budget allows. Thousands of domains each drip politely while the fleet as a whole runs flat out. Without this, raw parallelism concentrates on whatever domains happen to be in the queue and hammers them; with it, throughput and politeness stop being in conflict.

Pattern 3: deduplicate before you spend

At millions of URLs a day, doing the same work twice is a first-order cost, so deduplication becomes infrastructure rather than a nicety. You need a fast, scalable seen-set to answer "have we already queued or fetched this URL" before spending a request on it, and at this volume an exact set can be memory-prohibitive, so a probabilistic structure like a Bloom filter is often the right trade, accepting a tiny false-positive rate to keep the check cheap and bounded. Beyond URL dedup, content hashing lets you skip re-processing pages that have not changed since last time, and conditional requests using validators like ETag or last-modified let the origin tell you a page is unchanged before you download it at all. Every one of these is a way of not spending compute, bandwidth, or a request, and at ten million a day, not spending is the whole game.

Pattern 4: assume failure is the steady state

At small scale a failure is an event. At millions a day, thousands of things fail every hour, blocks, timeouts, malformed pages, dead proxies, and the system has to treat failure as normal operation rather than an exception to be surprised by. That means several habits working together:

def worker(msg, deps):
    domain = domain_of(msg.url)

    if deps.breaker.is_open(domain):          # circuit breaker: stop hammering a burned domain
        return deps.requeue_later(msg, domain)

    try:
        resp = deps.fetch(msg.url, identity=deps.identity_for(domain))
        kind = classify(resp)                 # DATA | BLOCK | CHALLENGE | THROTTLE
        if kind != "DATA":
            deps.breaker.record_failure(domain)
            return deps.retry_or_deadletter(msg, reason=kind)

        record = extract(resp)
        if not contract.valid(record):
            return deps.deadletter(msg, reason="failed_validation")

        deps.store.upsert(record.id, record)  # idempotent: at-least-once delivery is safe
        deps.breaker.record_success(domain)

    except Exception as e:
        deps.retry_or_deadletter(msg, reason=str(e))   # backoff; poison messages end in the DLQ
Enter fullscreen mode Exit fullscreen mode

Retries with exponential backoff handle the transient failures. A dead-letter queue catches the poison messages so one un-parseable page does not block a worker forever. Per-domain circuit breakers stop the fleet from throwing requests at a domain that is returning nothing but blocks, which both saves cost and lets the domain recover. Bulkheads, isolating each source's work so one pathological site cannot starve the workers every other source depends on, keep a local fire from becoming a fleet outage. And because queues deliver at-least-once, writes must be idempotent, keyed on a deterministic ID, so a redelivered message overwrites rather than duplicates. None of these is exotic; together they are what keeps a system upright when failure is constant.

Pattern 5: watch aggregates, not jobs

You cannot supervise ten million individual jobs, so at scale the unit of observability shifts from the job to the aggregate. The control surface is a handful of fleet-level signals: throughput in pages per minute, overall success rate, queue depth and lag per stage, block rate per domain, and validation pass rate. A healthy system is a set of steady lines; an incident is a line moving, queue lag climbing because a stage is starved, block rate spiking on a domain that just tightened its defences, validation rate dropping because a source redesigned. You manage the system by its vital signs, and you alert on deviations in those, because no human is reading logs at this volume.

Pattern 6: cost is an architectural constraint

At millions a day, cost stops being something you optimise later and becomes something the architecture is shaped around, because every inefficiency is multiplied by the volume. The expensive operations are rendering and bandwidth, so the design minimises both structurally: render only the fraction of pages that genuinely need a browser and fetch the rest over plain HTTP, block images, fonts, and media so you are not paying to download pixels no parser reads, reserve expensive residential proxies for the hard targets and serve the bulk on cheap datacenter IPs, and lean on the dedup and conditional-request patterns above so you never pay twice for the same unchanged page. Efficiency here is not tuning; it is the difference between a viable unit cost and one that makes the whole operation uneconomic.

Build, or consume

Assembled, this is a serious distributed system: a durable frontier, staged autoscaling workers, per-domain schedulers, a dedup layer, circuit breakers and dead-letter queues, fleet observability, and a cost-control regime, all maintained as targets and volumes change. It is entirely buildable, and some teams should build it because crawling at scale is their core competency. For many others it is undifferentiated heavy lifting, which is the case for treating large-scale crawling as a consumed data crawling service rather than a platform to operate: the patterns here become someone else's operational burden and you receive the data at the end. Either way, the point stands that at this volume the interesting engineering is the system, not the scraper.

The takeaway

Scaling to millions of pages a day is not about making a scraper faster, it is about making it a distributed system: decouple stages with queues so they scale independently and contain failure, schedule per domain so you are polite and parallel at once, deduplicate and use conditional requests so you never spend twice, treat failure as the steady state with retries and breakers and idempotent writes, observe the fleet by its aggregate vital signs, and shape the whole thing around cost because every waste is multiplied by the volume. Get those patterns right and adding another few million pages a day is a matter of adding workers, which is exactly the property you were trying to buy.

FAQ

How do you scrape millions of pages a day without getting everything blocked?

By separating global throughput from per-domain politeness. The fleet runs massively parallel in aggregate, but the frontier is sharded by domain and each domain is gated by its own rate limit and concurrency cap, so no single site is hammered even as total volume is enormous. This is layered on top of the usual per-request concerns like coherent identity and proxy management, but the scale-specific insight is that raw parallelism must be constrained per domain, or high global throughput simply concentrates into taking individual sites down and getting your fleet banned.

What is the most important architecture pattern for large-scale scraping?

Decoupling the pipeline with queues. Turning discovery, fetching, rendering, parsing, and storage into separate stages connected by durable queues is what lets each stage scale independently, lets workers be stateless and autoscale on queue depth, absorbs bursts, and contains failures rather than letting one bad stage stall everything. Almost every other scale pattern, per-domain scheduling, dead-letter queues, circuit breakers, idempotent writes, depends on that queue-based decoupling being in place first. Without it you have a big monolith; with it you have a system you grow by adding machines.

Top comments (0)