Menu

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.

  • Intermediate
  • 6 min read
  • Updated Oct 2026
On this page
  1. Goal
  2. Prerequisites
  3. The loader
  4. How it achieves each goal
  5. Try it
  6. Test it
  7. Validation
  8. Troubleshooting
  9. Interview relevance
  10. Next step

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:

  1. running the loader twice does not create duplicate rows;
  2. a corrected file replaces earlier values for the same order;
  3. malformed rows are skipped, counted and logged, never silently lost;
  4. 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_id is the primary key and the insert uses ON 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_rows catches only the expected parsing errors, logs each bad row with its line number, and load_audit records 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:

  • pytest reports 3 passed;
  • running python loader.py repeatedly 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. Check sqlite3.sqlite_version or upgrade.
  • Everything is skipped: the CSV header must contain order_id, customer and amount exactly.
  • 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.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Tested with Python 3.12 and SQLite 3.45 (UPSERT requires SQLite 3.24 or later)

Progress is saved in this browser only. No account needed.

Search
Filter by type