Most data pipelines do not fail because the transformation logic was wrong. They fail because a source returned a surprise, a job ran twice, or a run died halfway and left the table in a strange state. Building pipelines that survive those days comes down to a few design principles. This article walks through them with Python, PostgreSQL and Apache Airflow.
Separate extract, transform and load
Three small steps are easier to test, retry and reason about than one long script:
- Extract pulls raw data from the source and stores it untouched.
- Transform cleans, validates and reshapes it.
- Load writes the result to its destination.
Keeping the raw data lets you re-run the transform after fixing a bug without hitting the source again, which matters when the source is rate-limited, slow, or only keeps recent history.
Make every run idempotent
An idempotent job produces the same result whether it runs once or five times. That one property removes most operational anxiety, because "just re-run it" becomes a safe answer.
The most reliable pattern is to process data in partitions, usually by date, and to make each run replace its own partition:
from datetime import date
import psycopg
def load_daily_usage(conn: psycopg.Connection, day: date, rows: list[tuple]) -> None:
with conn.transaction():
conn.execute("DELETE FROM daily_usage WHERE usage_date = %s", (day,))
with conn.cursor() as cur:
cur.executemany(
"INSERT INTO daily_usage (usage_date, site_id, bytes_in, bytes_out) "
"VALUES (%s, %s, %s, %s)",
rows,
)
The delete and the insert happen in one transaction, so readers see either the old day or the new day, never half of each. If the job crashes, the transaction rolls back and nothing is left in between.
When you cannot replace a whole partition, use an upsert with a natural key instead:
INSERT INTO customers (customer_id, name, plan, updated_at)
VALUES (%s, %s, %s, now())
ON CONFLICT (customer_id)
DO UPDATE SET name = EXCLUDED.name,
plan = EXCLUDED.plan,
updated_at = now();
Running that twice with the same input leaves the table identical.
Validate at the boundary
Bad data is cheapest to catch the moment it enters your pipeline. Define what a valid row looks like and check every row against it. Pydantic works well for this:
from datetime import date
from pydantic import BaseModel, Field, ValidationError
class UsageRow(BaseModel):
usage_date: date
site_id: int = Field(gt=0)
bytes_in: int = Field(ge=0)
bytes_out: int = Field(ge=0)
def validate(records: list[dict]) -> tuple[list[UsageRow], list[dict]]:
good, bad = [], []
for record in records:
try:
good.append(UsageRow(**record))
except ValidationError as err:
bad.append({"record": record, "error": str(err)})
return good, bad
Do not silently drop the bad rows. Write them to a quarantine table with the reason, and raise an alert when the rejection rate crosses a threshold. A sudden jump from 0.1% to 20% rejected rows usually means the source changed its format.
Add table-level checks too: the row count is within an expected range, key columns have no nulls, and totals match the source. These cheap assertions catch the failures that row validation cannot, like a source that quietly returned half the data.
Retries, timeouts and backoff
Networks fail. Retry transient errors, but not forever, and not instantly. The tenacity library keeps that logic readable:
import httpx
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type
@retry(
stop=stop_after_attempt(5),
wait=wait_exponential(multiplier=2, max=60),
retry=retry_if_exception_type((httpx.TransportError, httpx.HTTPStatusError)),
reraise=True,
)
def fetch_page(client: httpx.Client, url: str, params: dict) -> dict:
response = client.get(url, params=params, timeout=30)
response.raise_for_status()
return response.json()
Always set a timeout, because a request with no timeout can hang a worker forever. Retrying is only safe because the load step is idempotent, which is one more reason to get that right first.
Be careful with Pandas at scale
Pandas is excellent for transformations that fit in memory and a poor fit when they do not. A few habits help:
- Read large files in chunks with
chunksizeinstead of loading everything. - Specify
dtypeexplicitly so a column of IDs is not inferred as floats. - Prefer vectorised operations to row-by-row
applycalls. - Push heavy aggregation into PostgreSQL, where it can use indexes, rather than pulling rows into Python to sum them.
If a single run starts to strain memory, that is often the signal to move that step into SQL or to a distributed engine, not to buy a bigger machine.
Orchestrate with Airflow
Airflow schedules the steps, retries failures and shows you what ran. A small DAG using the TaskFlow API looks like this:
from datetime import datetime, timedelta
from airflow.decorators import dag, task
@dag(
schedule="@daily",
start_date=datetime(2025, 1, 1),
catchup=False,
default_args={"retries": 3, "retry_delay": timedelta(minutes=5)},
tags=["usage"],
)
def daily_usage():
@task
def extract(ds=None) -> str:
# fetch the source data for `ds` (YYYY-MM-DD) and write it to storage
return f"s3://example-raw/usage/{ds}.json"
@task
def transform(raw_path: str) -> str:
# read raw data, validate rows, write clean data and a quarantine file
return raw_path.replace("raw", "clean")
@task
def load(clean_path: str, ds=None) -> None:
# replace the partition for `ds` in one transaction
...
load(transform(extract()))
daily_usage()
Two Airflow ideas matter more than any syntax. First, tasks should receive the logical date (ds) and compute everything from it, never from datetime.now(), so a re-run for last Tuesday processes last Tuesday. Second, tasks should pass references (a path or table name) between each other and not the data itself, because Airflow's task-to-task channel is meant for small values.
Observability
A pipeline you cannot see into is one you will not trust. At minimum, log row counts at each stage (extracted, valid, rejected, loaded), track how long each run takes, and alert on failures and on freshness, meaning "the table has not updated in 26 hours". Freshness alerts catch the silent failures that a green dashboard hides.
The short version
- Split extract, transform and load, and keep the raw data.
- Make loads idempotent with partition replacement or upserts.
- Validate every row, quarantine rejects and alert on spikes.
- Retry with backoff, always with timeouts.
- Use the logical date, and pass references between tasks.
- Monitor counts, durations and freshness.
Get these right and a pipeline stops being something you watch nervously and becomes something you can safely forget about.