Learn Labs
9. The Trouble with Distributed Systems

9.1 Faults and Partial Failures

The single computer is a lie we've agreed to believe:

There is no fundamental reason software on a single computer should be flaky. When hardware works correctly, the same operation always produces the same result (it is DETERMINISTIC). If there is a hardware problem, the consequence is usually a TOTAL system failure — kernel panic, blue screen, failure to start. An individual computer with good software is either fully functional or entirely broken, but not something in between.

This is a DELIBERATE CHOICE in the design of computers. If an internal fault occurs, we prefer a computer to CRASH COMPLETELY rather than returning a wrong result, because WRONG RESULTS ARE DIFFICULT AND CONFUSING TO DEAL WITH. Thus computers HIDE the fuzzy physical reality on which they are implemented and present an IDEALIZED SYSTEM MODEL that operates with mathematical perfection.

(As Ch 2 showed, this is not actually true — data does get silently corrupted and CPUs do sometimes silently return the wrong result — but it happens rarely enough that we can get away with ignoring it.)

Across a network, that abstraction collapses.

"In my limited experience I've dealt with long-lived network partitions in a single data center, PDU failures, switch failures, accidental power cycles of whole racks, whole-DC backbone failures, whole-DC power failures, and a hypoglycemic driver smashing his Ford pickup truck into a DC's HVAC system. And I'm not even an ops guy." — Coda Hale

PARTIAL FAILURE: some parts of the system are broken in an unpredictable way while other parts work fine.

The difficulty is that partial failures are NONDETERMINISTIC: if you try to do anything involving multiple nodes and the network, IT MAY SOMETIMES WORK AND SOMETIMES UNPREDICTABLY FAIL. You may not even know whether something succeeded.

The compensation: if a system CAN tolerate partial failures, that opens powerful possibilities — rolling upgrades, rebooting one node at a time while the system keeps working.

Fault tolerance therefore allows us to make distributed systems MORE RELIABLE THAN SINGLE-NODE SYSTEMS; we can build a reliable system from unreliable components.