Roadmap To Be A Data Engineer / Lesson 04

Lesson 04 Idempotency Fundamentals §5 Orchestration & Scheduling About 10 min read

The Payout That Ran Twice

A retry pays sellers twice. Safe re-runs and backfills.

1The situation

The payout job pays consignors 60% of what their item sold for. It runs at 03:00, reads yesterday’s sales, writes one payout row per sale. At 09:00 a file of those rows goes to the bank.

Night of 25 August: the job writes 1,247 of 2,000 rows and dies — the warehouse stage runs out of disk. On-call clears the disk and hits rerun. The job starts at the top, reads all 2,000 sales, and appends 2,000 more rows next to the 1,247 already committed.

Time What happened
03:00 Job starts on 2,000 sales from 25 Aug
03:12 Crash after 1,247 rows — already-written rows stay
07:40 On-call reruns; job appends all 2,000. Table holds 3,247
08:05 Pipeline green. Freshness passes; no test counts rows per sale
09:00 Bank file leaves: 3,247 payments, €486,745.80
14:20 Finance finds out — from a consignor asking why she was paid twice

€187,015.80 paid a second time, in one night. Note the nasty detail: the overpayment equals exactly the value of the rows the crashed run had already written — so the damage depends on how far the job got before it died, is different every time, and is never a round number anyone would spot.

2The mechanism

The bug is not the crash. Crashes are normal — networks blip, disks fill, APIs rate-limit, someone deploys mid-run. The bug is that the job produced a different result the second time it ran on the same input.

A job is idempotent when running it twice on the same input leaves the world exactly as running it once did. f(f(x)) = f(x)

That property is what makes a retry safe — and retries are the entire operational model of data engineering. Airflow retries. Kafka consumers redeliver (at-least-once means a record will arrive twice eventually). Humans rerun. If the job is idempotent, all of that is free. If not, every one of those mechanisms is a loaded gun.

Same crash, same two retries, two write patterns (verified output)

Append (INSERT INTO … SELECT) Merge (INSERT … ON CONFLICT DO UPDATE)
After run 1 (crashed, 1,247 rows) 1,247 rows · €187,015.80 1,247 rows · €187,015.80
After run 2 (full retry) 3,247 rows · €486,745.80 2,000 rows · €299,730.00
After run 3 (retry again) 5,247 rows · €786,475.80 2,000 rows · €299,730.00

2,000 / €299,730.00 is the truth. The merge job reaches it and stays there forever; the append job climbs by 2,000 every time someone presses the button.

3Four ways to write it

Idempotency is not a setting you enable — it is a property of how the job writes. Be able to name which shape any job uses.

Pattern Rerun-safe What it needs Use it when
Append — INSERT INTO … SELECT No Nothing, which is why it’s everyone’s first draft Never for money, counts, or anything feeding an aggregate. Only for an append-only event log deduplicated on read
Merge / upsert — MERGE, ON CONFLICT Yes A stable business key valid for all time (sale_id — not a row number, not a load timestamp) Row-level updates arriving continuously; CDC streams; anything keyed by an entity
Delete + insert by partition Yes A partition column (usually the date) and the discipline of never writing outside the partition you were given The workhorse for daily batch and backfills. Trivially parallel, and it fixes logic bugs — a merge only updates rows it sees
Insert with a dedupe key — unique constraint / ON CONFLICT DO NOTHING Yes A deterministic key derived from the payload — not uuid(), not now() At-least-once queues where a duplicate message should simply be dropped

The tell in a code review: does the job’s output depend on anything other than its inputs and parameters? now(), current_date, uuid(), random(), an auto-increment id, “the max timestamp currently in the target table” — each quietly makes the second run different from the first.

4The other failure: the hole

Duplicates are the loud version. The quiet one has the same root cause and is worse, because nobody emails about money they didn’t get.

-- looks harmless, is not
WHERE sale_date >= current_date - INTERVAL 1 DAY

The 25 August run fails; nobody gets to it until the 27th. The rerun asks for “yesterday” — the 26th. The 25th is never loaded, never retried, and never caught: the table has rows, freshness passes, the DAG is green. In this version of the story, €292,470.58 owed to 2,000 consignors is silently missing.

Pass the job the logical date it is processing as a parameter, and let it read nothing from the clock. A rerun of 25 August must still mean 25 August in March.

That is why Airflow, Dagster and Prefect hand every run an execution date instead of letting the task ask what time it is. That parameter is what makes reruns and backfills possible at all.

5Which is why backfills are the real test

A month later: premium items (≥ €400) pay 55%, not 60%. The rate has been wrong since launch. Thirty days need rebuilding.

A backfill is just the job, run again over a past window — so a backfill is only as safe as the job’s write pattern.

30 days · 60,000 sales · €14,983,500 GMV:

Backfill over 30 days Rows in payouts Total paid out Correct?
Before the fix — flat 60% 60,000 €8,990,100.00 Overpaid by €218,511.36
Append backfill — run once 120,000 €17,761,688.64 Every seller paid twice
Delete + insert — run once 60,000 €8,771,588.64 Correct
Delete + insert — run 3× 60,000 €8,771,588.64 Identical, byte for byte

The last row is the whole discipline: an idempotent job can be rerun by a nervous engineer at 2 a.m. who isn’t sure the first attempt finished, and nothing happens.

Backfill hygiene beyond idempotency: run into a staging table first and diff totals against production before swapping; chunk by partition so a failure at day 19 doesn’t cost you days 1–18; throttle it so the backfill doesn’t starve the nightly job of warehouse slots; tell whoever owns the downstream dashboard before yesterday’s revenue moves.

6Three questions to ask in a standup

  1. “If this job runs twice, what happens?” The right answer is instant and specific — “nothing, it merges on order_id“. Hesitation is the finding.
  2. “What does this job read to decide its window?” If the answer contains now() or current_date rather than a parameter the orchestrator passes in, reruns and backfills are already unreliable.
  3. “How would we backfill this for the last 90 days?” One sentence = an idempotent pipeline. “Carefully” = a manual one, and every historical bug will be expensive.

720-minute hands-on (DuckDB)

pip install duckdb. Every statement below was run as written; the output is the real one.

-- STEP 1 — a month of sales. Every item unique, quantity 1.
CREATE OR REPLACE TABLE sales AS
SELECT i AS sale_id,
       DATE '2026-08-01' + INTERVAL ((i-1)//2000) DAY AS sale_date,
       round(25 + ((i * 37) % 900) * 0.5, 2) AS price_eur
FROM range(1, 60001) t(i);

-- STEP 2 — two payout tables. The only difference is the write pattern.
CREATE OR REPLACE TABLE naive (sale_id BIGINT, sale_date DATE, amount_eur DOUBLE);
CREATE OR REPLACE TABLE merged(sale_id BIGINT PRIMARY KEY, sale_date DATE, amount_eur DOUBLE);

-- STEP 3 — the crashed run, then two full retries.
--          Run this block three times: LIMIT 1247 → 2000 → 2000.
INSERT INTO naive
  SELECT sale_id, sale_date, round(price_eur*0.60,2)
  FROM sales WHERE sale_date = DATE '2026-08-25' ORDER BY sale_id LIMIT 1247;

INSERT INTO merged
  SELECT sale_id, sale_date, round(price_eur*0.60,2)
  FROM sales WHERE sale_date = DATE '2026-08-25' ORDER BY sale_id LIMIT 1247
  ON CONFLICT (sale_id) DO UPDATE SET amount_eur = excluded.amount_eur;  -- ← the whole lesson

-- STEP 4 — compare after each run.
SELECT 'append' AS job, count(*) AS rows_, round(sum(amount_eur),2) AS paid FROM naive
UNION ALL
SELECT 'merge',  count(*),        round(sum(amount_eur),2)        FROM merged;
run 1 (LIMIT 1247):  append rows= 1247  paid=  187,015.80  |  merge rows= 1247  paid=  187,015.80
run 2 (LIMIT 2000):  append rows= 3247  paid=  486,745.80  |  merge rows= 2000  paid=  299,730.00
run 3 (LIMIT 2000):  append rows= 5247  paid=  786,475.80  |  merge rows= 2000  paid=  299,730.00

Then do the backfill. The contract fix is CASE WHEN price_eur >= 400 THEN 0.55 ELSE 0.60 END. Rebuild the month one partition at a time:

-- for each of the 30 days, in date order:
DELETE FROM payouts WHERE sale_date = :run_date;
INSERT INTO payouts
  SELECT sale_id, sale_date,
         round(price_eur * (CASE WHEN price_eur >= 400 THEN 0.55 ELSE 0.60 END), 2),
         CASE WHEN price_eur >= 400 THEN 0.55 ELSE 0.60 END
  FROM sales WHERE sale_date = :run_date;

Run the whole loop three times and check the total: it stays at €8,771,588.64. Now delete the DELETE line and run it twice: €17,761,688.64. One statement is the difference between a pipeline you can rerun and one you can only run once and hope.

8Takeaway

Pipelines do not fail gracefully because failure is rare. They fail gracefully because rerunning them is boring. Make every job a function of its parameters, key every write, and the retry button stops being frightening.

Vocabulary

idempotencyat-least-once / exactly-once deliverynatural (business) keysurrogate keyupsertMERGEdelete-insert by partitiondedupe keylogical / execution datebackfillreplaywatermarkpartition pruningstaging tablereconciliation

Next: Lesson 05 — star schemas, and why two analysts get two different revenue numbers from the same tables.

Back to top