Why Your Data Pipeline Needs to Survive Being Run Twice
A data pipeline will get run twice on the same input eventually, whether from a retry after a timeout, an operator replaying a failed batch, or a message queue delivering the same event more than once. If the pipeline isn't built to handle that, the result is duplicate rows, double-counted revenue in a dashboard, or a customer getting billed twice from a reprocessed event.
Idempotency, the property that running a job twice produces the same result as running it once, is the fix, and it has to be designed in, not bolted on after the first incident.
Where Duplication Actually Sneaks In
The most common source isn't a bug in your pipeline logic, it's the delivery guarantee of whatever's feeding it. Most message queues and event streams offer at-least-once delivery, not exactly-once, which means your pipeline should assume every message might arrive more than once. A network timeout during a write can also cause a retry that reprocesses data the first attempt actually succeeded at writing, just without the caller ever finding out. Both cases look identical from the pipeline's perspective: the same data, arriving again.
The Core Pattern: A Natural or Assigned Key
Idempotency depends on being able to answer one question cheaply: have I already processed this exact unit of work? That requires a stable key, either a natural one already present in the data, an order ID, an event ID, or one you assign at the point of ingestion if the source doesn't provide one. With that key, an upsert, insert if the key doesn't exist, update or skip if it does, replaces a blind insert, and reprocessing the same record becomes a no-op instead of a duplicate row.
Handling Aggregations, Not Just Raw Inserts
Upserts solve duplication for row-level ingestion, but aggregation steps, summing revenue into a daily total, incrementing a counter, need a different approach, because blindly adding a value twice is exactly the bug idempotency is supposed to prevent. The fix is to make the aggregation itself derivable from the underlying idempotent records rather than an incrementing counter: recompute the daily total from the deduplicated rows instead of adding a delta to a running number. It costs more compute per run, but it means a reprocessed batch can never silently inflate a total, since the total is recalculated from source data that's already protected against duplication.
A Worked Example: A Reprocessed Payment Webhook
Say your payment provider's webhook fires for a successful charge, your pipeline starts processing it, and the connection drops before your service can send back an acknowledgment. The provider, having not received confirmation, resends the same webhook a minute later. If your handler just inserts a new revenue record on receipt, you now have the same charge recorded twice. If it instead upserts on the payment provider's event ID, checking whether that ID has already been recorded before writing anything, the second delivery becomes a harmless no-op, and your revenue numbers stay correct without anyone having to notice or intervene.
Testing Idempotency Instead of Assuming It
The only reliable way to know a pipeline is actually idempotent is to run the same batch through it twice, on purpose, in a test or staging environment, and diff the output. Assuming idempotency because the code looks like it should behave that way is how these bugs end up in production. Make replaying a batch a standard part of your testing for any new ingestion job, not a special case you only think of after the first duplicate-row incident.
A quick idempotency checklist for any ingestion job:
- Assume every message can arrive more than once, since most queues and event streams offer at-least-once delivery rather than exactly-once.
- Give each unit of work a stable key, natural or assigned at ingestion, and upsert or skip based on it.
- Derive aggregates from the idempotent underlying records instead of blindly incrementing a counter.
- Replay the same batch twice in staging and diff the output before trusting the job in production.
- For a pipeline you did not build this way, add a deduplication table checked before any write happens.
What to Do About Pipelines You Didn't Build This Way
Retrofitting idempotency into an existing pipeline is harder than building it in from the start, but it's rarely as bad as rewriting the whole thing. Start by adding a deduplication table that records the key of every unit of work already processed, checked before any write happens, even if the rest of the pipeline logic stays exactly as it is. This is less elegant than redesigning the pipeline around upserts throughout, but it closes the actual gap, duplicate processing, without requiring you to touch every downstream step at once. Once that safety net is in place, you can migrate individual steps to a cleaner upsert-based pattern over time instead of under pressure during an incident.
What Good Looks Like
A safe pipeline means every ingestion and aggregation step can be rerun on the same input without changing the result, verified by actually replaying a batch, not assumed from the code's structure.
Building The Capability (5-Stage Skill Ladder)
How to Get Started
Frequently Asked Questions
Does using a message queue with exactly-once delivery mean we don't need idempotency?
Not entirely. Exactly-once delivery guarantees reduce the frequency of duplicates but rarely eliminate every source, since a retry at the application layer, a replayed batch, or a bug that resends data can still produce a duplicate that the queue's delivery guarantee never sees. Idempotent processing is cheap insurance even with a strong delivery guarantee upstream.
What if the source data doesn't have a natural unique key we can use?
Combine fields that should be unique together, such as a source system ID, a timestamp, and an amount, or assign your own key at the point of ingestion. The key does not need to come from the source system, as long as it is applied consistently every time that same unit of work is seen.
Is recomputing aggregates from scratch too slow for a large dataset?
For very large datasets, recomputing an entire aggregate on every run can get expensive, but you rarely need to recompute everything. Recompute just the time window that could have been affected by the reprocessed batch, a day or an hour, rather than the full history, which keeps the cost bounded while still avoiding the double-counting risk.
About the numbers
This guide doesn't quote a sourced benchmark. Figures in it are estimates or general guidance, so check them against your own numbers.
Related Guides
How to Ship a Risky Change Without a 2am Rollback
A concrete walkthrough of how to plan a risky production deployment: how to split it, what to watch, and when to decide the rollback trigger.
Build or Buy for Verifying Every Device That Connects?
How to split device identity from device posture checking, what building either one in house actually costs, and where a platform earns its keep instead.
Making Your Data Pipeline Safe to Rerun
A nightly ETL job fails halfway through, someone reruns it, and revenue gets double counted. A worked example of building a pipeline safe to replay.
Diagnosing Slow Requests Before You Blame the Database
A step-by-step way to find out whether a slowdown is the network, the app, or the database, before you add caching or upgrade infrastructure to fix it.
Where Your Data Actually Lives, and Why It Matters
How to figure out which of your data actually falls under residency or sovereignty rules, and what to check before assuming your cloud region is enough.
The Data Pipeline Bug That Only Shows Up After a Retry
Why a retried job silently duplicates data in most pipelines, and the idempotency key pattern that makes a pipeline safe to rerun from any failure point.