Hands-on: an incremental pipeline you can rerun without duplicating data
Clients will give you three years of old data and a nightly job. Both are only safe if the second run gives exactly the same result as the first.
In brief
- A watermark only helps select data; protection against duplicates belongs in the write step.
- MERGE on a business key, or DELETE then INSERT a whole partition in one transaction, makes reruns leave results unchanged.
- Once each partition is idempotent, backfill is just rerunning each day in order.
- 1Receive batch_dateScheduler passes in the logical date; each run processes one partition
- 2Select dataFilter by date; optionally move the watermark back to catch late rows
- 3Write idempotentlyMERGE on key or overwrite the partition in a transaction; use batch_date, not NOW()
- 4Rerun to verifyRun the same date twice, count duplicates and compare results
- 5BackfillLoop over each historical date, calling the same nightly job
Once the write step is idempotent, reruns and backfills give the same result.
Graphic: FDE Times
Picture your first week at a retail client. Last night’s order-loading job died halfway through, an operator hit rerun, and this morning the dashboard shows revenue doubled. Nobody wrote a wrong line of SQL. The problem is that the job was never designed to run twice.
If you work as an FDE, plan on meeting this situation.
Sooner or later the client will ask: “Can you load our data from the start of last year?” If your pipeline can rerun without creating duplicates, you can say yes without losing a night’s sleep.
This guide builds a daily orders pipeline with three properties: it loads incrementally, it is safe to rerun, and it backfills using the same code that runs every night. The SQL below is simplified; exact syntax will vary with the client’s data warehouse.
What do you need before you start?
IBM describes a data pipeline as three stages: ingesting raw data from multiple sources, transforming it, and storing it in a warehouse. The pipeline here runs in batch mode, loading data in chunks at set intervals. Each batch is triggered by a job scheduler, software that runs background jobs on a schedule without anyone watching over them.
You need a database that supports transactions and MERGE, plus two tables: staging_orders, holding freshly pulled raw data, and orders, the target table. Each order has order_id as its business key and order_date as its event date. A small Python script to call the job for each day is all you need on top.
Before writing a line, fix this definition in mind: an operation is idempotent if applying it many times leaves the result unchanged. That is what lets it be retried without side effects or corrupted data.
One point is often missed: idempotence is a property of each individual operation, not something that automatically holds for the whole system. Every step in the pipeline has to be designed for it separately.
Step 1: make the daily partition the unit of work
Each run processes exactly one logical day, called batch_date. The scheduler passes this date in as a parameter; the job never works out for itself what “today” is.
-- Every statement in the job filters on this parameter
SELECT order_id, customer_id, amount, updated_at, order_date
FROM staging_orders
WHERE order_date = :batch_date;
Check: run the query above for 2026-03-01 twice. The row count must be the same. If it differs, the source is changing while you read it, and you need to know that before going further.
Step 2: use MERGE so reruns produce nothing new
A plain INSERT is usually behind the doubled revenue dashboard. Say 2026-03-01 has 1,000 orders. Run INSERT twice and the target table holds 2,000 rows. MERGE (upsert) is different: new keys are inserted, existing keys are updated.
MERGE INTO orders AS t
USING (
SELECT order_id, customer_id, amount, updated_at, order_date
FROM staging_orders
WHERE order_date = :batch_date
) AS s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET
customer_id = s.customer_id,
amount = s.amount,
updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT
(order_id, customer_id, amount, updated_at, order_date)
VALUES (s.order_id, s.customer_id, s.amount, s.updated_at, s.order_date);
When the same batch is rerun, each row is updated to the value it already holds, so the second run changes nothing. In the example above, the target table keeps 1,000 rows however many times the job runs.
Check: run the MERGE twice, then look for duplicates. The result must be empty.
SELECT order_id, COUNT(*)
FROM orders
GROUP BY order_id
HAVING COUNT(*) > 1;
Step 3: with no key, overwrite the whole partition
Some tables have no stable business key, such as a daily revenue summary. The approach here is to treat the whole partition as the unit of work and replace it entirely. Rerunning a day replaces that day rather than doubling it.
BEGIN;
DELETE FROM daily_revenue WHERE order_date = :batch_date;
INSERT INTO daily_revenue (order_date, total_amount, order_count)
SELECT order_date, SUM(amount), COUNT(*)
FROM orders
WHERE order_date = :batch_date
GROUP BY order_date;
COMMIT;
BEGIN and COMMIT are the most important part of this block. DELETE followed by INSERT is only safe inside a single transaction. If the job crashes between the two statements, the partition is left empty, and the next morning the client sees zero revenue for that day.
Check: in a test environment, make the job stop right after the DELETE (by injecting a deliberate error, for example), then confirm the day’s old data is still intact.
Step 4: take NOW() out of the write
Many jobs write loaded_at = NOW() into every row. A pipeline that does this produces different data on every rerun, which means it is no longer idempotent. Use the event time or the batch’s logical date instead. This is still the write from step 3; only the INSERT inside the BEGIN ... COMMIT block changes:
-- Before: ... SELECT order_date, SUM(amount), COUNT(*), NOW()
-- After: write the logical date passed in by the scheduler
INSERT INTO daily_revenue (order_date, total_amount, order_count, batch_date)
SELECT order_date, SUM(amount), COUNT(*), :batch_date
FROM orders
WHERE order_date = :batch_date
GROUP BY order_date;
Check: run one day twice, export the results to files and diff them. If they differ in even one column, there is still a source of non-determinism somewhere in the job.
Why won’t a watermark save you?
The familiar approach to incremental loading is to store a watermark, such as the highest updated_at loaded so far, and next time fetch only newer rows. An analysis by Algoscale observes that watermark logic looks simple but is surprisingly fragile.
It breaks on timestamps sitting exactly on the boundary, on late-arriving rows and on time-zone mismatches, and in one case on Fabric it produced duplicate data.
So do not make the watermark responsible for correctness. Once the write is idempotent, with one row per business key, a row selected twice does not create a duplicate. According to Algoscale, duplicate selection then no longer corrupts the data; at worst it costs a little extra reading.
In practice, once the write is a MERGE, you can deliberately move the watermark back a little to catch late-arriving rows. The price is reading slightly more data; in return you stop worrying about losing rows.
Step 5: backfill is just a loop
At this point, backfilling the client’s historical data becomes simple: rerun each partition. Data Vidhya writes that teams with idempotent pipelines do backfills on an ordinary Tuesday afternoon. The script below is simplified; run_partition is the function that calls steps 1 to 4 for a given date.
from datetime import date, timedelta
def run_partition(batch_date):
# MERGE into orders, then overwrite daily_revenue for batch_date
...
d = date(2025, 1, 1)
end = date(2025, 12, 31)
while d <= end:
run_partition(d)
d += timedelta(days=1)
Check: run the loop twice over a short range, say three days. Row counts and totals must be identical between the two runs. If the loop dies on day 200, you simply rerun from that day with nothing to clean up.
Common pitfalls
The most common mistake is choosing the wrong business key. If the client’s source system reuses order_id across branches, MERGE will overwrite the wrong orders, and the correct key is the pair of branch and order number. Ask the client about this during customer discovery; don’t guess.
It is also common to find DELETE and INSERT in two separate jobs, or running in autocommit mode, which throws away the transactional safety net from step 3. And many teams test only the happy path. A pipeline that has never been run twice for the same day should be considered untested.
When you are on site with the client
Your first task at a new client is not to write a new pipeline but to ask: “What happens if we rerun last night’s job?” The answer tells you whether their system can survive retries and backfills.
Before agreeing to load three years of old data, check whether each target table is written with INSERT, MERGE or a partition overwrite.
If you are applying for FDE roles, look out for job descriptions that mention “data ingestion”, “backfill” or “pipeline reliability”. On a CV, “built ETL pipelines” says almost nothing. Write instead that you moved a job from INSERT to MERGE on a business key, and as a result backfilled a year of data with no duplicates.
Every pipeline breaks eventually. What is worth showing off is a pipeline that, once it breaks, only needs to be rerun, with no one left cleaning up data.
Was this article useful?
Thanks for the feedback!