Making a Data Ingestion Pipeline Safe to Retry Without Duplicating Records
A data ingestion pipeline is safe to retry when every write is idempotent, meaning a repeated attempt produces the same result as the first one. Network retries, manual replays, and consumers re-reading from an earlier offset are inevitable, and without idempotency each one risks a double counted metric or an inflated total.
Idempotency isn't a single technique, it's a property you design for at each step where a retry is possible, and the right mechanism differs depending on whether you're writing to a database, publishing to a queue, or calling an external API.
How do you choose an idempotency key that is actually unique?
The most common mistake is using a key that isn't actually guaranteed unique per logical event: a timestamp that two events can share, or an auto incrementing ID from a source system that resets or collides across partitions. A reliable idempotency key usually combines a stable source identifier with something that uniquely identifies the specific occurrence, like a source system's own event ID, or a hash of the event's immutable fields if the source doesn't provide one.
Whatever key you choose, treat it as part of your schema, not an implementation detail. Document where it comes from and what happens if a source system ever changes how it generates that identifier, since a silent change upstream can quietly break your deduplication without any obvious error.
Why deduplicate at write time instead of afterward?
Deduplication is far cheaper to enforce at write time, with a unique constraint or an upsert keyed on your idempotency key, than to detect and clean up after the fact. An upsert that writes a record if it doesn't exist and safely no ops if it does, keyed on the idempotency key, makes retries genuinely safe rather than merely usually safe.
For an append only event log where you can't upsert in place, track processed idempotency keys separately, checking against that set before writing, and expire old keys on a schedule that matches how far back a real retry could plausibly reach. Keeping every key forever is unnecessary and expensive; keeping too few risks a legitimate late retry being wrongly rejected.
Handling partial failures without losing idempotency
The trickiest case isn't a full retry of an untouched batch, it's a batch that partially succeeded before failing: some records were written, then the job crashed before finishing. A naive retry of the whole batch double writes the records that already succeeded. Design your batch processing so each individual record write is independently idempotent, not just the batch as a whole, so a partial retry safely re processes only what didn't complete.
Log which records in a batch succeeded before a failure, not just that the batch failed, so a retry can skip ahead instead of blindly reprocessing everything from the start. This also makes debugging a partial failure much faster, since you know exactly where the batch stopped.
Make partial failures safe to retry with these steps:
- Make each individual record write idempotent on its own, not just the batch as a whole.
- Key every write on the idempotency key, so a retried record that already exists is a safe no op.
- Once each write is idempotent, retry the whole batch after a crash, since records that already succeeded are skipped instead of written twice.
- Run the same batch through a test environment twice and diff the resulting state to confirm nothing doubles.
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 in a test environment and diff the resulting state. It's a small amount of test infrastructure that catches a surprising number of bugs, especially around aggregate tables or counters that get incremented rather than upserted, which are the most common place idempotency quietly breaks.
Make this test part of your normal test suite for any new ingestion path, not a one time manual check. Pipelines change over time, and a code change that breaks idempotency months after the original design rarely announces itself until a retry actually happens in production.
Cleaning up a duplicate that already happened
Even a well designed pipeline can end up with historical duplicates from before the idempotency work was in place, and cleaning those up is a separate task from preventing new ones. Write a one time reconciliation job that identifies duplicates by the same idempotency key logic you use going forward, and run it against a copy of the data first so you can verify the identified duplicates are actually duplicates before deleting anything from production.
Keep the reconciliation job's output as an audit log, listing exactly what was removed and why, in case a downstream consumer or a finance team needs to explain a change in historical totals later. A silent cleanup that nobody can explain months later tends to cause more distrust than the original duplicate did.
What Good Looks Like
A genuinely idempotent pipeline uses a verified unique key, deduplicates at write time rather than after the fact, and handles partial batch failures without double writing whatever already succeeded.
Building The Capability (5-Stage Skill Ladder)
How to Get Started
Frequently Asked Questions
What makes a good idempotency key for a data ingestion pipeline?
A value that's guaranteed unique per logical event, not just per batch. A stable source system event ID is ideal; if the source doesn't provide one, a hash of the event's immutable fields works. Avoid timestamps or auto incrementing IDs alone, since both can collide across sources or partitions in ways that aren't obvious until they happen.
Is it enough to make the whole batch idempotent, or does each record need to be idempotent individually?
Each record needs to be independently idempotent. A batch that partially succeeds before failing will otherwise double write the records that already completed on retry. Design writes so a partial retry only affects what didn't finish, and log progress within the batch so a retry can pick up where it actually stopped.
How do we know our pipeline is actually idempotent instead of assuming it is?
Run the same batch through your pipeline twice in a test environment and compare the resulting state. This is the most reliable way to catch idempotency bugs, particularly around counters or aggregate tables that get incremented instead of upserted, which is the most common place the property quietly breaks.
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
Verifying Every Service That Talks to Your Pipeline
Which parts of zero-trust verification to build and which to buy, so every producer and consumer on a streaming pipeline proves its identity.
Blue-Green, Canary, or Rolling: Deploying Stream Processors
A decision guide to rolling, blue-green, and canary deploys for stateful stream processors, plus the rollback plan most teams never actually test.
Deciding Where Your Event Pipeline Can Store Data
A decision guide for handling data residency and sovereignty requirements in a real-time event pipeline that spans more than one region.
Where Latency Actually Hides in a Growing Data Pipeline
A walkthrough of where latency hides as a real-time pipeline grows, from producer batching to consumer lag, so you can find your own bottleneck fast.
How to Run a Security Audit on a Real-Time Data Pipeline
A step by step way to check access, encryption, and patch timelines on your event streams before an incident or an auditor finds the gap first.
Handling GDPR Erasure Requests in a Streaming Pipeline
Answers to the privacy questions a real-time pipeline actually raises: erasure across replicated topics, data minimization, and cross-border transfer.