Production pipelines run automatically, fail sometimes, get retried, and occasionally need backfills. The pipeline should handle all of this without intervention. Three concepts:
Concept 1 — Idempotency
A task is idempotent if running it twice (or N times) produces the same result as running it once.
Bad (not idempotent):
@task
def append_daily_summary():
today = datetime.utcnow().date()
count = db.query("SELECT COUNT(*) FROM orders WHERE date = ?", today)
db.execute("INSERT INTO summary (date, count) VALUES (?, ?)", today, count)
Run this twice and you have two rows. Run it 5 times and you have 5.
Good (idempotent):
@task(execution_date_in_task=True)
def upsert_daily_summary(execution_date):
count = db.query("SELECT COUNT(*) FROM orders WHERE date = ?", execution_date)
db.execute("DELETE FROM summary WHERE date = ?", execution_date)
db.execute("INSERT INTO summary (date, count) VALUES (?, ?)", execution_date, count)
DELETE-then-INSERT keyed by date. Re-running replaces previous output.
The idempotency patterns
Pattern 1: DELETE + INSERT
For tables partitioned by date:
DELETE FROM target WHERE date = '{{ ds }}';
INSERT INTO target SELECT * FROM source WHERE date = '{{ ds }}';
Most common; simple; works for most batch workloads.
Pattern 2: MERGE / UPSERT
For tables with primary keys:
MERGE INTO target USING source
ON target.id = source.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...) VALUES (...);
Handles cases where the same row might be updated.
Pattern 3: Overwrite the whole partition
For lakehouse tables (Iceberg/Delta):
df.write_to("target_table") .overwrite_partitions(year=2026, month=5, day=25)
Overwrites just the partition for this run.
Pattern 4: Idempotent message processing
For event consumers:
def process_event(event):
if already_processed(event.id):
return # idempotent skip
process(event)
mark_processed(event.id)
Track processed IDs (dedupe table); skip already-seen.
Concept 2 — Backfills
Backfill = running the pipeline for historical dates after the fact.
Why you need them:
- New table needs historical data.
- Bug fix; reprocess affected days.
- Source data was delayed; catch up.
How they work in Airflow:
airflow dags backfill daily_pipeline --start-date 2026-04-01 --end-date 2026-04-30
Each day runs as if it were that day, in parallel or sequential.
Idempotency is what makes backfill safe. Without it, backfill creates duplicates.
Backfill patterns
For a 90-day backfill on a daily pipeline:
- Sequential: one day at a time. Slower; easier to monitor.
- Parallel: 30 days concurrently. Faster; can overload upstream.
- Chunked: 10 days per run. Middle ground.
In Airflow:
airflow dags backfill ... --reset-dagruns --max-active-runs 10
max-active-runs caps concurrency.
Backfill anti-patterns
- Backfilling without testing on a single day first. Run 1 day, verify output, then backfill the range.
- Backfilling on top of partial data. If today's run partially succeeded, backfilling creates inconsistent state.
- No way to track which days have been backfilled. Re-runs by mistake.
Defensive practice: log every successful run to a metadata table; check before triggering backfill.
Concept 3 — Scheduling
The schedule defines when DAGs fire. Cron expression or interval.
schedule="0 3 * * *" # 3am every day
schedule="@hourly"
schedule="@daily"
schedule="@weekly"
schedule=timedelta(hours=2)
Important: logical_date vs actual_run_time
Airflow's logical_date (formerly execution_date) is NOT when the DAG runs — it's the START of the data interval.
A DAG with schedule="@daily" fires at midnight; logical_date is yesterday's midnight. Confusing.
Convention: use logical_date for the data range you process (i.e., "process data for the day ending at logical_date + 1 day").
@task
def process_data(logical_date):
# data for logical_date through logical_date + 1 day
process(start=logical_date, end=logical_date + timedelta(days=1))
Catchup
If you create a DAG with start_date=2024-01-01 and catchup=True, Airflow runs every missed day on first deploy.
For 90% of new DAGs, set catchup=False. You don't want history reprocessed automatically.
@dag(start_date=datetime(2026, 1, 1), catchup=False, ...)
Putting it all together — a robust pipeline
from airflow.decorators import dag, task
from datetime import datetime, timedelta
@dag(
start_date=datetime(2026, 1, 1),
schedule="0 3 * * *",
catchup=False,
max_active_runs=1,
default_args={"retries": 3, "retry_delay": timedelta(minutes=5)},
)
def customer_pipeline():
@task
def extract(logical_date):
return fetch_data(date=logical_date)
@task
def load(data, logical_date):
# Idempotent: delete-then-insert
db.execute("DELETE FROM raw.customers WHERE _date = ?", logical_date)
db.execute(
"INSERT INTO raw.customers SELECT *, ? AS _date FROM ...",
logical_date,
)
@task
def run_dbt(logical_date):
# dbt is idempotent by default
run_command(f"dbt run --vars '{{date: {logical_date}}}'")
@task
def validate(logical_date):
count = db.query("SELECT COUNT(*) FROM marts.dim_customers WHERE _date = ?", logical_date)
if count == 0:
raise ValueError(f"No data for {logical_date}")
data = extract()
load(data)
run_dbt()
validate()
dag = customer_pipeline()
Has:
- Idempotent load (DELETE + INSERT by date).
- Retries on failure.
- Validation step (catches silent failures).
- Single active run (no race conditions).
- No catchup (no historical surprises).
Backfill-safe; retry-safe; self-healing.
Common mistakes
- Non-idempotent loads. Most production bugs.
- Catchup=True by default. Surprise historical runs.
- No validation step. Pipelines succeed silently with empty output.
- Logical_date misunderstanding. Process the wrong day's data.
- Backfill without single-day test first. Spread the bug across 90 days.
Takeaway
Idempotency = re-running produces same result. Patterns: DELETE+INSERT, MERGE, partition overwrite, dedup tables. Backfill = re-run history; idempotency makes it safe. Schedule = when. logical_date = data interval start. catchup=False unless you want history. Together: pipelines that self-heal across retries and re-runs.