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.
| step | step | what happened | records at the destination | did anything report a problem | what happened |
|---|---|---|---|---|---|
| 1 | 1 | sender posts a batch of 500 records | 0 | no | An ordinary request. The sender starts a timer, because it has to do something if no answer comes back. |
| 2 | 2 | receiver writes all 500 and commits | 500 | no | The work is done and durable. Everything is correct at this instant. |
| 3 | 3 | the response is lost on the way back | 500 | no | A network hiccup, a load balancer restart, a connection reset. The receiver believes it succeeded, because it did. |
| 4 | 4 | the sender times out and retries, as des | 500 | no | This 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. |
| 5 | 5 | receiver writes all 500 again and commit | 1000 | no | Both 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. |
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.
- records moved per unit of time
- time between reconciliations
- the fraction of records affected
- wrong records sitting at the destination when somebody finally checks
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.
| can lose | can duplicate | can reorder | can corrupt | |
|---|---|---|---|---|
| nightly file copy | 1 | 0 | 0 | 1 |
| single ordered log, one | 1 | 1 | 0 | 0 |
| queue with several worke | 1 | 1 | 1 | 0 |
| at-least-once delivery, | 0 | 1 | 0 | 0 |
| two systems written by t | 1 | 1 | 1 | 1 |
| an interface called over | 1 | 1 | 0 | 1 |
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.
| count at each end | sum of a column | hash per record | sequence numbers | |
|---|---|---|---|---|
| a record was lost | 1 | 1 | 1 | 0 |
| a record arrived twice | 1 | 1 | 1 | 0 |
| records arrived out of o | 0 | 0 | 0 | 1 |
| a value arrived changed | 0 | 1 | 1 | 0 |
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