By Mert Dönmezler12 min read

Data Contracts and Idempotency: Two Reliability Primitives for Pipelines

Most pipeline incidents trace back to two missing primitives: data contracts enforced at the producer boundary and idempotent writes at the consumer boundary.

  • data-engineering
  • reliability
  • data-contracts
  • research
  • idempotency
  • pipelines
  • schema-evolution

Pipeline incidents are rarely exotic distributed-systems failures. The argument here is that two missing primitives, data contracts enforced at the producer boundary and idempotent writes at the consumer boundary, account for most preventable data incidents between them. Schema-on-read only postpones the cost of validation, and exactly-once delivery is achievable in practice only as effectively-once semantics built on idempotent operations. The sections below show how to write a contract as an executable artefact, pick a compatibility direction on purpose and define an idempotency key for a given write path. They also cover planning for backfills at design time, so that a backfill does not turn into an operational emergency.

The producer–consumer coupling problem

Every pipeline has a contract, whether or not anyone has written it down. When a producing service emits an event or a table and a consuming job reads it, the consumer necessarily encodes assumptions: that a field is present, that a timestamp is UTC, that an amount is in minor currency units, that one row means one order and not one order line. Those assumptions are rarely declared anywhere. Someone inferred them from a sample of data on the day the consumer was written.

This is the mechanism behind Hyrum's Law, an observation attributed to Hyrum Wright: with enough consumers, every observable behaviour of a system becomes something someone depends on, regardless of what was promised. In a data platform the observable surface is unusually wide. Column order, string padding, the accidental sort order of a file, a nullable column that has never actually been null: all of these can be observed, so someone comes to depend on them, and they end up as part of the interface whether anyone meant them to or not.

The robustness principle, be conservative in what you send and liberal in what you accept (Postel, 1981), served internet protocols well at the transport layer, where the alternative was systems that could not talk to each other at all. Inside a data platform its liberal half does real damage. A consumer that quietly accepts a malformed record takes a producer bug that would have been obvious and easy to trace, and turns it into a data-quality problem spread across the platform that shows up weeks later in a dashboard. Part of what the relational model (Codd, 1970) got right, and what has lasted, is that structure is declared once and enforced independently of any program that reads it, so readers do not each have to work it out again.

So validation belongs at the boundary, facing the producer, and the reason is economic. At the boundary a violation has a single owner and an obvious fix, and the damage is contained. Once it has spread downstream, n teams each own a piece of it and nobody can fix it. The same argument makes automated validation at the point of production preferable to inspection after the fact.

Why schema-on-read defers cost rather than eliminating it

The case for schema-on-read has merit. Ingestion is cheap, producers are not blocked, and interpretation waits until someone has a use for the data. What the case leaves out is that interpretation is then pushed onto every reader, every time they read. If a raw stream has no enforced schema, each of k consumers has to discover a field's meaning on its own, write its own coercion and decide for itself what to do when the field is missing. Total validation work grows with k instead of staying at one, and the k implementations drift apart. Two dashboards that compute "active venues" from the same stream with different null handling raise no error and simply show two different numbers. Someone calls a meeting about it, and the trust people lose in the platform costs far more than the ingestion friction the design was meant to avoid. Schema-on-read still makes sense for real exploration and for archival capture when nobody knows what the data will be used for, but not for a stream with several production consumers.

Data contracts as executable artefacts

A contract in a wiki page is documentation. It becomes an artefact when it is machine-readable, version-controlled and enforced by the same pipeline that ships the code, which amounts to continuous-integration discipline applied to data. Three classes of obligation belong in it.

Structural obligations

Field names, physical types, nullability, cardinality, primary and foreign keys, permitted ranges and enumerations. These are the cheapest to check and the easiest to generate tests from. They can be verified statically against a schema registry and dynamically against live samples.

Semantic obligations

These are the obligations structural checks cannot express: the unit of a numeric field, the timezone of a timestamp, the grain of a row, the difference between a null meaning "unknown" and one meaning "not applicable", and whether a monotonic identifier is guaranteed to be monotonic or just usually is. Semantic drift is the most damaging failure class, because it passes every type check.

Operational obligations

A named owning team, a freshness objective, a documented deprecation window, and a change protocol stating how a breaking change is proposed, announced and rolled out. Without the owner, nobody enforces any of the rest.

In our own work, the clearest example of a deliberate contract surface is BrewX, a venue operating system in development. Its data model is a 46-table schema, and clients never touch those tables. All writes and most reads go through more than 80 remote procedures, and the contract is the set of procedure signatures. That indirection costs real effort, because every new capability needs a procedure where a query would otherwise do. In return, internal schema changes do not automatically break clients, and validation, authorisation and idempotency logic sit in one place where they can be enforced, instead of being reimplemented in each caller.

Schema evolution and compatibility directions

Contracts will change, so the useful question is which direction of compatibility you guarantee. That choice sets the deployment order, and getting it wrong causes outages that look like data bugs.

DirectionGuaranteeTypical safe changesDeploys first
BackwardNew readers read old dataAdd optional field, widen type, append enum valueConsumer
ForwardOld readers read new dataAdd optional field, remove optional fieldProducer
FullBoth directions holdAdd or remove optional fields onlyEither
NoneNo guaranteeRename, retype, change semanticsCoordinated migration

Two rules follow. Adding a required field without a default is a breaking change, even though it feels like an addition, because no existing writer supplies it. Changing what a field means while keeping its name and type is the worst option available. No automated compatibility check will see it. Only someone who remembers the old meaning could catch it, and usually nobody on the team still does. If the meaning has to change, add a new field and deprecate the old one on a stated timetable. Carrying two fields for a quarter almost always costs less than silently reinterpreting history.

Delivery semantics: at-most-once, at-least-once and exactly-once

Three guarantees are commonly named. At-most-once delivers a message zero or one times, so a failure loses data. At-least-once delivers it one or more times, and failures create duplicates. Exactly-once delivers each message precisely once.

In the general case, exactly-once delivery across an unreliable network between processes that fail independently cannot be achieved. The reason is simple. A sender that transmits and gets no acknowledgement cannot tell a lost request from a lost acknowledgement, and it is unsafe both to retry and to hold back. No ordering of events removes the ambiguity, because there is no global notion of simultaneity to fall back on (Lamport, 1978). The same tension shows up in the consistency–availability trade-off under partition described by Brewer (2000), where a system has to choose what to give up when it cannot tell a slow peer from a dead one.

What can be achieved, and what Kleppmann (2017) treats as the accurate way to state the goal, is effectively-once: at-least-once delivery plus an idempotent apply step, so a duplicate delivery leaves observable state exactly as a single delivery would. Formally the apply function must satisfy f(f(s, m), m) = f(s, m) for any state s and message m. Stated this way, the engineering work moves off the transport, which is allowed to be sloppy, and onto the write, which is not. The transactional machinery for this is old and well documented: atomic commit, durable logs and recovery semantics (Gray & Reuter, 1993). What is new in modern pipelines is that the write side is often a system without multi-statement transactions, which pushes idempotency out of the database and into the application.

Idempotency keys, natural keys and merge semantics

An idempotent write needs a stable identity for the unit of work. There are three options, from most to least preferred.

A natural key is best where one exists. It is the tuple of business attributes that uniquely identifies the entity, such as source system plus order number plus line number, and it stays stable across replays because it comes from the data and not from the delivery.

Next comes a producer-assigned idempotency key, a value generated once and carried through every retry:

key = hash(source_system, natural_key, logical_version)

What matters is that the key is computed before the first attempt and reused on every retry. If the key is generated inside the retry loop, each attempt gets a new identity for the same work, and nothing gets deduplicated.

A delivery-assigned identifier, such as an offset, message ID or file name, is the last resort. It deduplicates retries of the same delivery, but it misses the same business event arriving again by a different path, which is exactly what happens in a backfill.

With a key, the write becomes a merge on that key with an explicit conflict policy. Last-writer-wins is acceptable only when "last" is defined by a logical version carried in the payload. Arrival or ingestion time will not do, because arrival order is not a total order over causally related events. Where updates are partial, the merge has to state which columns it may overwrite, so that two producers writing different fields of the same row do not silently erase each other's work.

Our ERP-to-fulfilment workflow is a small but useful case. It pulls order data from an ERP, formats shipping labels and pushes them to the fulfilment system, and it runs more than 800 times a day. At that rate transient failures are routine, and operationally the question that matters is what a retry does. The unit of work is keyed by order identity and not by run. A retried step therefore produces the same label again and never creates a second shipment, whatever the transport does. Retry behaviour comes up again and again in the broader failure taxonomy for robotic process automation.

Backfills and replay as first-class requirements

A backfill is the real test of a pipeline design, because it exercises every assumption at once: contract compatibility with historical data, idempotency when records are redelivered through a non-standard path, and the correctness of every late-arrival rule. Until a pipeline has been through a backfill, nobody knows whether those assumptions hold.

Designing for replay has concrete requirements. Transformations have to depend on their declared inputs and parameters and not on the wall clock, so that reprocessing a past window gives the past answer. Writes need the keys described above, so that a replay overlapping live traffic converges instead of duplicating rows. Partition by event time rather than ingestion time, and a bounded replay will touch a bounded set of partitions. It also helps if the pipeline can backfill into a shadow destination for comparison before it touches production. In practice, a design is not finished until someone can say straight away what happens if the last thirty days are reprocessed while the stream is live.

Observability: freshness, volume, distribution, lineage

Contracts prevent one class of defects, and monitoring has to catch the rest. Four families of signal cover most of it. Freshness, the age of the newest record measured against its objective, catches stalls, which are the most common failure and the quietest. Volume catches partial and duplicate loads. Distribution (null rate, cardinality, category mix, numeric quantiles) catches the semantic drift that structural checks let through. Lineage answers the question every incident gets to sooner or later: what else downstream is affected.

Set thresholds by arithmetic. Here is an illustrative example that uses assumed figures, not measured ones. Suppose a table receives roughly 10,000 rows per hour, hourly volume is approximately normal with a coefficient of variation of 15%, and an alert fires when volume deviates by more than two standard deviations. A symmetric two-sigma rule then fires on about 5% of intervals from noise alone. With hourly checks that is roughly one false alarm a day, near 30 a month for a single table. With ten such tables, the alert makes enough noise that people start ignoring it. The numbers are hypothetical, but the point holds generally. An alert threshold is a statistical choice with a false-positive rate you can compute, and thresholds picked by feel are how monitoring ends up as decoration. The control-chart literature worked out this trade-off long ago, as the discussion of statistical process control under machine vision sets out.

Practical decision rules for contract and write design

Enforce a contract at any boundary that crosses an ownership line, and prefer rejecting an invalid record to leniently accepting it. Choose and document one compatibility direction per interface, and take the deployment order from it instead of from habit. Assume at-least-once delivery everywhere and make every write idempotent under a key computed from the data. Never define last-writer-wins by arrival time. A backfill rehearsal is part of the definition of done. Monitor freshness before anything else, and compute the false-positive rate of every threshold you set.

Limitations

This article argues from mechanism and from our own projects. It is not an empirical study. We have not measured what share of incidents in pipelines generally comes down to contracts and idempotency, and the claim that these two primitives dominate is an untested hypothesis that fits our experience. We also offer no evidence on the cost of adopting contracts, which is real and falls more heavily on producing teams, or on when that cost outweighs the benefit. For single-consumer or short-lived pipelines it plausibly does.

The BrewX and ERP examples describe design decisions and operating volumes. They are not controlled comparisons, and we cannot say what those systems would have cost or how they would have failed with a different design. The observability arithmetic is a calculation from stated assumptions, and no false-positive rate was observed. Finally, the article leaves out the organisational question of how contract ownership gets negotiated when a producer has no incentive to serve a consumer. In practice, that decides whether any of this gets adopted at all.

References

  • Brewer, E. (2000). Towards robust distributed systems.
  • Codd, E. F. (1970). A relational model of data for large shared data banks.
  • Gray, J., & Reuter, A. (1993). Transaction processing: concepts and techniques.
  • Kleppmann, M. (2017). Designing data-intensive applications.
  • Lamport, L. (1978). Time, clocks, and the ordering of events in a distributed system.
  • Postel, J. (1981). Transmission Control Protocol.