14 parts · 22 chapters

Distributed Systems

A distributed system is one in which the failure of a computer you did not know existed can make your own computer unusable. Leslie Lamport's line is still the best definition, because the defining property is partial failure: some parts work, some do not, and you often cannot tell which. Every hard problem in this course, from ordering events to agreeing on a leader to moving money between two services exactly once, comes from that one fact plus a network that loses, delays, duplicates and reorders messages.

Fourteen parts, in the order the ideas build on each other. Why distributed is different; time and ordering without a shared clock; consensus and Raft step by step; replication in its three shapes; partitioning; the consistency models and what each costs; failure detection and idempotency; logs, queues and streams; transactions across services with sagas and the outbox; and the papers that proved what is possible, read for what each one changed. Four further parts follow: CRDTs and client sync, caching across the system, membership with gossip and consistent hashing, and testing distributed systems.

partial failure · time and order · consensus · replication · partitioning · consistency · idempotency · logs · transactions · the papers · CRDTs and sync · caching · gossip and hashing · testingsenior → staff · anyone designing or debugging systems with more than one machine
partial failureThe fallacies of distributed computing, and why "did it work?" can have no answer.
time and orderPhysical clocks and their drift, Lamport clocks, vector clocks, causality, hybrid logical clocks.
consensusPaxos and Raft, leader election, log replication, what consensus buys and costs.
replicationLeader-follower, multi-leader and leaderless replication; lag, conflicts and quorums.
partitioningKey and hash partitioning, rebalancing, hot spots, secondary indexes.
consistencyLinearizable, sequential, causal, eventual, and the CAP and PACELC trade-offs.
idempotency and logsHeartbeats, timeouts, exactly-once in effect; Kafka's model and stream processing.
transactions and papersSagas, the outbox, 2PC; and Dynamo, Spanner, Raft, Calvin, Chubby, Kafka, ZooKeeper.
The theory beneath Cloud, Infra and SREManaged databases, queues, Kubernetes' etcd, multi-region failover and every payment flow in the Trust course are applications of this course. Read it after Cloud Engineering, or alongside Designing Data-Intensive Applications.