Designing a Web Crawler and Search Engine
Difficulty: Advanced Prerequisites:Message Queues, Consistent Hashing, and Bloom Filters
TL;DR
Two coupled systems in one question: a crawler that discovers and downloads billions of pages, and a search service that indexes them and answers queries in milliseconds. They are joined only by storage β the crawler writes pages, the indexer reads them β which is what lets them scale independently.
The shape of the answer: a two-level URL frontier for politeness, a Bloom filter for dedup, an offline indexing pipeline, and a sharded inverted index fronted by a query aggregator.
The single most important framing: crawling is a write-heavy batch problem, serving is a read-heavy latency problem. Treat them as one system and you will design both badly.
Understanding the Problem
Two coupled systems: a crawler that discovers and downloads billions of web pages, and a search service that indexes that content and answers user queries in milliseconds. The hard parts: crawling politely without hammering any single site, avoiding re-crawling unchanged content, building a massive inverted index, and serving ranked results at sub-200ms latency.
Crawling is throughput-bound and can be arbitrarily slow per page as long as aggregate rate holds. Serving is tail-latency-bound, where a single slow shard ruins the query. Because they share nothing but the page store, the crawler can fall hours behind without the search path noticing.
Naive First Cut
flowchart LR
SEED["Seed URLs"]:::client
CRAWLER["Single Crawler"]:::service
DB[("One DB<br/>pages + index")]:::data
USER["Search User"]:::client
SEED --> CRAWLER
CRAWLER --> DB
USER --> DB
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Why this breaks:
- Single crawler β at 1 page/sec, crawling 10B pages takes 317 years
- No URL deduplication β the same page is fetched thousands of times via different links
- No politeness β hammering a single domain causes your IP to be blocked
- One database canβt hold billions of documents AND serve as a search index
- No ranking β results are returned in insertion order, not by relevance
The rest of the doc splits this into two independently scalable halves.
Prior Art Weβre Drawing From
- Mercator / the original Google crawler β introduced the two-level frontier (front queues for priority, back queues for per-domain politeness) that remains the standard answer to βhow do you crawl fast without hammering anyone.β (The Anatomy of a Large-Scale Hypertextual Web Search Engine)
- Apache Nutch / Heritrix β open-source crawlers that show the frontier, fetcher, parser, and dedup as separable services with durable handoffs between them. (Apache Nutch)
- Googleβs index serving β document-sharded index with a query aggregator scatter-gathering across shards, which is why P99 depends on the slowest shard rather than the average. (Web Search for a Planet)
- Bigtable β built in large part to store crawled pages and their versions at web scale; the canonical answer for βwhere do billions of pages live.β (Bigtable paper)
Functional Requirements
Core (top 3)
- Crawl the web β discover, download, and store web pages at scale (1B+ pages)
- Build an inverted index β map every word to the pages that contain it
- Serve search queries β return the top 10 most relevant results for a query in <200ms
Below the Line
- Image/video indexing, real-time freshness (news), autocomplete, personalization, ads
Non-Functional Requirements
- Crawl throughput β 10K pages/second sustained
- Index freshness β popular pages re-crawled within hours; long-tail within weeks
- Query latency β <200ms P99 for search results
- Politeness β respect robots.txt; max 1 request/second per domain
Below the Line
- Real-time indexing (seconds-fresh news)
- Personalized or session-aware ranking
- JavaScript rendering for SPA-heavy sites (a large, separate problem)
Scale Estimation (Back-of-Envelope)
Crawl side:
- Throughput: 10K pages/sec β 864M pages/day β 1B pages in ~1.2 days, 10B in ~12 days
- Fleet size: at ~1 page/sec per worker-domain pair with network waits dominating, 10K pages/sec needs on the order of 10K concurrent fetches β a few hundred workers with high concurrency each
- Raw storage: 10B pages Γ ~75KB compressed HTML β 750TB, which is object-storage/Bigtable territory, not a database
- Bandwidth: 10K pages/sec Γ 100KB β 1GB/sec β 8 Gbps sustained ingress
- Dedup filter: 10B URLs in a Bloom filter at 1% false-positive β 12GB RAM β fits in memory, which is the whole reason to use one
Serve side:
- Index size: ~10B docs; posting lists after compression β 100s of TB, sharded across thousands of machines
- Query QPS: assume 100K queries/sec; top 1% of queries drive ~30% of traffic, so cache hit rate is high
- Fan-out per query: a query hits every document shard, so P99 is governed by the slowest shard, not the average β this is why tail latency dominates the serving design
The 750TB figure is the one to state early: it immediately rules out βstore pages in Postgresβ and justifies object storage plus a separate index.
Core Entities
- URL β address to crawl, domain, last crawled timestamp, crawl priority
- Page β raw HTML content, extracted text, outgoing links, content hash
- Inverted Index Entry β word β list of (page ID, position, frequency)
- Query β user search terms, results with relevance scores
API / System Interface
POST /v1/crawl/seed
Body: { urls: ["https://example.com", ...] }
Response: { queued: 150 }
GET /v1/search?q=distributed+systems&page=1
Response: { results: [{ url, title, snippet, score }], total: 15000, took: "45ms" }
GET /v1/crawl/status
Response: { pagesIndexed: 1200000000, crawlRate: "9800 pages/sec" }
High-Level Design
FR1: Crawl the Web
A crawler is a loop. Take a URL, download it, save what came back, harvest the links on the page, and put the ones you have not seen into the pile. Read literally, the requirement asks for that loop and somewhere to keep the pile.
Build exactly that. One table as the pile, a fleet of workers running the loop, and object storage for the HTML. No Bloom filter, no priority scoring, no per-domain scheduling. Every one of those answers a non-functional requirement β throughput, freshness, politeness β so each is earned in a deep dive rather than asserted here.
New components we need:
- Crawler Workers β run the loop. Claim a URL, fetch it over HTTP, write the HTML, parse out the links, and feed the unseen ones back.
urlstable (Postgres) β the frontier. One row per URL with astatusofpendingordone. A unique index on the URL is what stops us crawling the same address twice.- Page Store (object storage) β the raw HTML, keyed by a hash of the URL. 10B pages at ~75KB compressed is roughly 750TB, which is object storage territory and not a database.
flowchart LR
SEED["Seed URLs"]:::client
DB[("Postgres urls table<br/>pending and done")]:::data
W["Crawler Workers"]:::service
STORE[("Page Store<br/>object storage")]:::data
SEED -->|"1. Insert as pending"| DB
DB -->|"2. Claim oldest pending row"| W
W -->|"3. Download and write HTML"| STORE
W -->|"4. Insert each unseen link"| DB
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step flow:
- Seed URLs are inserted into
urlswithstatus = pending - A worker claims the oldest pending row and marks it
in_progress, so two workers do not fetch the same page - The worker issues an HTTP GET and writes the response body to the Page Store under
hash(url) - It parses the HTML and extracts every outbound link
- For each link it attempts an insert into
urls; the unique index rejects anything already known, and that rejection is our deduplication - The original row flips to
doneand the worker claims the next one
Why store the raw HTML rather than just the extracted text? Parsing rules change. If tokenisation improves, or we later want to pull out tables and captions we ignored the first time, re-deriving from stored HTML costs a batch job. Re-deriving from the live web means crawling 10B pages again, which takes twelve days and annoys every site we hit.
What we have deliberately left broken. For a few million pages this is a real crawler and it works. Three things are wrong with it at the scale the requirements state, and none of them is a missing feature:
- Step 5 is the bottleneck. At 10K pages/sec with roughly 50 links per page, that is about 500K insert attempts per second against a table holding 10B rows, purely to answer βhave I seen this?β. No index makes that cheap, and the answer does not need to be exact. That is Deep Dive 1.
- Nothing is polite. Step 2 hands out whatever is oldest, so if a page has 200 links to the same site, 200 workers hit that host at once. Real sites read that as an attack and block our IP, and we lose the domain entirely. The requirements ask for one request per second per domain and this design cannot express that. That is Deep Dive 2.
- Nothing ever gets re-crawled. A URL goes
pending β doneand is never revisited, so the index decays from the moment it is written. That is Deep Dive 5.
FR2: Build the Inverted Index
Stored HTML is useless for search. Nobody wants to grep 750TB per query, so we invert it:
instead of page β words, keep word β pages. Then a search for a term is a lookup rather
than a scan.
π‘ Inverted index = a map from each word to the list of pages containing
it, called a posting list. It is the same idea as the index at the back of a textbook.
New components we need:
- Indexer β reads pages the crawler has stored but not yet indexed, splits the text into terms, and appends each page to that termβs posting list.
- Inverted Index β one store mapping term β posting list of
(page_id, term_frequency, positions).
No message bus here. The urls table already records which pages are indexed, so the
Indexerβs work queue is a query against a column we are keeping anyway.
flowchart LR
STORE[("Page Store")]:::data
IDX["Indexer"]:::service
INV[("Inverted Index<br/>single store")]:::data
DB[("Postgres urls table")]:::data
DB -->|"1. Find rows done but unindexed"| IDX
STORE -->|"2. Read the stored HTML"| IDX
IDX -->|"3. Append page to each term"| INV
IDX -->|"4. Mark the row indexed"| DB
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step flow:
- The Indexer queries
urlsfor rows wherestatus = doneandindexed_at is null - For each, it reads the stored HTML from the Page Store
- It strips markup, lowercases, and splits into terms, counting frequency and recording positions
- For every term it appends
(page_id, frequency, positions)to that termβs posting list - It sets
indexed_at, so the same page is not indexed twice
Why keep positions and not just a term count? Positions are what make phrase search
possible. Without them "new york times" matches any page containing those three words
anywhere, including a page about new times in York.
What we have deliberately left broken. This produces a correct index and it will not survive contact with the numbers:
- The index does not fit anywhere. Posting lists over 10B documents run to hundreds of terabytes after compression. There is no single machine to put that on, so the store above is a fiction the moment the crawl gets past a few hundred million pages. How to split it, and what that costs at query time, is Deep Dive 4.
- One Indexer cannot keep up. The crawler is producing 10K pages/sec and a single consumer reading them serially will fall permanently behind, so the index is always stale by a growing margin. Splitting the index in Deep Dive 4 is also what lets us parallelise this.
FR3: Serve Search Queries
The requirement is ten relevant results for a set of words. With the index built, that is a lookup per term, an intersection, an ordering, and a cut at ten.
New components we need:
- Query Service β parses the query, reads one posting list per term, intersects them, scores the survivors, and returns the top ten.
That is the whole addition. No cache: until we have said why reading the index directly is
not good enough, a cache is machinery without a justification. Ranking is TF-IDF only,
because it is the ordering that falls out of data we already stored in FR2.
π‘ TF-IDF =
term frequency Γ inverse document frequency. A word counts for more if it appears often on
a page and rarely across the whole corpus, which is why matching on βtheβ tells you
nothing.
flowchart LR
USER["User"]:::client
Q["Query Service"]:::service
INV[("Inverted Index")]:::data
USER -->|"1. Submit query terms"| Q
Q -->|"2. Read one posting list per term"| INV
INV -->|"3. Return the posting lists"| Q
Q -->|"4. Intersect score and return ten"| USER
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step flow:
- User submits
distributed systems; the Query Service lowercases and splits it into two terms - It reads the posting list for each term
- It intersects them, keeping only pages that contain both
- It scores each surviving page by TF-IDF and sorts
- It returns the top ten with a snippet pulled from the stored HTML
Why intersect rather than union? A union returns every page containing either word,
which for distributed systems is most of the English-language web. Intersection is what
makes a multi-word query narrower than a single-word one, which is what users expect.
What we have deliberately left broken. This answers queries and the answers are bad:
- TF-IDF alone is trivially gamed. It rewards a page for repeating a word, so a page that says βdistributed systemsβ four hundred times outranks the canonical paper on the subject. Nothing here measures whether anyone considers the page worth linking to. That is Deep Dive 3.
- Every query touches every termβs full posting list. At 100K queries/sec against a 200ms P99 budget, reading and intersecting hundreds of megabytes of postings per query is not close to affordable, and the popular terms are the expensive ones. That is Deep Dive 4, alongside the tail latency that splitting the index introduces.
Technology Choices
| Tier | Purpose | Primary Pick | Alternatives | Why Primary Wins Here |
|---|---|---|---|---|
| URL Frontier | Priority queue + politeness queue | Redis Sorted Sets + per-domain queues | Kafka, RabbitMQ, custom on-disk queue | ZSET for priority ordering; per-domain keys enforce politeness naturally |
| URL Dedup | Avoid re-crawling same URL | Bloom Filter (in-memory) | Redis set, RocksDB, HyperLogLog | 10B URLs in ~12GB RAM with 1% false positive; no disk I/O on hot path |
| Page Store | Raw crawled HTML pages | S3 or Bigtable | HDFS, GCS, Cassandra | 750TB of pages needs cheap, durable object storage β not a database |
| Inverted Index | Term β document posting lists | Custom sharded index (Lucene-based) | Elasticsearch, Solr, Vespa | Document-sharded for scatter-gather; Lucene segment format is proven at web scale |
| Content Dedup | Detect near-duplicate pages | SimHash / MinHash | Exact MD5 hash, locality-sensitive hashing | Catches near-duplicates (same content + different ads) not just exact copies |
| Crawl Metadata | URL state, last crawl time, priority | RocksDB (embedded) or Postgres | DynamoDB, Cassandra, LevelDB | Fast local lookups for crawler workers; Postgres for coordinator state |
Why Bloom Filter over a hash set for URL dedup? At 10B URLs, a hash set needs ~400GB RAM (40 bytes/entry). A Bloom filter with 1% false-positive needs only 12GB. The trade-off (occasionally re-crawling 1% of pages) is acceptable since the crawl is already idempotent.
Data Modeling
Redis Sorted Sets (URL Frontier β priority queue):
Key: "frontier:priority" β Sorted Set (score = priority_score, member = url_hash)
Key: "frontier:domain:{domain}" β List (per-domain FIFO queue for politeness)
Key: "frontier:domain_delay:{domain}" β String (next_allowed_crawl_time)
Bloom Filter (URL Dedup β 10B URLs in 12GB RAM):
In-memory: 10-element hash Bloom filter
Capacity: 10 billion URLs
False positive rate: 1%
Memory: ~12GB
Operations: probe before adding to frontier; insert after crawl succeeds
S3 / Bigtable (Page Store β raw crawled HTML):
Key: "pages/{url_hash}" β Raw HTML content (avg 50KB per page)
Metadata: { url, content_hash, crawled_at, http_status, content_type, size_bytes }
Total: ~750TB for 10B pages
Custom Sharded Index (Inverted Index β search results):
Segment format (Lucene-based):
Term β PostingList [ (doc_id, term_frequency, positions[]) ]
DocValues: pagerank_score, freshness, domain_authority
Sharded by document_id_range across N machines (scatter-gather on query)
Access Patterns:
| Query | Data Source | How |
|---|---|---|
| Get next URL to crawl | Redis | ZPOPMIN frontier:priority β highest priority URL, then check domain delay |
| Check if URL already crawled | Bloom Filter | Probe in-memory β βdefinitely notβ β enqueue; βmaybe yesβ β skip |
| Store crawled page | S3/Bigtable | PUT pages/{url_hash} with raw HTML + metadata |
| Search query | Inverted Index | Scatter query to N shards, each returns top-K scored docs, merge globally |
| Detect near-duplicate pages | SimHash | Compute SimHash of page content β compare Hamming distance with known pages |
How the URL Frontier Manages Politeness and Priority:
- New URLs discovered from a crawled page β check Bloom filter (reject already-seen)
- Compute priority score:
base_priority + domain_authority + freshness_bonus - depth_penalty ZADD frontier:priority <score> <url_hash>β priority queue insert O(log N)- Crawler worker:
ZPOPMIN frontier:priorityβ gets highest-priority URL - Before fetching: check
frontier:domain_delay:{domain}β if current_time < next_allowed, re-enqueue and pick another - After successful crawl: set
frontier:domain_delay:{domain}= now + robots.txt crawl-delay (default 1s) - This ensures no domain gets hit more than once per second (politeness) while high-priority URLs are crawled first globally
Deep Dives
1) How do we answer βhave I already crawled this?β 500K times a second?
Bad: What FR1 built β let the urls table decide, by attempting an insert and letting a unique index reject duplicates. At 10K pages/sec and roughly 50 links per page that is about 500K insert attempts per second against a table of 10B rows. Every one is a B-tree descent through an index far too large to cache, so nearly all of them become disk reads, and the failed inserts still burn a round trip each. Swapping it for an in-memory hash set does not rescue it: 10B URLs at ~100 bytes each is about 1TB of RAM, so it does not fit on one machine either. There is also a correctness gap underneath the performance one β example.com/a, example.com/a/, and example.com/a#top are the same page and all three are separate rows.
Good: Use a Bloom filter for fast membership testing (probabilistic: may say βseenβ for an unseen URL, but never misses a seen one). False positive rate of 1% at 10B URLs needs ~12GB β fits in memory. Normalize URLs (lowercase, remove fragments, sort params) before checking.
Great: Combine the Bloom filter with content-based dedup. After downloading a page, compute a SimHash (locality-sensitive hash) of the content. Two pages with similar content (mirror sites, syndicated articles) get the same SimHash β dedup at the content level, not just URL level. This eliminates near-duplicate pages that have different URLs but identical content, reducing index bloat by 30-40%.
2) A page links to the same site 200 times. How do we stop 200 workers hitting it at once?
Bad: What FR1 built β workers claim whatever row is oldest, with no notion of which host a URL belongs to. Links cluster by domain, because that is what site navigation is, so a single page with 200 internal links puts 200 rows for one host at the front of the queue together. The requirement is one request per second per domain; this design issues 200 in whatever time the fleet takes to drain them, which is a fraction of a second. From the far end that is indistinguishable from a denial-of-service attack, and the response is a 429 and then an IP ban. We do not lose one page, we lose the whole domain, and the ban outlives the mistake.
Good: Stop treating the pile as one queue. Replace the flat urls table with a URL
Frontier that keeps a separate queue per domain, each carrying a not_before timestamp,
and let workers claim only from domains that are currently eligible. Politeness stops being
something workers have to remember and becomes a property of what the frontier will hand
out. Fetch and cache each siteβs robots.txt, and honour its Crawl-delay where it sets
one rather than assuming our own one-second default is welcome.
Great: Use a two-level frontier. Back queue: per-domain queues with rate limiting (politeness). Front queue: priority queue that selects which domain to crawl next based on importance (PageRank of domain, freshness requirements). This ensures high-value domains (news sites, Wikipedia) are crawled frequently while staying polite. Assign each worker a set of domains via consistent hashing β this ensures DNS caching efficiency and persistent connections per worker-domain pair.
3) Why does a page repeating a phrase 400 times outrank the paper that defined it?
Bad: What FR3 built β TF-IDF and nothing else. Every signal it uses is written by the author of the page being ranked, which makes the ranking a measure of how badly someone wants to rank rather than of whether the page is any good. Repeating a phrase is free, so the ordering is decided by whoever is least embarrassed to do it. This is not a theoretical weakness; it is the entire business model of keyword-stuffing, and it dominated search results until ranking started using signals the author does not control.
Good: Combine TF-IDF with PageRank β a pageβs authority is proportional to the number and quality of pages linking to it. This boosts authoritative sources (Wikipedia, official docs) above spam. PageRank is computed offline as a batch job over the link graph.
Great: Multi-signal ranking: TF-IDF (text relevance) + PageRank (authority) + freshness (prefer recent content for time-sensitive queries) + click-through rate (learn from user behavior over time). The scoring formula is a weighted combination, tuned via ML. For query latency, pre-compute static scores (PageRank) and combine with query-time scores (TF-IDF) during serving. Cache results for popular queries (top 1% of queries account for 30% of traffic) with a 1-minute TTL.
4) The index is hundreds of terabytes. Do we split it by term or by document?
Bad: What FR2 built β one inverted index in one store. Posting lists over 10B documents come to hundreds of terabytes compressed, so this stops being deployable somewhere around a few hundred million pages, which the crawler reaches in under a week at 10K pages/sec. It also caps indexing throughput at whatever one writer can absorb while the crawler produces 10K pages/sec, so the index falls behind by a margin that only grows. The single store was worth writing down because it makes the next question precise: the index has to be split, and the only real decision is along which axis.
Good: Term-partitioned (each shard owns a subset of words). A query for βdistributed systemsβ touches only the two shards owning those terms, so per-query fan-out is small. The problem: posting lists for common words are enormous and unevenly distributed, so the shard owning βtheβ is a permanent hotspot, and multi-term queries must ship huge posting lists across the network to be intersected.
Great: Document-partitioned (each shard owns a subset of pages and indexes all their words). Every query goes to every shard, each computes its own local top-K, and an aggregator merges the results. Fan-out is wide but each shard does a small, bounded amount of work with no cross-shard data movement, and load is naturally even. This is what real search engines do.
The cost is tail latency: with 1000 shards, the query is as slow as the slowest one. Mitigate with hedged requests β if a shard hasnβt replied by the 95th-percentile mark, send a duplicate request to a replica and take whichever answers first. Combined with per-shard tiering (put high-PageRank documents in a small βtop tierβ that is searched first, and only fall through to the long-tail tier if there arenβt enough good results), most queries never touch the full index at all.
5) A URL goes pending then done and is never seen again. How do we stop the index rotting?
Bad: What FR1 built β nothing re-crawls. A URL reaches done and is finished forever,
so the index is a photograph of whenever each page happened to be fetched. Prices, headlines
and documentation drift out from under it, and dead pages stay in the results permanently
because nothing ever revisits them to notice the 404. The obvious first fix, re-crawling
everything on a fixed cycle of say 30 days, trades one failure for two: most of the budget
goes on pages that have not changed since the last sweep, and a news homepage that changes
every few minutes is still 30 days stale. The requirements ask for hours on popular pages
and weeks on the long tail, and a single global interval cannot express both.
Good: Tier by observed change rate. Track a content hash per URL across crawls; pages that changed last time get a shorter interval, unchanged pages get a longer one, with exponential backoff up to a cap. Use HTTP conditional requests (If-Modified-Since / ETag) so unchanged pages cost a 304 response instead of a full download β often a 10x bandwidth saving on the re-crawl path.
Great: Treat it as budget allocation against a freshness objective. Estimate each pageβs change rate (a Poisson model fits well) and weight it by importance (PageRank, observed query and click traffic). Spend the crawl budget where P(changed) Γ importance is highest, so a news homepage is crawled every few minutes while an archived page drops to yearly. Sitemaps and lastmod hints, plus push protocols like IndexNow, let cooperative sites tell you what changed instead of you polling to find out.
Design Self-Audit
| Question | Answer |
|---|---|
| Crawler traps? | Infinite calendars and session-ID URLs generate unbounded links. Defenses: cap URL depth and length, cap pages per domain, strip known session params during normalization, and detect near-duplicate content via SimHash. |
| A worker dies mid-fetch? | The frontier uses visibility-timeout semantics (like SQS): an unacknowledged URL becomes eligible again after a timeout. Worst case a page is fetched twice, which is harmless β the crawl is idempotent. |
| Bloom filter false positive? | A never-crawled URL is wrongly marked βseenβ and is silently skipped. At 1% thatβs acceptable for web crawling β those pages are almost always reachable via other links later. It is a deliberate trade of completeness for 12GB instead of 1TB. |
| Bloom filter canβt delete? | Correct β standard Bloom filters donβt support removal, and it fills over time. Use a rotating/scalable variant, or rebuild periodically from the authoritative URL store. |
| Indexing falls behind? | Search keeps serving the older index; only freshness suffers. Because indexing is offline and decoupled via the page store, crawler and indexer backlogs never affect query availability. |
| Index rebuild without downtime? | Build the new index generation alongside the live one and atomically flip an alias once it passes validation. Never mutate the serving index in place. |
| robots.txt fetch fails? | Fail closed β donβt crawl the domain until robots.txt is retrievable, and cache it with a TTL. Guessing βallowedβ risks legal and reputational problems. |
Final Architecture
flowchart TB
SEED["Seed URLs"]:::client
FRONTIER["URL Frontier<br/>front priority queues<br/>back per-domain queues"]:::async
FETCH["Fetcher Workers"]:::service
ROBOTS[("robots.txt<br/>cache")]:::data
PARSE["Parser<br/>extract text and links"]:::service
DEDUP["Dedup<br/>Bloom filter plus SimHash"]:::service
STORE[("Page Store<br/>object storage 750TB")]:::data
KAFKA["Kafka<br/>new and changed pages"]:::async
INDEXER["Indexer"]:::service
PAGERANK["PageRank<br/>batch job"]:::async
INDEX[("Inverted Index<br/>document sharded")]:::data
USER["Search User"]:::client
AGG["Query Aggregator"]:::service
CACHE[("Redis<br/>hot queries")]:::data
SEED --> FRONTIER
FRONTIER -->|"1. Dequeue eligible URL"| FETCH
FETCH -->|"2. Check crawl rules"| ROBOTS
FETCH -->|"3. Store raw page"| STORE
FETCH -->|"4. Hand off content"| PARSE
PARSE -->|"5. Check if seen"| DEDUP
DEDUP -->|"6. Enqueue new URLs"| FRONTIER
STORE -->|"7. Change events"| KAFKA
KAFKA -->|"8. Tokenize and index"| INDEXER
PAGERANK -->|"9. Static authority scores"| INDEXER
INDEXER -->|"10. Write posting lists"| INDEX
USER -->|"11. Submit query"| AGG
AGG -->|"12. Check hot cache"| CACHE
CACHE -.->|"miss: scatter gather"| INDEX
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
| Color | Meaning |
|---|---|
| π£ Indigo | Clients (Seed URLs, Search User) |
| π’ Green | Services |
| π‘ Yellow | Data stores |
| πͺ Violet | Queues / async batch jobs |
How it works end-to-end (crawl path):
- Frontier selects the next URL β front queues decide which domain deserves attention by priority; back queues enforce when that domain may next be hit
- Fetcher checks crawl rules β robots.txt is fetched once per domain and cached, honoring
Crawl-delay - Raw page lands in object storage β 750TB of HTML belongs in a blob store, not a database
- Parser extracts text and outgoing links β the links are the crawlerβs fuel; this is how discovery continues
- Dedup filters whatβs already known β Bloom filter for URL-level dedup, SimHash for near-duplicate content
- New URLs return to the frontier, closing the discovery loop
How it works end-to-end (index and serve path):
- Page store changes publish to Kafka β this queue is the seam that decouples crawling from indexing, so either side can lag without breaking the other
- Indexer tokenizes and builds posting lists β entirely offline, never in a userβs request
- PageRank supplies static authority scores β computed as a batch job over the link graph and folded in at index time, so it costs nothing at query time
- Posting lists written to a document-sharded index β each shard owns a slice of the corpus
- Query aggregator receives the user query
- Hot query cache answers ~30% of traffic outright; on a miss the aggregator scatter-gathers across shards, merges each shardβs local top-K, and returns the global top 10
Key Technologies
| Term | What it is |
|---|---|
| URL Frontier | The crawlerβs scheduler. Two-level: front queues for priority, back queues for per-domain rate limiting. The heart of the crawl design. |
| Bloom Filter | Probabilistic set membership in constant space. Answers βseen this URL?β for 10B URLs in 12GB instead of 1TB. Can false-positive, never false-negative. |
| SimHash | A locality-sensitive hash where similar content yields similar hashes. Catches mirrored and syndicated pages that URL dedup misses. |
| Inverted Index | Maps each term to a posting list of documents containing it. The core search data structure. |
| Document Sharding | Each shard indexes a subset of pages. Wide fan-out per query but even load and no cross-shard traffic. |
| TF-IDF / BM25 | Text relevance scoring: term frequency weighted down by how common the term is corpus-wide. |
| PageRank | Link-graph authority β a page is important if important pages link to it. Computed offline as a static per-page score. |
| Hedged Requests | Duplicate a slow request to a replica and take the first response. Cuts P99 when a query fans out to many shards. |
Whatβs Expected at Each Level
| Level | Expectations |
|---|---|
| Mid | URL Frontier + crawler fleet. Bloom filter for dedup. Inverted index concept. Basic TF-IDF ranking. Robots.txt politeness. |
| Senior | Per-domain rate limiting with two-level frontier. Content-based dedup (SimHash). PageRank for authority. Index sharding by term or document. |
| Staff+ | Consistent hashing for worker-domain affinity. Multi-signal ML ranking. Incremental index updates (not full rebuild). Freshness-based re-crawl prioritization. Cache strategy for query serving. |
π― Key Takeaways
- These are two systems, not one. Crawling is write-heavy and throughput-bound; serving is read-heavy and tail-latency-bound. The page store and Kafka are the seam that lets each scale and fail independently
- The frontier is the design. Two levels β priority in front, per-domain politeness behind β is the answer to crawling fast without getting blocked
- Bloom filter trades 1% completeness for 100x memory (12GB vs 1TB). Knowing what you gave up matters more than naming the structure
- Dedup twice: by URL (Bloom) and by content (SimHash), because mirrors and syndicated articles have different URLs and identical text
- Shard the index by document, not by term β wide fan-out but even load; then fight the resulting tail latency with hedged requests and tiering
- Do expensive work offline. PageRank, indexing, and re-crawl scheduling are batch jobs; query time only combines pre-computed scores
- Re-crawling is budget allocation, not a fixed cycle β spend where
P(changed) Γ importanceis highest
Related Designs
- Search Autocomplete - prefix serving in front of a search system
- Q&A Forum - full-text search and relevance ranking at smaller scale
- News Aggregator - content ingestion and freshness-based ranking
- Job Scheduler - the scheduling machinery behind re-crawl timing
- Rate Limiter - the per-domain politeness mechanism, generalized
Related Concepts
Understand the building blocks used in this design:
- Bloom Filters β β URL dedup for 10B URLs in 12GB of RAM
- Message Queues β β the frontier and the crawl-to-index seam
- Consistent Hashing β β pins domains to workers for DNS and connection reuse
- Object Storage β β holds 750TB of raw pages cheaply
- Database Sharding β β document-partitioning the inverted index
- Batch vs Stream β β PageRank and indexing as offline jobs
- Rate Limiting β β per-domain crawl politeness
- Retry & Backoff β β handling fetch failures and unreachable hosts
Discussion
Newest first