0%
Data-Intensive Applications
Foundations of Data Systems
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Operational Patterns
The Trouble with Distributed Systems
A program running on one computer has a comforting property: it either works or it crashes. If the CPU produced a wrong answer or memory returned a different value than was written, the machine would be considered broken and sent back. Hardware is engineered so that software can treat failure as total and deterministic. Your function returns, or your process dies and you find out about it.
Distributed systems do not offer that deal. Some components work while others are broken, and the working components cannot reliably tell which is which. This is partial failure, and every other difficulty in this lesson is a consequence of it.
The Six Fates of an Unanswered Request
You send a request to another node and no response arrives. Here is everything that could have happened:
- The request was lost in the network and never arrived.
- The request is sitting in a queue somewhere and will be delivered later.
- The remote node crashed before it saw the request.
- The remote node processed the request completely, then crashed before replying.
- The remote node processed the request and replied, and the response was lost.
- The remote node processed the request and replied, and the response is delayed behind a queue.
From where you are sitting, all six are the same observation: silence. You cannot distinguish them, not with a better client library, not with more logging, not with a faster network.
Look at which ones are which. In cases 1, 2, and 3 the work did not happen. In cases 4, 5, and 6 the work did happen and you have no idea. A timeout is not evidence of failure. A timeout is evidence of ignorance.
This single fact is why idempotency is not an optimization in distributed systems but a structural requirement. If you cannot tell whether a payment was captured, your only safe options are to make capturing it twice equivalent to capturing it once, or to never retry and accept losing the ones that silently succeeded. There is no third option in which the caller simply finds out.
Idempotency appears in several lessons in this course because it is the answer to several different questions. Here it is the answer to ambiguity: the caller cannot learn whether the work happened, so the only escape is to make the answer not matter. The implementation details, client-generated keys stored in the same transaction as the effect, are covered where delivery guarantees are discussed; what this lesson supplies is the reason the requirement is structural rather than defensive.
Networks Fail in Stranger Ways Than You Expect
The mental model most engineers carry is that networks either work or are cut. Field studies of production datacenters, most famously the survey work behind Kingsbury and Bailis's The Network Is Reliable, describe something messier.
Asymmetric partitions. Node A can send to B, but B cannot send to A. Each node draws a different and equally confident conclusion about who is alive. Failure detectors built on one-way pings will disagree about reality in a way that is stable rather than transient.
MTU mismatches. A misconfigured interface silently drops packets over a certain size while letting small ones through. Health checks are small and pass continuously. Real traffic is large and fails. The monitoring says the node is healthy and the users say it is not, and both are reporting accurately.
Gray failure. A link drops 3 percent of packets. Nothing is down. TCP retransmits, so everything still works, at ten times the latency, and every timeout you tuned for the healthy case starts firing.
Switch and firmware faults. A switch reboots and takes 90 seconds to reconverge. A firmware bug forwards traffic to the wrong VLAN. A misapplied BGP change blackholes a rack. None of these is your code and all of them are your problem.
Inside a single datacenter, with no cloud provider and no internet in the path, partitions lasting tens of seconds are a normal operational event rather than a rare one. Any design whose correctness assumes they will not happen is not conservative; it is wrong on a schedule.
Latency Is a Property of Queues
It is tempting to think of network latency as a physical constant for a given topology. It is not. The propagation delay is a constant; almost all of the latency you actually measure is time spent in queues.
A packet waits in the sending host's NIC buffer if the link is busy. It waits in each switch's buffer if the output port is contended. It waits in the receiving host's kernel buffer if no core is free. Then the request waits in the application's thread pool or accept queue. Each of those is a queue whose depth is a function of load, and each of them grows without bound as utilization approaches saturation.
This is why tail latency degrades so violently under load, and why an average-based timeout is close to useless. The system you tuned at 40 percent utilization behaves like a different system at 85 percent.
The underlying cause is a deliberate design choice. A circuit-switched telephone network reserves a fixed slice of bandwidth for the duration of a call, which gives it a hard bound on latency and jitter and wastes the capacity when nobody is talking. Datacenter networks statistically multiplex instead: they let bursty traffic share the whole link, which is why you can afford them, and which is exactly what makes latency unbounded. You bought throughput and price with a guarantee, and the missing guarantee is the one you now have to engineer around.
Fail-Slow Is Worse Than Fail-Stop
A node that crashes is the easy case. Its connections reset, its health check stops answering, the load balancer removes it, and the system heals.
A node that is merely slow stays in the pool. It answers health checks, because a health check is a cheap endpoint that does not touch the failing disk or the exhausted connection pool. It accepts its full share of traffic and takes 30 seconds to serve each request. Every caller that dispatches to it has a thread or a connection blocked for 30 seconds, so the slow node consumes caller capacity in proportion to how bad it is. One degraded replica out of twenty can exhaust the thread pools of every service that depends on it, and the dashboard for those services shows them failing while the actual cause is elsewhere and still reporting green.
This asymmetry is the reason bulkheads, deadlines, and outlier detection exist, and why a health check that does not exercise the dependency it is claiming to validate is close to worthless. The patterns for containing it belong to designing for failure; what matters here is understanding why the slow case is harder than the dead case rather than easier.
There Is No "Is It Alive" Primitive
You cannot ask the network whether a node is up. You can only send something and wait, which means every failure detector in existence is built on a timeout, and every timeout inherits the ambiguity above.
A few mechanisms narrow it. A TCP RST tells you a port has nothing listening, which distinguishes "process is gone" from "process is slow", though only if the machine is still up to send the reset. A hypervisor or a rack management controller can report that a VM is gone with authority the network cannot provide. A coordinated shutdown lets a node announce its own departure, turning an unplanned ambiguity into a planned handoff, which is why graceful drain is worth the engineering.
None of them remove the ambiguity in the general case. The realistic goal is not to know the truth about remote nodes. It is to build systems that behave correctly while remaining uncertain: give every remote call an explicit deadline, make every mutation safe to repeat, and decide in advance what the caller does with "I do not know" instead of discovering that it silently treats it as "no".