Distributed Systems Roadmap

Start at Remote Calls and Partial Failure and follow the stages in order: each one names what it needs first and what you should be able to do before moving on, from the machine boundary to running the system in production. Progress is stored locally in your browser.

Where to start

0 / 170 lessons masteredNot started 170Learning 0Practicing 0Mastered 0
  1. 1

    Remote Calls and Partial Failure

    Start here
    0/18

    What actually changes when a call crosses a machine boundary: no shared memory, no global clock, and a timeout that says nothing about whether the work happened. Then the failure models — crash, slow, Byzantine, correlated — and why a node can never be sure another one is down. Everything else in the domain depends on this shift, so it comes first.

    Before moving on: Take a timeout from a remote call, list the five things that could have happened at the callee, and explain why the caller can rule out none of them — including the three in which the work was done.

  2. 2

    Replication and Consistency

    0/18

    Copies buy availability and cost you agreement. Leader-based, multi-leader and leaderless replication each pay that cost differently, and quorums show where R + W > N holds and where it quietly does not. The consistency models — linearizable, causal, eventual, session guarantees — give you the vocabulary to name exactly what a replicated system promises, which is why CAP and PACELC come at the end of this stage rather than the start.

    Before moving on: State the guarantee of a replicated store with its scope (for example "linearizable for single-key operations") and explain which reads it can and cannot serve during a partition.

  3. 3

    Time, Ordering and Conflict

    0/14

    Physical clocks skew, so wall-clock timestamps cannot order events across machines; causality can. Lamport and vector clocks make that ordering explicit, and total-order broadcast shows what it costs to get a single order everyone agrees on. When two writes are concurrent and no order exists, someone must decide: last-write-wins, version vectors, application merges or CRDTs, each with a different idea of what convergence means.

    Before moving on: Read a spacetime diagram and say which events are causally ordered and which are concurrent, then pick a conflict-resolution rule for a given write pattern and name what it silently loses.

  4. 4

    Consensus and Coordination

    0/18

    Agreement under failure: what consensus solves, the system model it assumes, and how leader election, terms and the replicated log make Raft safe when a partitioned leader does not know it has been replaced. Fencing tokens, leases and stale lock holders are the same problem at the application level. Between the two sits the question the stage keeps returning to — do you need consensus at all — because the cheapest agreement is the one you did not need.

    Before moving on: Explain why a lock service without fencing tokens cannot protect an invariant, trace a Raft election through a partition, and say whether a given invariant needs consensus or can be protected without it.

  5. 5

    Messaging and Streams

    0/16

    Brokers, queues and logs, and the delivery guarantee each one actually provides. Acknowledgement, visibility timeouts, poison messages and dead-letter queues are the timeout-ambiguity problem applied to a queue; partitioned logs, consumer groups and offsets are ordering applied to a stream. Event time versus processing time, late events and watermarks close the stage, because a stream processor lives with two clocks at once.

    Before moving on: Say what a given broker guarantees about ordering and duplicates (and within what scope), and explain what a consumer must do so that a redelivered message does not change the result.

  6. 6

    Partitioning and Membership

    0/14

    Splitting data across nodes without splitting correctness: hash and range partitioning, ring hashing with virtual nodes, hot partitions, rebalancing and the cross-partition operations that quietly bring coordination back. Then membership: how nodes discover each other, how gossip and anti-entropy spread state, how Merkle trees find what differs, and why heartbeats can only produce suspicion, never certainty.

    Before moving on: Choose a partitioning scheme for a given key distribution and query shape, explain what happens to a hot partition and to a rebalance in flight, and say what a missed heartbeat does and does not prove.

  7. 7

    Workflows and Idempotency

    0/13

    Atomicity across services you do not control. Two-phase commit and how it blocks; sagas, orchestrated and choreographed, and why compensation is not rollback. The second half makes retries safe: retry ambiguity, idempotent operations and their scope, and the honest reading of delivery semantics — exactly-once is at-least-once plus deduplication, and the deduplication has a boundary.

    Before moving on: Design a multi-service write as a saga with a compensation for each step, and give every side-effecting call an idempotency key whose scope you can state.

  8. 8

    Overload, Deadlines and Caching

    0/18

    Bounded behaviour when demand exceeds capacity: backpressure, load shedding, admission control, and the retry amplification that turns a slow dependency into an outage unless budgets and jitter contain it. Deadlines propagate a time budget through a call graph, and hedged requests trade extra load for a shorter tail. Distributed caches are copies of copies, so they bring back every staleness, stampede and hot-key problem from earlier stages at a new layer.

    Before moving on: Set a retry budget and a propagated deadline for a three-hop call, explain which component sheds load first and why, and say what a cache stampede on a hot key looks like and how to bound it.

  9. 9

    Storage, Compute and Geography

    0/19

    The layers beneath a distributed database: file and object storage, durability, checkpoint-and-log recovery and consistent snapshots. Distributed compute moves work to data and pays for the shuffle, with stragglers setting the finish time. Multi-region designs then meet physics: the speed of light sets a floor on coordination, so active-passive versus active-active is a choice about which writes may conflict and where a region failure leaves you.

    Before moving on: Pick a multi-region write model for a given consistency requirement and explain its floor on write latency, its behaviour during a region failure, and what a data-residency rule does to the design.

  10. 10

    Operating Distributed Systems

    0/22

    How the system fails in production and how it recovers: detect, contain, recover, reconcile, verify. Cascading failure and dependency blast radius, chaos engineering and fault injection, debugging across logs that were never in one place. Then the boundaries — microservices as distribution, the distributed monolith, shared databases, data ownership and sources of truth — and finally agent workflows, which are distributed systems with every failure mode above and none of the tooling.

    Before moving on: Take an incident, say what each node could actually know at the time, name the containment step that limits the blast radius, and write the reconciliation that restores the invariant afterwards.