0%
Data-Intensive Applications
Foundations of Data Systems
Distributed Data
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Building Reliable Data Pipelines
Data pipelines fail. Networks drop packets, services crash mid-processing, upstream systems send the same message twice. The question is never whether your pipeline will see a duplicate or a retry, but whether your pipeline produces the correct result when it does. Building reliable pipelines means designing for three guarantees: no data loss (every record reaches its destination), no unintended duplicates (or at least idempotent handling of duplicates), and recoverability (when something breaks, you can fix it and reprocess without manual data surgery).
Reliability starts with idempotency: processing the same record twice produces the same result as processing it once.
An operation is idempotent if applying it multiple times has the same effect as applying it once. Formally, f(x) = f(f(x)) for the same input x. Setting a user's email to "[email protected]" is idempotent: doing it ten times leaves the same email. Incrementing a counter by 1 is not idempotent: doing it ten times adds 10 instead of 1. Inserting a row into a table without a unique constraint is not idempotent: each execution adds another row. Upserting a row by primary key is idempotent: the final state is the same regardless of repetitions.
This distinction determines whether your pipeline can safely retry on failure or whether a retry corrupts your data. Every pipeline stage should be analyzed for idempotency before it goes to production. If a stage is not idempotent, either redesign it (replace INSERT with UPSERT, replace append with overwrite) or add explicit deduplication logic upstream of that stage.
Idempotency Keys
The standard technique for making non-idempotent operations safe is an idempotency key: a unique identifier attached to each message or record. Before processing a record, the pipeline checks whether that key has already been processed. If yes, skip it. If no, process it and record the key.
The key must be generated by the producer, not the consumer. If the consumer generates the key (for example, by hashing the message body), two genuinely different messages with the same content would collide. If the producer generates the key (a UUID, a sequence number, or a composite of entity ID plus event timestamp), each logical event gets exactly one key regardless of how many times it is delivered.
The storage for processed keys matters. An in-memory set is fast but lost on restart, causing the pipeline to reprocess everything it handled before the crash. A database table is durable but adds latency to every record. A common middle ground is a write-ahead approach: the pipeline writes the idempotency key and the processing result in the same transaction. If the key already exists, the transaction is a no-op. If the pipeline crashes before committing, neither the key nor the result is persisted, so the retry processes the record from scratch and produces the correct outcome.
The idempotency key store also needs a retention policy. Storing every key forever means unbounded growth. If your pipeline processes 10 million messages per day, the key store grows by 10 million entries daily. After a year, that is 3.65 billion entries. The solution is a TTL (time-to-live) on key entries. If duplicates always arrive within 24 hours, a 48-hour TTL provides safety margin while keeping the store bounded. Keys older than the TTL are automatically purged, and any duplicate arriving after the TTL would be treated as a new message. This is acceptable because the probability of a duplicate arriving days later is vanishingly small in well-designed systems.
The whole mechanism is one table and one statement:
Two details carry the guarantee. The key insert and the business write are in the same transaction, so there is no window in which the key is recorded but the work is not, which is the failure that turns a duplicate-suppression table into a silent data-loss mechanism. And the ON CONFLICT DO NOTHING returning zero rows is the duplicate test itself, rather than a separate SELECT beforehand: a check-then-act pair leaves a race between two concurrent consumers holding the same redelivered message, and the unique index does not.
The TTL sweep is the part that gets forgotten until it becomes an incident. At 10 million messages a day the index alone grows past 10 GB in a year, and the DELETE that finally runs against it locks for long enough to stall the pipeline it was protecting. Partition the table by day and drop partitions instead, for the same reason developed in partitioning and sharding.
At-Least-Once Plus Idempotency Equals Exactly-Once
Distributed message systems offer two delivery guarantees: at-most-once (fire and forget, messages may be lost) and at-least-once (messages are retried until acknowledged, but may arrive more than once). True exactly-once delivery across network boundaries is impossible without coordination overhead that most systems cannot afford.
The practical solution is at-least-once delivery combined with idempotent processing. The message system guarantees every message is delivered at least once. The consumer guarantees that processing the same message multiple times has the same effect as processing it once. Together, they achieve effectively exactly-once semantics without the coordination cost of true exactly-once protocols.
This pattern appears everywhere in production systems. Kafka consumers commit offsets after processing. If the consumer crashes after processing but before committing, the message is redelivered on restart. Without idempotent processing, this redelivery creates duplicates. With idempotent processing (checking the idempotency key before writing), the redelivery is harmless. The consumer detects that the key was already processed and skips the record.
The distinction between delivery semantics and processing semantics matters. Kafka's "exactly-once semantics" (EOS) feature, introduced in version 0.11, provides exactly-once within a Kafka-to-Kafka pipeline by using transactional producers and read-committed consumers. But the moment data leaves Kafka and enters an external system (a database, a search index, an API), Kafka's EOS guarantees end. The external write might succeed while the Kafka transaction aborts, or vice versa. At that boundary, you are back to at-least-once delivery plus idempotent processing as the practical solution. Understanding where Kafka's transactional guarantees end and where your idempotency logic must begin is essential for designing correct end-to-end pipelines.
Mid-level engineers implement idempotency checks for individual pipeline stages. Senior engineers design end-to-end exactly-once semantics by combining at-least-once delivery with transactional idempotency key storage. Staff engineers build platform-level abstractions (idempotency middleware, transactional outbox patterns) that make every pipeline idempotent by default, so individual teams do not need to reinvent the pattern.
Transactional Outbox Pattern
A common failure scenario occurs when a pipeline needs to update a database and publish a message to a queue. If the database write succeeds but the message publish fails, the system is inconsistent. If the message publishes but the database write fails, downstream consumers act on data that does not exist.
The transactional outbox pattern solves this by writing the outgoing message to an "outbox" table in the same database transaction as the business data. A separate process reads the outbox table and publishes messages to the queue. Because the business data and the outbox entry are in the same transaction, they either both commit or both roll back. The publishing process is idempotent: if it crashes after publishing but before deleting the outbox entry, it re-publishes the message on restart, and the downstream consumer's idempotency key check handles the duplicate.
This pattern converts a distributed coordination problem (two-phase commit across database and queue) into a local transaction problem (single database transaction) plus an idempotent publishing problem (safe to retry). The tradeoff is added complexity in the form of the outbox table and the publishing process, but the correctness guarantee is worth it for pipelines where data loss or duplication has business impact.
FOR UPDATE SKIP LOCKED is what makes the relay horizontally scalable without a leader election: each instance claims a disjoint block of rows and never waits on another instance's block. ORDER BY id inside that claim preserves per-table ordering, though not ordering across aggregates, which is a distinction worth stating out loud before someone downstream assumes otherwise.
The relay is deliberately at-least-once. A crash between the broker accepting the message and the UPDATE committing republishes the event on the next pass, which is why the consumer's idempotency key is not optional infrastructure but the other half of this pattern. Polling is also only one implementation: reading the database's own write-ahead log instead of the outbox table removes the poll interval and the published_at update entirely, which is change data capture applied to a table you created for the purpose.
Common Idempotency Pitfalls
Several operations appear idempotent but are not. A pipeline that writes processed_at = NOW() generates a different timestamp on every execution, so the output changes on rerun even though the input is the same. Replace NOW() with the batch's logical timestamp (the scheduled run time, not the wall clock time). A pipeline that generates UUIDs for output records creates different IDs on rerun. Replace random UUIDs with deterministic identifiers derived from the input data (a hash of the natural key fields).
External API calls are another trap. A pipeline that sends a confirmation email during processing sends it again on retry. A pipeline that calls a payment gateway charges the customer twice. The fix is to separate pure data transformation (which can be made idempotent) from side effects (which must be guarded by their own idempotency mechanisms, such as the payment gateway's own idempotency key parameter).
Database operations also vary in idempotency. An INSERT is not idempotent if the table lacks a unique constraint on the natural key, because repeated inserts create duplicate rows. An UPSERT (INSERT ON CONFLICT UPDATE) is idempotent because it converges to the same state regardless of how many times it runs. A DELETE WHERE condition is idempotent because deleting zero rows (already deleted) and deleting one row produce the same final state. Understanding which SQL operations are naturally idempotent helps you design pipeline stages that are safe to retry without additional deduplication logic.
Partition-based overwrite is one of the most reliable idempotency techniques for batch pipelines. If your output table is partitioned by date, each pipeline run overwrites exactly one partition (today's date). Rerunning the pipeline for the same date replaces the same partition with identical data. Other partitions remain untouched. This gives you idempotency at the partition level without merge logic, deduplication keys, or coordination with downstream consumers. In Spark, this is a single configuration: df.write.mode("overwrite").partitionBy("date"). In SQL warehouses, it is a CREATE OR REPLACE TABLE ... AS SELECT scoped to the target partition.