Building a Data Pipeline for Email List Processing
On this page
Why Build a Pipeline
When email list processing is manual, the same steps repeat for every batch: extract addresses from source files, remove duplicates, clean formatting, validate deliverability, filter out junk and load the results into a CRM or email tool. Each step takes time, introduces human error and creates bottlenecks when volume increases.
A data pipeline automates this sequence. Each stage takes input from the previous stage, applies a transformation and passes the result forward. Once built, the pipeline processes 100 addresses or 100,000 addresses with the same effort.
Pipeline Architecture
A typical email list processing pipeline has six stages:
Ingestion -> Extraction -> Deduplication -> Validation -> Enrichment -> Storage
| Stage | Input | Output | Purpose |
|---|---|---|---|
| Ingestion | Raw files (CSV, XLSX, PDF, text) | Staged files in a standard location | Collect source files from all channels |
| Extraction | Staged files | Raw email list (one address per row) | Pull email addresses from various formats |
| Deduplication | Raw email list | Unique email list | Remove duplicates (case-insensitive) |
| Validation | Unique email list | Validated list with status flags | Check syntax, domain, deliverability |
| Enrichment | Validated list | Enriched list with metadata | Add name, company, source, date |
| Storage | Enriched list | Database or CRM records | Store for use in campaigns |
Each stage is independent. If validation fails for a batch, the earlier stages do not need to re-run. If you add a new source file, it enters at ingestion and flows through the existing stages.
Stage 1: Ingestion
Ingestion collects source files and stages them for processing.
Sources
Email addresses arrive in many formats from many places:
| Source | Typical format | Challenge |
|---|---|---|
| CRM exports | CSV or XLSX | May contain thousands of columns; email column name varies |
| Web scraping results | CSV, JSON or text files | Mixed quality; may contain false positives |
| Event attendee lists | XLSX or CSV | Often contains names, roles and other fields alongside email |
| Business card scans | VCF (vCard) | Multiple fields; email may be in any of several vCard properties |
| Document archives | PDF, DOCX | Email addresses embedded in unstructured text |
| Email exports | EML, MSG | Addresses in headers and body |
| Manual collection | Text files, spreadsheets | Inconsistent formatting |
File staging
Create a standard directory structure:
/pipeline/
/inbox/ # New files land here
/processing/ # Files currently being processed
/completed/ # Processed files (kept for audit trail)
/failed/ # Files that could not be processed
/output/ # Final results
A simple ingestion script watches the inbox directory and moves files into processing:
import os
import shutil
from pathlib import Path
from datetime import datetime
INBOX = Path("/pipeline/inbox")
PROCESSING = Path("/pipeline/processing")
def ingest():
"""Move files from inbox to processing with timestamp."""
for filepath in INBOX.iterdir():
if filepath.is_file():
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
dest = PROCESSING / f"{timestamp}_{filepath.name}"
shutil.move(str(filepath), str(dest))
print(f"Ingested: {filepath.name}")
Handling large files
For large files (over 25 MB), split them before processing:
import pandas as pd
def split_csv(filepath, chunk_size=50000):
"""Split a large CSV into chunks."""
for i, chunk in enumerate(
pd.read_csv(filepath, chunksize=chunk_size)
):
chunk_path = filepath.with_name(
f"{filepath.stem}_part{i}{filepath.suffix}"
)
chunk.to_csv(chunk_path, index=False)
Stage 2: Extraction
Extraction pulls email addresses from source files regardless of format.
Format-specific extraction
Different file formats require different extraction approaches:
| Format | Extraction method |
|---|---|
| CSV/TSV | Read as structured data; identify the email column by header name or content pattern |
| XLSX/XLS | Read with a spreadsheet library; same column identification as CSV |
| JSON | Traverse the structure; extract string values matching email patterns |
| Extract text content; apply regex to find email patterns (note: scanned PDFs without OCR text layers will not yield results) | |
| DOCX | Extract text from paragraphs; apply regex |
| VCF | Parse vCard format; extract EMAIL property values |
| EML/MSG | Parse email headers (From, To, CC, Reply-To) and body |
| Plain text | Apply regex directly |
| HTML | Extract from visible text and mailto: links |
Using Email Extractor for batch extraction
For one-off or small-batch processing, Email Extractor handles extraction from 19 file formats directly in the browser:
- Upload files (up to 25 MB per file, 100 MB per batch).
- The tool extracts email addresses from all uploaded files.
- Results are deduplicated (case-insensitive).
- Download as CSV with source information.
This is useful when you want to process a batch without writing code, or when you need to verify your pipeline's extraction results against a known-good implementation.
Regex-based extraction
For automated pipelines, regex extraction handles most plain-text formats:
import re
EMAIL_PATTERN = re.compile(
r"[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}"
)
def extract_emails_from_text(text):
"""Extract email addresses from text content."""
return EMAIL_PATTERN.findall(text)
Column identification for structured data
When processing CSV or XLSX files, identify the email column automatically:
import pandas as pd
EMAIL_COLUMN_NAMES = {
"email", "e-mail", "email_address", "emailaddress",
"email address", "mail", "contact_email", "work_email",
"personal_email", "user_email",
}
def find_email_column(df):
"""Find the column most likely to contain email addresses."""
# Check column names first
for col in df.columns:
if col.strip().lower() in EMAIL_COLUMN_NAMES:
return col
# Fall back to content analysis
email_pattern = re.compile(
r"^[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}$"
)
for col in df.columns:
sample = df[col].dropna().head(20).astype(str)
matches = sample.apply(
lambda x: bool(email_pattern.match(x.strip()))
)
if matches.mean() > 0.5:
return col
return None
Stage 3: Deduplication
Deduplication removes duplicate addresses so each person appears only once in the output.
Case-insensitive deduplication
Email addresses are case-insensitive by RFC specification. User@Example.com and user@example.com deliver to the same mailbox.
def deduplicate(emails):
"""Remove duplicate email addresses (case-insensitive)."""
seen = set()
unique = []
for email in emails:
normalised = email.strip().lower()
if normalised not in seen:
seen.add(normalised)
unique.append(normalised)
return unique
Cross-source deduplication
When processing files from multiple sources, track which source each address came from. Keep the earliest source as the primary record:
from collections import OrderedDict
def deduplicate_with_sources(email_source_pairs):
"""Deduplicate across sources, keeping first occurrence."""
records = OrderedDict()
for email, source in email_source_pairs:
normalised = email.strip().lower()
if normalised not in records:
records[normalised] = {
"email": normalised,
"primary_source": source,
"all_sources": [source],
}
else:
records[normalised]["all_sources"].append(source)
return list(records.values())
Plus-address handling
Decide whether to treat plus-addressed emails as duplicates:
user+newsletter@example.comanduser@example.comdeliver to the same inbox.- Some organisations deliberately use plus addressing for routing.
def normalise_plus_address(email):
"""Remove plus addressing for deduplication."""
local, domain = email.split("@")
local = local.split("+")[0]
return f"{local}@{domain}"
Whether to normalise plus addresses depends on your use case. For cold outreach, normalising avoids sending multiple emails to the same inbox. For newsletter lists where users deliberately signed up with different plus addresses, keep them separate.
Stage 4: Validation
Validation checks whether each address is real and deliverable.
Validation layers
| Layer | What it checks | Speed | Accuracy |
|---|---|---|---|
| Syntax | Format matches email specification | Instant | Catches typos, not fake addresses |
| Domain | Domain exists and has MX records | Fast (DNS lookup) | Catches non-existent domains |
| Mailbox | SMTP check to verify the mailbox exists | Slow (network call) | Catches non-existent users |
| Deliverability | Full verification including catch-all detection | Slowest (API call) | Most accurate |
Syntax validation
import re
def validate_syntax(email):
"""Check if email has valid syntax."""
pattern = r"^[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}$"
return bool(re.match(pattern, email))
Domain validation
import dns.resolver
def validate_domain(email):
"""Check if the email domain has MX records."""
domain = email.split("@")[1]
try:
dns.resolver.resolve(domain, "MX")
return True
except (dns.resolver.NXDOMAIN, dns.resolver.NoAnswer):
return False
except Exception:
return None # Inconclusive
API-based validation
For production pipelines, use an email verification API for the mailbox check:
def validate_batch(emails, api_key):
"""Validate a batch of emails through a verification API."""
results = []
for email in emails:
# Syntax and domain checks first (free)
if not validate_syntax(email):
results.append({"email": email, "status": "invalid_syntax"})
continue
if not validate_domain(email):
results.append({"email": email, "status": "invalid_domain"})
continue
# API check for mailbox verification (costs per check)
status = api_verify(email, api_key)
results.append({"email": email, "status": status})
return results
Running syntax and domain checks before the API call reduces the number of paid API calls.
Stage 5: Enrichment
Enrichment adds metadata to each email address to make it more useful for outreach.
Source-based enrichment
Information already available from the extraction stage:
| Field | Source |
|---|---|
| Source file | Tracked during ingestion |
| Source URL | If scraped, the page URL where the address was found |
| Extraction date | Timestamp from the ingestion stage |
| Co-located name | If the source file had name and email in the same row |
| Company domain | Derived from the email domain for business addresses |
Third-party enrichment
External data providers can add information based on the email address:
| Data point | Common providers |
|---|---|
| Full name | Clearbit, Hunter.io, Apollo.io |
| Job title | Clearbit, Apollo.io, LinkedIn (via API) |
| Company name | Clearbit, derived from domain |
| Company size | Clearbit, Crunchbase |
| Industry | Clearbit, Apollo.io |
| Location | Clearbit, IP-based (if web scraping) |
| Social profiles | Clearbit, FullContact |
Enrichment is the most expensive stage. Prioritise enrichment for validated addresses only. There is no value in enriching an address that will bounce.
Stage 6: Storage
The final stage loads processed records into your working system.
Storage options
| System | Best for | Integration method |
|---|---|---|
| CRM (HubSpot, Salesforce, Pipedrive) | Sales teams running outbound campaigns | API or CSV import |
| Email marketing platform (Mailchimp, Brevo) | Newsletter and marketing campaigns | API or CSV import |
| Cold email tool (Instantly, Smartlead, Lemlist) | Cold outreach sequences | CSV import or API |
| Database (PostgreSQL, MySQL) | Custom applications and reporting | Direct insert |
| Spreadsheet (Google Sheets, Excel) | Small teams, manual review | CSV export |
Idempotent loading
Design the storage stage to be idempotent: running it twice with the same data should not create duplicate records.
def load_to_crm(records, crm_client):
"""Load records into CRM, skipping duplicates."""
for record in records:
existing = crm_client.find_by_email(record["email"])
if existing:
# Update existing record with new source info
crm_client.update(existing["id"], {
"sources": existing["sources"] + [record["source"]],
"last_seen": record["extraction_date"],
})
else:
crm_client.create(record)
Scheduling and Monitoring
Running the pipeline
For recurring processing (new files arrive daily or weekly):
| Approach | Complexity | Best for |
|---|---|---|
| Cron job | Low | Simple, time-based scheduling |
| File watcher | Low | Event-driven processing (process files as they arrive) |
| Task queue (Celery, RQ) | Medium | Parallel processing with retry logic |
| Workflow orchestrator (Airflow, Prefect) | High | Complex pipelines with dependencies and monitoring |
Monitoring
Track these metrics for each pipeline run:
| Metric | What it tells you |
|---|---|
| Files processed | Pipeline is running and picking up new files |
| Addresses extracted | Extraction stage is working |
| Duplicate rate | Percentage of duplicates across sources |
| Validation pass rate | Percentage of addresses that are deliverable |
| Enrichment coverage | Percentage of addresses with complete metadata |
| Processing time | Whether the pipeline is keeping up with input volume |
| Error count | Files or addresses that failed processing |
Error handling
Each stage should handle errors without stopping the entire pipeline:
def process_file(filepath):
"""Process a single file through the pipeline."""
try:
emails = extract(filepath)
unique = deduplicate(emails)
validated = validate(unique)
enriched = enrich(validated)
load(enriched)
move_to_completed(filepath)
except ExtractionError:
move_to_failed(filepath, "extraction_failed")
except ValidationError:
# Partial results are still useful
load(unique) # Load unvalidated
move_to_completed(filepath, "validation_skipped")
except Exception as e:
move_to_failed(filepath, str(e))
Scaling Considerations
| Volume | Architecture | Tools |
|---|---|---|
| Under 1,000 addresses/day | Single script, cron job | Python script, CSV files |
| 1,000-10,000 addresses/day | Task queue with workers | Celery or RQ, PostgreSQL |
| 10,000-100,000 addresses/day | Distributed pipeline | Airflow or Prefect, cloud workers |
| Over 100,000 addresses/day | Streaming pipeline | Apache Kafka or AWS Kinesis, distributed processing |
For most B2B prospecting use cases, a simple Python script with a cron job handles the volume. Over-engineering the pipeline is a common mistake.