Lesson 16 / 25

Idempotent, Partitioned Tasks

Write each run's output to a deterministic, date-based location so re-runs overwrite rather than append.

Same input, same output, any number of times

Tasks fail, get retried and are re-run for past dates, so every task should be idempotent: running it twice for the same logical date leaves the same final state as running it once. Techniques: write to a path or partition derived from the run date (.../dt=2026-10-01/) and overwrite it; use MERGE/upsert or DELETE then INSERT for the date's rows instead of blind INSERT; write to a temporary location and atomically swap; never append without a de-duplication key. Idempotency is what makes retries and backfills safe.

Safe to run twice

Good pipelines can be re-run for any date without duplicating or losing data.

Three principles: idempotent, partitioned, lightweight.
Figure 5.1 — Idempotent, partitioned and lightweight.

Deterministic output path, run

I ran this: the path depends only on the logical date, so a retry or backfill writes to exactly the same place and replaces the old file.

def out_path(ds):
    return f"s3://lake/orders/dt={ds}/part.parquet"

print(out_path("2026-10-01"), out_path("2026-10-01") == out_path("2026-10-01"))

Output:

s3://lake/orders/dt=2026-10-01/part.parquet True

Delete-then-insert for a date (SQL)

Running this twice for the same date leaves one copy of that day's rows. {{ ds }} is templated by Airflow.

BEGIN;
DELETE FROM analytics.daily_orders WHERE order_date = '{{ ds }}';
INSERT INTO analytics.daily_orders
SELECT order_date, customer_id, SUM(amount) AS revenue
FROM staging.orders
WHERE order_date = '{{ ds }}'
GROUP BY order_date, customer_id;
COMMIT;

Quick check: Why should a task be idempotent?

  • Because Airflow forbids writing files
  • To make it slower
  • So retries and re-runs do not duplicate or corrupt data
  • To avoid using dates
Answer

So retries and re-runs do not duplicate or corrupt data — Safe repetition is what lets orchestrators retry and backfill confidently.