Build a Simple ETL Pipeline: Extract, Transform, Load
Data rarely arrives clean and ready to use. Sales come from a CSV, customer info from an API, and inventory from a legacy system. Before anyone can analyze this data, someone has to pull it together, clean it up, and store it in a usable format. That process is called ETL - Extract, Transform, Load.
ETL pipelines are the backbone of data engineering. Every company that uses data has them running behind the scenes, often on a schedule. In this project, you'll build a complete ETL pipeline that extracts data from multiple simulated sources, validates and transforms it, and loads it into a SQLite database.
Step 1: Extract Data from Multiple Sources
The "E" in ETL stands for Extract. In the real world, this means pulling data from databases, APIs, files, and services. Since we can't access real APIs in the browser, we'll simulate multiple data sources using Python data structures - dictionaries for API-like data and multi-line strings for CSV-like data.
The key principle: extraction should be source-agnostic. Your pipeline should convert every source into a common format (a list of dictionaries) so downstream steps don't need to know where the data came from.
Create two extraction functions:
1. extract_from_api(data) - takes a list of dictionaries, prints "Extracted N records from API", returns the list.
2. extract_from_csv(csv_string) - takes a CSV string, parses it into a list of dicts (converting id to int and amount to float), prints "Extracted N records from CSV", returns the list.
Combine results from both sources into all_records. Print the total count as "Total records: N" and print each record's id and name.
# Source 1: API-like data
api_data = [
{'id': 1, 'name': 'Alice Johnson', 'email': 'alice@example.com', 'amount': 150.00},
{'id': 2, 'name': 'Bob Smith', 'email': 'bob@example.com', 'amount': 230.50},
{'id': 3, 'name': 'Charlie Brown', 'email': '', 'amount': 75.00},
]
# Source 2: CSV-like data
csv_data = """id,name,email,amount
4,Diana Prince,diana@example.com,320.00
5,Eve Wilson,,180.25
6,Frank Miller,frank@example.com,95.50"""
def extract_from_api(data):
pass
def extract_from_csv(csv_string):
pass
# Combine and print results
Step 2: Validate the Extracted Data
Data from the real world is messy. Emails might be missing, amounts might be negative, names could be empty strings. Before transforming data, you should validate it and flag any problems. Some issues are warnings (missing email), while others are errors (negative amount).
Write a function validate_records(records) that checks each record for:
1. name must not be empty or missing - if invalid, add to errors list.
2. amount must be a positive number - if invalid, add to errors list.
3. email should not be empty - if empty, add to warnings list.
Return a tuple: (valid_records, errors, warnings) where valid_records excludes records with errors (warnings are OK).
Test with sample data that includes at least one invalid record (empty name or negative amount). Print "Valid: N", "Errors: N", "Warnings: N".
def validate_records(records):
valid = []
errors = []
warnings = []
# Check each record
pass
return valid, errors, warnings
# Test data with some problems
test_data = [
{'id': 1, 'name': 'Alice Johnson', 'email': 'alice@example.com', 'amount': 150.00},
{'id': 2, 'name': '', 'email': 'bad@example.com', 'amount': 230.50},
{'id': 3, 'name': 'Charlie Brown', 'email': '', 'amount': 75.00},
{'id': 4, 'name': 'Diana Prince', 'email': 'diana@example.com', 'amount': -50.00},
{'id': 5, 'name': 'Eve Wilson', 'email': 'eve@example.com', 'amount': 180.25},
]
valid, errs, warns = validate_records(test_data)
Step 3: Transform - Clean and Standardize
The "T" in ETL is often the most complex step. Transformation includes cleaning (fixing inconsistencies), standardizing (making formats uniform), and enriching (adding derived fields). Let's start with cleaning and standardizing.
Common cleaning tasks include trimming whitespace, standardizing capitalization, fixing data types, and handling missing values. The goal is to make every record follow the same format.
Write a function clean_records(records) that transforms each record:
1. Trim whitespace from name and email.
2. Convert name to title case (e.g., "alice johnson" becomes "Alice Johnson").
3. Convert email to lowercase.
4. Round amount to 2 decimal places.
5. If email is empty, set it to "unknown@placeholder.com".
Return the list of cleaned records. Test with messy data and print the cleaned results.
def clean_records(records):
cleaned = []
# Clean each record
pass
return cleaned
# Messy test data
messy_data = [
{'id': 1, 'name': ' alice johnson ', 'email': ' ALICE@Example.COM ', 'amount': 150.456},
{'id': 2, 'name': 'BOB SMITH', 'email': 'Bob@Example.com', 'amount': 230.5},
{'id': 3, 'name': 'charlie brown', 'email': '', 'amount': 75.0},
]
cleaned = clean_records(messy_data)
for r in cleaned:
print(f'{r["id"]}: {r["name"]} | {r["email"]} | ${r["amount"]:.2f}')
Step 4: Transform - Derive New Fields
Beyond cleaning, transformations often add new fields that make the data more useful for analysis. These derived fields are calculated from existing data. For example, you might categorize amounts into tiers, extract the domain from email addresses, or add processing timestamps.
Write a function enrich_records(records) that adds these derived fields to each record:
1. amount_tier - "low" if amount < 100, "medium" if 100–250, "high" if > 250.
2. email_domain - the part after @ in the email (e.g., "example.com").
3. name_parts - number of words in the name.
4. processed - set to True.
Return the enriched list. Test with sample data and print each record showing all new fields.
def enrich_records(records):
enriched = []
# Add derived fields
pass
return enriched
sample = [
{'id': 1, 'name': 'Alice Johnson', 'email': 'alice@example.com', 'amount': 150.00},
{'id': 2, 'name': 'Bob Smith', 'email': 'bob@company.org', 'amount': 50.00},
{'id': 3, 'name': 'Charlie Dean Brown', 'email': 'charlie@startup.io', 'amount': 320.00},
]
enriched = enrich_records(sample)
for r in enriched:
print(f'{r["name"]}: ${r["amount"]:.2f} -> tier={r["amount_tier"]}, domain={r["email_domain"]}, parts={r["name_parts"]}')
Step 5: Load into SQLite
The "L" in ETL means loading transformed data into its final destination. For us, that's a SQLite database. We'll create a table that matches our enriched record structure and insert all records in a batch.
Write a function load_to_database(records) that:
1. Creates an in-memory SQLite database.
2. Creates a customers table with columns: id (INTEGER PRIMARY KEY), name (TEXT), email (TEXT), amount (REAL), amount_tier (TEXT), email_domain (TEXT).
3. Inserts all records using executemany().
4. Commits and returns the connection.
Then query the database to verify: print the row count and all rows. Print "Loaded N records into database".
import sqlite3
def load_to_database(records):
# Create database, table, insert records
pass
# Sample enriched records
records = [
{'id': 1, 'name': 'Alice Johnson', 'email': 'alice@example.com', 'amount': 150.00, 'amount_tier': 'medium', 'email_domain': 'example.com'},
{'id': 2, 'name': 'Bob Smith', 'email': 'bob@company.org', 'amount': 230.50, 'amount_tier': 'medium', 'email_domain': 'company.org'},
{'id': 3, 'name': 'Charlie Brown', 'email': 'charlie@startup.io', 'amount': 75.00, 'amount_tier': 'low', 'email_domain': 'startup.io'},
{'id': 4, 'name': 'Diana Prince', 'email': 'diana@example.com', 'amount': 320.00, 'amount_tier': 'high', 'email_domain': 'example.com'},
]
conn = load_to_database(records)
Step 6: Query and Verify the Loaded Data
After loading data, you should always verify it. Run queries to check row counts, look for unexpected NULLs, and generate summary statistics. This is the quality assurance step that catches loading errors before anyone uses the data.
Study the code below. It loads data into SQLite and runs verification queries. What will be printed? Pay close attention to the SQL aggregations and GROUP BY results.
import sqlite3
conn = sqlite3.connect(':memory:')
cursor = conn.cursor()
cursor.execute('CREATE TABLE customers (id INTEGER PRIMARY KEY, name TEXT, email TEXT, amount REAL, amount_tier TEXT, email_domain TEXT)')
data = [
(1, 'Alice', 'alice@a.com', 150.0, 'medium', 'a.com'),
(2, 'Bob', 'bob@b.com', 50.0, 'low', 'b.com'),
(3, 'Charlie', 'charlie@a.com', 320.0, 'high', 'a.com'),
(4, 'Diana', 'diana@b.com', 80.0, 'low', 'b.com'),
]
cursor.executemany('INSERT INTO customers VALUES (?, ?, ?, ?, ?, ?)', data)
conn.commit()
cursor.execute('SELECT COUNT(*) FROM customers')
print(f'Total: {cursor.fetchone()[0]}')
cursor.execute('SELECT amount_tier, COUNT(*), SUM(amount) FROM customers GROUP BY amount_tier ORDER BY amount_tier')
for tier, count, total in cursor.fetchall():
print(f'{tier}: {count} customers, ${total:.0f}')Step 7: Build the Complete ETLPipeline Class
Now let's assemble everything into a single ETLPipeline class. A class-based pipeline is easier to configure, test, and rerun. It tracks statistics from each step so you can monitor the pipeline's health.
Build an ETLPipeline class with these methods:
1. __init__(self) - initializes stats dict: {extracted: 0, validated: 0, errors: 0, transformed: 0, loaded: 0}.
2. extract(self, api_data, csv_data) - combines both sources, updates stats, returns records.
3. validate(self, records) - filters out records with empty name or non-positive amount, updates stats, returns valid records.
4. transform(self, records) - cleans (strip, title case names, lowercase emails) and adds amount_tier, updates stats, returns transformed records.
5. load(self, records) - creates SQLite :memory: db, loads records, stores connection as self.conn, updates stats.
6. run(self, api_data, csv_data) - runs all steps in order, prints a summary report.
Test by running the pipeline with sample data. The summary should print "=== PIPELINE SUMMARY ===" and each stat.
import sqlite3
class ETLPipeline:
def __init__(self):
self.stats = {'extracted': 0, 'validated': 0, 'errors': 0, 'transformed': 0, 'loaded': 0}
self.conn = None
def extract(self, api_data, csv_data):
pass
def validate(self, records):
pass
def transform(self, records):
pass
def load(self, records):
pass
def run(self, api_data, csv_data):
pass
# Test data
api_data = [
{'id': 1, 'name': 'Alice Johnson', 'email': 'ALICE@example.com', 'amount': 150.00},
{'id': 2, 'name': '', 'email': 'bad@test.com', 'amount': 100.00},
{'id': 3, 'name': 'Charlie Brown', 'email': '', 'amount': 75.00},
]
csv_data = """id,name,email,amount
4,Diana Prince,diana@example.com,320.00
5,Eve Wilson,eve@test.com,-50.00
6,Frank Miller,frank@example.com,95.50"""
pipeline = ETLPipeline()
pipeline.run(api_data, csv_data)
Project Complete!
You've built a complete ETL pipeline from scratch. Here's what you accomplished: