Building a Web Scraping Pipeline: Architecture, Scheduling and Data Processing
On this page
Beyond One-Off Scripts
A scraping script that runs once and dumps results into a CSV is fine for a one-time data pull. But when you need data regularly, from multiple sources, with consistent quality and reliable delivery, you need a pipeline.
A scraping pipeline is an automated system that collects data from websites on a schedule, processes and cleans it, stores it in a usable format and alerts you when something breaks.
Pipeline Architecture
Components
A production scraping pipeline has six stages:
- Scheduler. Triggers scraping jobs on a defined schedule.
- Crawler. Fetches web pages and handles navigation, pagination and authentication.
- Parser. Extracts structured data from raw HTML.
- Processor. Cleans, normalises and validates the extracted data.
- Storage. Saves processed data to a database, file system or downstream service.
- Monitor. Tracks job status, data quality and alerts on failures.
Flow
Scheduler -> Crawler -> Parser -> Processor -> Storage
|
Monitor
Each stage is a separate concern. When the website changes its HTML, only the parser needs updating. When you add a new data source, you add a new crawler and parser without changing the processor or storage.
Stage 1: Scheduling
Options
Cron (Linux). Simple, reliable, built into every Linux server.
# Run every day at 2 AM
0 2 * * * /usr/bin/python3 /opt/scraper/run.py
Task schedulers. Airflow, Prefect, Dagster. These provide dependency management, retries, logging and a web interface.
Cloud schedulers. AWS EventBridge, Google Cloud Scheduler, Azure Logic Apps. Trigger serverless functions or containers.
Scheduling decisions
Frequency. How often does the source data change? A job board updates hourly. A company directory updates weekly. A government database updates quarterly. Match your schedule to the source's update frequency.
Time of day. Schedule during off-peak hours for the target website (typically 2-5 AM in their timezone). This reduces impact on their servers and reduces the chance of rate limiting.
Staggering. If you scrape multiple sources, stagger the start times. Running all scrapers simultaneously creates resource spikes.
Retry policy. Define what happens when a scrape fails. Retry immediately? Retry after 30 minutes? Alert and wait for manual intervention?
Stage 2: Crawling
Simple HTTP requests
For static HTML pages:
import requests
from time import sleep
class Crawler:
def __init__(self, base_url, delay=2):
self.session = requests.Session()
self.session.headers.update({
'User-Agent': 'CompanyBot/1.0 (contact@example.com)'
})
self.delay = delay
def fetch(self, url):
sleep(self.delay)
response = self.session.get(url, timeout=30)
response.raise_for_status()
return response.text
def fetch_paginated(self, url_template, max_pages=100):
pages = []
for page in range(1, max_pages + 1):
url = url_template.format(page=page)
try:
html = self.fetch(url)
pages.append(html)
if self._is_last_page(html):
break
except requests.HTTPError as e:
if e.response.status_code == 404:
break
raise
return pages
Headless browser for JavaScript-rendered pages
When pages load content via JavaScript:
from playwright.sync_api import sync_playwright
class BrowserCrawler:
def __init__(self, headless=True):
self.playwright = sync_playwright().start()
self.browser = self.playwright.chromium.launch(headless=headless)
self.context = self.browser.new_context()
def fetch(self, url, wait_selector=None, timeout=30000):
page = self.context.new_page()
try:
page.goto(url, timeout=timeout)
if wait_selector:
page.wait_for_selector(wait_selector, timeout=timeout)
return page.content()
finally:
page.close()
def close(self):
self.browser.close()
self.playwright.stop()
Handling rate limits and blocks
Delays. Wait 1-5 seconds between requests. Randomise the delay.
Rotating proxies. When scraping at volume, use rotating residential or datacenter proxies to distribute requests across IP addresses.
Request headers. Set realistic User-Agent, Accept, Accept-Language and Referer headers.
Session management. Reuse sessions (cookies, headers) across requests to the same domain.
Exponential backoff. When you get a 429 (Too Many Requests) or 503 (Service Unavailable), wait and retry with increasing delays.
import time
import random
def fetch_with_retry(session, url, max_retries=3):
for attempt in range(max_retries):
try:
response = session.get(url, timeout=30)
if response.status_code == 429:
wait = (2 ** attempt) + random.uniform(0, 1)
time.sleep(wait)
continue
response.raise_for_status()
return response.text
except requests.RequestException as e:
if attempt == max_retries - 1:
raise
time.sleep(2 ** attempt)
return None
Respecting robots.txt
Check robots.txt before scraping. Respect Crawl-delay directives.
from urllib.robotparser import RobotFileParser
def check_robots(url, user_agent='*'):
parser = RobotFileParser()
parser.set_url(f"{url}/robots.txt")
parser.read()
return parser.can_fetch(user_agent, url)
See Web Scraping Ethics and Best Practices.
Stage 3: Parsing
HTML parsing with Beautiful Soup
from bs4 import BeautifulSoup
class DirectoryParser:
def parse(self, html):
soup = BeautifulSoup(html, 'html.parser')
results = []
for listing in soup.select('.business-listing'):
record = {
'name': self._text(listing, '.business-name'),
'email': self._text(listing, '.email a'),
'phone': self._text(listing, '.phone'),
'address': self._text(listing, '.address'),
'website': self._attr(listing, '.website a', 'href'),
}
results.append(record)
return results
def _text(self, element, selector):
found = element.select_one(selector)
return found.get_text(strip=True) if found else None
def _attr(self, element, selector, attr):
found = element.select_one(selector)
return found.get(attr) if found else None
Separating crawling from parsing
Keep parsing logic separate from crawling logic. This lets you:
- Re-parse cached HTML when you change the parser (no need to re-crawl).
- Test parsers against saved HTML samples.
- Replace one without changing the other.
Handling parser failures
When the website changes its HTML structure, parsers break. Build in validation:
def validate_record(record):
"""Check that parsed data looks reasonable."""
issues = []
if not record.get('name'):
issues.append('missing name')
if record.get('email') and '@' not in record['email']:
issues.append('invalid email format')
if record.get('phone') and len(record['phone']) < 7:
issues.append('phone too short')
return issues
If validation failures spike, the parser is likely broken. Alert and pause until fixed.
Stage 4: Processing
Data cleaning
After parsing, apply normalisation:
class Processor:
def process(self, records):
cleaned = []
for record in records:
record['email'] = self._clean_email(record.get('email'))
record['phone'] = self._clean_phone(record.get('phone'))
record['name'] = self._clean_name(record.get('name'))
if record['email']: # Only keep records with valid email
cleaned.append(record)
return self._deduplicate(cleaned)
def _clean_email(self, email):
if not email:
return None
return email.strip().lower()
def _deduplicate(self, records):
seen = set()
unique = []
for record in records:
if record['email'] not in seen:
seen.add(record['email'])
unique.append(record)
return unique
Email extraction from scraped content
When scraping pages that contain email addresses mixed with other text (contact pages, team pages, forum posts), you need to extract the emails.
For scraped HTML files saved locally, upload them to Email Extractor. It extracts email addresses from HTML, TXT, CSV, JSON and 15 other file formats, and deduplicates the results.
For programmatic extraction within the pipeline:
import re
EMAIL_PATTERN = re.compile(
r'[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}'
)
def extract_emails(text):
return list(set(EMAIL_PATTERN.findall(text)))
See Regex for Email Extraction.
Incremental processing
Do not re-process the entire dataset every run. Track what changed.
New records: Items that did not exist in the previous run. Updated records: Items that existed but changed (different phone, new address). Deleted records: Items that existed in the previous run but are gone now.
def compare_runs(previous, current, key='email'):
prev_map = {r[key]: r for r in previous}
curr_map = {r[key]: r for r in current}
new = [curr_map[k] for k in curr_map if k not in prev_map]
deleted = [prev_map[k] for k in prev_map if k not in curr_map]
updated = [
curr_map[k] for k in curr_map
if k in prev_map and curr_map[k] != prev_map[k]
]
return new, updated, deleted
Stage 5: Storage
File storage (simple)
For small pipelines, CSV or JSON files with timestamps:
data/
2026-10-08/
raw/
source-a.html
source-b.html
processed/
contacts.csv
contacts.json
logs/
scrape.log
Database storage (scalable)
For larger pipelines, a database provides querying, deduplication and history:
CREATE TABLE contacts (
id SERIAL PRIMARY KEY,
email VARCHAR(255) UNIQUE NOT NULL,
name VARCHAR(255),
company VARCHAR(255),
phone VARCHAR(50),
source VARCHAR(100),
first_seen TIMESTAMP DEFAULT NOW(),
last_seen TIMESTAMP DEFAULT NOW(),
last_updated TIMESTAMP DEFAULT NOW()
);
Downstream delivery
After processing, deliver data to where it needs to go:
- CRM import. Push to HubSpot, Salesforce or your CRM via API.
- Email platform. Add to a list in Mailchimp, ActiveCampaign or your ESP.
- Data warehouse. Load into BigQuery, Snowflake or Redshift for analysis.
- File delivery. Drop a CSV in a shared folder or S3 bucket.
Stage 6: Monitoring
What to monitor
Job status. Did the scrape run? Did it complete? How long did it take?
Data volume. How many records were scraped? Is the count within expected range? A sudden drop (site changed, got blocked) or spike (pagination error, duplicate pages) indicates a problem.
Data quality. How many records passed validation? What percentage have email addresses? What is the error rate?
Website changes. Did the HTML structure change? Are selectors still matching? A spike in parse failures is the early warning.
Alerting
Set up alerts for:
- Job failure (did not start, crashed, timed out).
- Data volume anomaly (more than 20% deviation from average).
- Validation failure spike (more than 10% of records failing validation).
- IP blocked (429 or 403 responses).
- Scheduled run missed.
Logging
Log every step:
- URLs fetched and response codes.
- Records parsed per page.
- Records after cleaning and deduplication.
- Records delivered to storage.
- Errors and exceptions with full context.
Scaling Considerations
Multiple sources
When scraping 5+ sources:
- Use a task queue (Celery, RQ, Dramatiq) to manage jobs.
- Each source gets its own crawler and parser.
- Shared processor and storage.
- Centralised monitoring dashboard.
Distributed scraping
For very large scraping operations:
- Run multiple crawler instances across different servers/IPs.
- Use a queue (Redis, RabbitMQ) to distribute URLs.
- Central storage receives results from all crawlers.
Legal considerations
Before building a scraping pipeline:
- Check each target website's Terms of Service.
- Review robots.txt.
- Understand privacy laws applicable to the data you collect.
- Consult legal counsel for commercial scraping operations.
See Is Web Scraping Legal? and Web Scraping vs API Access.