ContentsThe library

Moving Data Without Losing Any

Nothing Crashed, Nothing Alerted, and the Total Is Short by Nine Hundred

Data moving between systems fails in exactly four ways, all of them quiet. Knowing which ones a given transport can produce is what tells you which checks are worth building.

Moving data from one system to another is the simplest-sounding thing in this

library and the place where the most undetected damage happens. The reason is

in the title of this lesson: when it goes wrong, nothing goes wrong.

Loss, duplication, reordering, corruption

There are four failures. Every data-movement incident is one of them or a

combination, and the list is short enough to keep in your head.

Loss. A record that existed at the source does not exist at the destination. A

batch file was written while the reader was mid-scan, a worker crashed between

reading and writing, a message expired in a queue nobody was draining, a filter

excluded more than it was meant to.

Duplication. A record exists twice at the destination. Almost always caused by

a retry: the receiver processed the message and the acknowledgement was lost,

so the sender, correctly, sent it again.

Reordering. Records arrive in a different order than they were produced. Two

workers pulling from the same queue finish at different speeds, so an update

lands before the insert it depends on, or an old value overwrites a new one.

Corruption. A record arrives changed. A number truncated by a narrower column,

a date reinterpreted in another time zone, text mangled by an encoding

mismatch, a decimal rounded by a transport that only carries floating point.

FIG 1A duplicate, made entirely by correct behaviour
stepstepwhat happenedrecords at the destinationdid anything report a problemwhat happened
11sender posts a batch of 500 records0noAn ordinary request. The sender starts a timer, because it has to do something if no answer comes back.
22receiver writes all 500 and commits500noThe work is done and durable. Everything is correct at this instant.
33the response is lost on the way back500noA network hiccup, a load balancer restart, a connection reset. The receiver believes it succeeded, because it did.
44the sender times out and retries, as des500noThis is not a bug. A sender that gave up here would cause loss instead, which is worse. There is no configuration of the sender that avoids both.
55receiver writes all 500 again and commit1000noBoth parties behaved correctly and the destination now holds double. Nothing in any log says so, and the only trace is a count that does not match.
5 steps
Five steps, no mistakes, and the destination is wrong. This is the single most important trace in the course, because it shows that duplication is not caused by a defect to be fixed but by an unavoidable ambiguity: a sender that receives no response cannot tell a lost request from a lost response. Choosing to retry produces duplicates and choosing not to produces loss, so the problem has to be solved at the destination instead, which is the next three lessons.

Why nothing goes red

Consider what would have to notice. The sender saw a timeout, which it is built

to handle and did handle. The receiver saw two valid requests and processed

both correctly. The database accepted two valid inserts. The monitoring saw

successful responses and normal latency.

Every component did its job. The failure exists only in the relationship

between the source and the destination, and nothing in the path is looking at

that relationship, which is exactly the end-to-end argument: the only place

this can be checked is at the ends.

FIG 2How much accumulates before anybody looks
records moved per unit of time
time between reconciliations
the fraction of records affected
wrong records sitting at the destination when somebody finally checks
The term that decides the damage is the middle one, and it is the only one under your control. A fault rate of one in two thousand is invisible in every dashboard and produces three thousand wrong records a day at a modest rate. The difference between reconciling daily and reconciling never is not the fault rate, which is unchanged, but whether the wrong records number in the thousands or in the millions by the time anybody notices, and whether the source still has the data needed to repair them.

What your path can actually do

The four failures are not equally available to every transport, and that is

useful, because it tells you which checks are worth building.

FIG 3What each transport can do to your data
can losecan duplicatecan reordercan corrupt
nightly file copy1001
single ordered log, one 1100
queue with several worke1110
at-least-once delivery, 0100
two systems written by t1111
an interface called over1101
The second marked cell is empty and is the one property worth designing for: a single ordered log read by one consumer cannot reorder, which removes an entire class of failure from the system and is most of the argument for that shape. The first marked cell is the row to be frightened of. Writing the same fact to two systems from one piece of code can produce every failure on the list, because the second write can fail after the first has committed, and nothing anywhere is responsible for the pair.

The row worth dwelling on is the nightly file copy, because it looks the safest

and can lose an entire day in silence. If the job does not run, nothing

complains: there is no error, no partial result, and no record anywhere of a

run that did not happen. A missing day looks exactly like a quiet day until

somebody adds up the month.

Which check catches which

Finally, match each failure to what would find it. This is the list the last

lesson of the course turns into a reconciliation.

FIG 4Each failure, and the cheapest check that finds it
count at each endsum of a columnhash per recordsequence numbers
a record was lost1110
a record arrived twice1110
records arrived out of o0001
a value arrived changed0110
The first marked cell is the one failure with exactly one answer: nothing except a position carried on each record can detect reordering, because every other check here is order-blind by construction. The second marked cell is the cheapest check in the table and catches half the rows, which is why it is the one to build first even when it is not sufficient. Note that the count column cannot distinguish its two rows: one lost and one duplicated record net to zero, so a matching count is weaker evidence than it feels.

A count at each end catches loss and duplication, and cannot tell them apart,

because one lost and one duplicated record net to zero. This is the cheapest

check and the one to build first anyway.

A sum of some numeric column catches both separately and catches most

corruption of that column as a bonus. Two checks, a count and a sum, cover a

surprising share of everything in this lesson.

A hash of each record catches corruption anywhere in it, including in columns

nobody thought to sum, and is the only thing that catches a value silently

truncated or re-encoded.

And a sequence number is the only thing that catches reordering, because every

other check is order-blind by construction. If the transport can reorder and

the order matters, the records have to carry their position.

One closing point on practice. These checks should be designed before the

pipeline rather than after the incident, for a reason that has nothing to do

with discipline: the check often changes the design. If the only way to detect

loss is a count, the pipeline needs to know how many records it was supposed to

move, and knowing that is a property of how it reads the source.

The next lesson takes the ambiguity from the trace above and derives what

delivery guarantees are actually available, which turns out to be fewer than

most products claim.

Recap

  • There are four failures and no others: a record is lost, a record arrives twice, records arrive in the wrong order, or a record arrives changed.
  • None of the four raises an error, because each of them is a normal outcome of a mechanism working as designed, which is why they are found by counting rather than by monitoring.
  • Each transport can produce some of the four and not others, so the check you build should match what your path can actually do to the data.

This is the reading half

Starting the course gives you your own copy of it. Every idea on every page has problems standing under it, marked with a reason rather than a tick, and any sentence you do not believe can be opened and argued with. None of that can happen on a page nobody owns.

The contents

NextAt Least Once, At Most Once →

The rest of this course

  1. 01Nothing Crashed, Nothing Alerted, and the Total Is Short by Nine Hundredyou are here
  2. 02The Sender Cannot Tell Which Half of the Journey Failedopening only
  3. 03Run It Twice on Purpose and See What Breaksopening only
  4. 04Two Writes, and the Order Between Them Is the Whole Designopening only
  5. 05The Row Was Saved at 10:04 and Became Visible at 10:06opening only
  6. 06Start the Stream Before You Read the Old Rowsopening only
  7. 07Tuesday's Number Arrived on Friday and Tuesday Is Already Publishedopening only
  8. 08The Check Takes Four Minutes a Day and Almost Nobody Runs Itopening only

Read alongside