Python courseLesson 5 of 5
Python course · Tutorial · Lesson 5 of 5
Build an Idempotent CSV Loader in Python
Build a small Python loader that reads CSV files, skips and counts bad rows, and loads into a database so running it twice gives the same result.
On this page
What you will learn
- Load CSV data into a database inside one transaction
- Make a load safe to run twice with an upsert
- Count and log bad rows instead of failing silently or crashing
- Test the loader with pytest
A pipeline step is idempotent when running it twice for the same input leaves the system in the same state as running it once. That single property makes retries, reruns and backfills safe. In this tutorial you build a small loader that has it.
Goal
Load order CSV files into a database table so that:
- running the loader twice does not create duplicate rows;
- a corrected file replaces earlier values for the same order;
- malformed rows are skipped, counted and logged, never silently lost;
- a failure part-way through leaves no half-loaded file.
Prerequisites
Python 3.10 or later and pytest (pip install pytest). We use SQLite so there is nothing to install, but the same pattern works for Postgres, Snowflake or any database that supports upserts.
The loader
Save this as loader.py.
import csv
import logging
import sqlite3
from pathlib import Path
logging.basicConfig(level=logging.INFO, format="%(levelname)s %(message)s")
log = logging.getLogger("loader")
SCHEMA = """
CREATE TABLE IF NOT EXISTS orders (
order_id INTEGER PRIMARY KEY,
customer TEXT NOT NULL,
amount REAL NOT NULL
);
CREATE TABLE IF NOT EXISTS load_audit (
file_name TEXT PRIMARY KEY,
rows_read INTEGER NOT NULL,
rows_bad INTEGER NOT NULL
);
"""
def parse_rows(path):
"""Return (good_rows, bad_count). Bad rows are logged and skipped, never silently dropped."""
bad = 0
good = []
with path.open(newline="") as f:
for line_no, row in enumerate(csv.DictReader(f), start=2):
try:
good.append((int(row["order_id"]), row["customer"].strip(), float(row["amount"])))
except (KeyError, ValueError, AttributeError):
bad += 1
log.warning("skipping bad row at line %d: %r", line_no, row)
return good, bad
def load_file(conn, path):
"""Load one file. Safe to run twice: the file is replaced, not appended."""
rows, bad = parse_rows(path)
with conn: # one transaction: all of it commits, or none of it does
conn.executemany(
"""INSERT INTO orders (order_id, customer, amount) VALUES (?, ?, ?)
ON CONFLICT(order_id) DO UPDATE SET customer = excluded.customer, amount = excluded.amount""",
rows,
)
conn.execute(
"""INSERT INTO load_audit (file_name, rows_read, rows_bad) VALUES (?, ?, ?)
ON CONFLICT(file_name) DO UPDATE SET rows_read = excluded.rows_read, rows_bad = excluded.rows_bad""",
(path.name, len(rows) + bad, bad),
)
log.info("%s: loaded %d rows, skipped %d", path.name, len(rows), bad)
def main(data_dir="data", db_path="warehouse.db"):
conn = sqlite3.connect(db_path)
conn.executescript(SCHEMA)
for path in sorted(Path(data_dir).glob("*.csv")):
load_file(conn, path)
conn.close()
if __name__ == "__main__":
main()
How it achieves each goal
- No duplicates on rerun.
order_idis the primary key and the insert usesON CONFLICT ... DO UPDATE(an upsert). Loading the same row again updates it to identical values. - Corrections win. Because the upsert overwrites, a corrected file changes the stored values rather than being ignored.
- Bad rows are visible.
parse_rowscatches only the expected parsing errors, logs each bad row with its line number, andload_auditrecords how many rows were read and how many were bad. Compare these in monitoring. - All or nothing.
with conn:wraps the inserts and the audit write in one transaction. If anything raises, the database rolls back and the file is not half loaded.
Try it
mkdir data
printf 'order_id,customer,amount\n1,Asha,50\n2,Ben,20\nx,Bad,1\n' > data/orders_1.csv
python loader.py
python loader.py # run it a second time
Both runs log the same result:
WARNING skipping bad row at line 4: {'order_id': 'x', 'customer': 'Bad', 'amount': '1'}
INFO orders_1.csv: loaded 2 rows, skipped 1
Inspect the table afterwards: it holds exactly two orders after one run and after two runs, and load_audit shows 3 rows read and 1 bad.
Test it
Save this as test_loader.py and run pytest.
import sqlite3
from pathlib import Path
import loader
def make_conn():
conn = sqlite3.connect(":memory:")
conn.executescript(loader.SCHEMA)
return conn
def write(tmp, name, text):
p = Path(tmp) / name
p.write_text(text)
return p
def test_running_twice_gives_same_result(tmp_path):
p = write(tmp_path, "orders_1.csv", "order_id,customer,amount\n1,Asha,50\n2,Ben,20\n")
conn = make_conn()
loader.load_file(conn, p)
loader.load_file(conn, p)
assert conn.execute("SELECT COUNT(*) FROM orders").fetchone()[0] == 2
assert conn.execute("SELECT SUM(amount) FROM orders").fetchone()[0] == 70
def test_bad_rows_are_counted_not_loaded(tmp_path):
p = write(tmp_path, "orders_2.csv", "order_id,customer,amount\n1,Asha,50\nx,Ben,20\n3,Chen,abc\n")
conn = make_conn()
loader.load_file(conn, p)
assert conn.execute("SELECT COUNT(*) FROM orders").fetchone()[0] == 1
assert conn.execute("SELECT rows_read, rows_bad FROM load_audit").fetchone() == (3, 2)
def test_corrected_file_updates_existing_rows(tmp_path):
conn = make_conn()
loader.load_file(conn, write(tmp_path, "a.csv", "order_id,customer,amount\n1,Asha,50\n"))
loader.load_file(conn, write(tmp_path, "a.csv", "order_id,customer,amount\n1,Asha,55\n"))
assert conn.execute("SELECT amount FROM orders WHERE order_id = 1").fetchone()[0] == 55
The first test is the important one: it loads the same file twice and asserts the row count and total are unchanged. That is the definition of idempotency expressed as a test.
Validation
You have succeeded when:
pytestreports 3 passed;- running
python loader.pyrepeatedly never changes the row count; - editing a value in the CSV and rerunning updates the stored value.
Troubleshooting
near "ON": syntax error: your SQLite is older than 3.24. Checksqlite3.sqlite_versionor upgrade.- Everything is skipped: the CSV header must contain
order_id,customerandamountexactly. - Rows missing after a crash: that is the transaction working. The whole file rolled back; fix the cause and rerun.
Interview relevance
“How would you design an idempotent load?” is a common question. This loader demonstrates the three standard techniques: a natural key, an upsert, and a single transaction. See also the Airflow retries guide for why this matters.
Next step
Wire this loader into a scheduler, then extend it with data-quality checks. The CSV to warehouse pipeline project builds on this foundation.
Progress is saved in this browser only. No account needed.