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
Distributed Systems
10 stages · 0/170 lessonsCorrectness, availability and recovery when computation and data are spread across machines that talk over an unreliable network.
- 10/18
Remote Calls and Partial Failure
Start hereWhat 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
timeoutfrom 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.- What Actually Makes a System Distributed
- A Remote Call Is Not a Function Call
- The Network Changes Everything
- No Shared Memory: Every Node Sees a Copy
- There Is No Global Clock
- Why Distribute At All
- When Not to Distribute
- Where the Boundary Goes
- Name the Invariant Before You Choose the Protocol
- Partial Failure: The Founding Condition
- A Timeout Tells You Nothing About Whether It Happened
- Failure Models: What You Are Allowed to Assume
- No Heartbeat Does Not Mean Dead
- Crashed or Just Slow: The Distinction You Cannot Make
- Byzantine Failures, and Why You Probably Do Not Assume Them
- Fault Domains: What Fails Together
- Correlated Failure: The Independence Assumption Is Usually False
- Redundancy Is Not Resilience
- 20/18
Replication and Consistency
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 > Nholds 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.
Needs first:Remote Calls and Partial Failure- Why Replicate: What a Second Copy Buys You
- Leader-Based Replication: Buying Order With a Single Writer
- Synchronous Replication: Paying Latency for a Durability Guarantee
- Asynchronous Replication: The Loss Window You Chose
- Multi-Leader Replication: Accepting Writes in More Than One Place
- Read-After-Write: Letting a User See Their Own Change
- Monotonic Reads: Never Let Time Run Backwards
- Quorums: What R + W > N Does and Does Not Buy
- Leaderless Replication: Every Replica Accepts Writes
- Consistency Models: What "Consistent" Has To Mean
- Linearizability: An Operation Is an Interval, Not a Point
- Serializability vs Linearizability: Two Different Properties
- Eventual Consistency: If Updates Stop, Replicas Converge
- Causal Consistency: Never Show an Effect Before Its Cause
- Session Guarantees: The Underrated Middle Ground
- CAP: What the Theorem Actually Says
- PACELC: The Trade-off That Exists When Nothing Is Broken
- Choosing a Consistency Model: Start From the Invariant
- 30/14
Time, Ordering and Conflict
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.
- Two Timestamps Are Not an Ordering
- Clock Skew: The Gap You Cannot Measure From Inside
- Never Measure a Duration With the Wall Clock
- Happens-Before: The Only Ordering You Actually Have
- Lamport Clocks: Consistent With Causality, Blind to Concurrency
- Vector Clocks: Buying Concurrency Detection at O(N)
- Four Orderings, Four Prices
- Total Order Broadcast Is Consensus Wearing a Different Hat
- Two Writes, No Order, One Answer Required
- Last Write Wins Is Data Loss You Chose by Default
- Version Vectors: Making the Conflict Visible
- Only the Application Knows What the Merge Means
- CRDTs: Deterministic Merge, Not Correct Merge
- What "Eventually Converges" Actually Requires
- 40/18
Consensus and Coordination
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.
- What Consensus Actually Solves
- Consensus Is Not Magic: The Assumptions It Runs On
- Leader Election: Choosing One, and Knowing You Were Chosen
- Terms and Epochs: Making Stale Leaders Harmless
- Split-Brain: Two Nodes, Both Certain They Are In Charge
- Fencing Tokens: Making the Stale Actor Safe, Not Just Unlikely
- Raft: Elections, Terms and Three States
- The Raft Log: Commit Index, Divergence and Reconciliation
- Paxos and the Other Protocols: What They Share and Where They Differ
- Do You Actually Need Consensus?
- Coordination Couples Availability
- Coordination Avoidance: Restructuring the Problem Instead of Paying for It
- Start From the Invariant, Not From the Architecture
- Distributed Locks: What They Are Actually For
- Leases: Authority With an Expiry Date
- The Stale Lock Holder: A Paused Process Does Not Know It Was Paused
- Coordination Services: The Primitives, Not the Product
- Distributed Uniqueness: One Name, Many Shards
- 50/16
Messaging and Streams
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.
- What a Broker Actually Buys You
- Work Queues: One Task, One Worker, Competing Consumers
- Publish/Subscribe: One Event, Many Independent Readers
- Queue or Pub/Sub: Answer the Question in One Sentence
- Acknowledgement: The Two-Line Protocol That Decides Your Delivery Semantics
- Visibility Timeout: The Message Is Hidden, Not Yours
- Poison Messages: The One That Fails Every Time, Forever
- A Dead-Letter Queue Is a Workflow, Not a Bin
- The Log Is Not a Queue
- A Topic Is Not One Log: Ordering Lives Inside a Partition
- Consumer Groups: Queue Semantics Inside, Pub/Sub Semantics Across
- Rebalancing: Everyone Stops So the Partitions Can Move
- Commit Before or After: There Is No Third Option
- Two Clocks: When It Happened and When You Saw It
- Late Events: The Window Already Fired
- Watermarks: A Guess About Time, Made Precise Enough to Act On
- 60/14
Partitioning and Membership
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.
- Why Partition: Four Ceilings, Four Different Answers
- Hash Partitioning and the Modulo Trap
- Range Partitioning: Scans You Keep, Hotspots You Inherit
- The Ring: Keeping the Mapping Stable When Membership Changes
- Virtual Nodes: Many Positions per Machine, and Why It Is Not Optional
- Hot Partitions: The Skew Hashing Cannot Fix
- Rebalancing: A Load Spike You Schedule for Yourself
- Cross-Partition Operations: Paying for What the Split Took Away
- Discovering Services: The Registry Is a Distributed System Too
- Cluster Membership: A Belief, Not a Fact
- Gossip: Epidemic Spread Instead of Everyone Telling Everyone
- Anti-Entropy: Repairing Divergence Nobody Reported
- Merkle Trees: Finding the Difference Without Reading the Data
- From Alive-or-Dead to a Suspicion Level
- 70/13
Workflows and Idempotency
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.
- Atomicity Stops at the Process Boundary
- Two-Phase Commit: Buying Atomicity With a Promise
- The Blocking Window: When 2PC Stops and Waits
- Sagas: Trading Isolation for Availability
- A Refund Is Not a Rollback
- Orchestration: One Component Owns the Workflow
- Choreography: The Workflow Nobody Wrote Down
- The Retry Is a Decision, Not a Reflex
- Idempotent Is a Property of the Whole Effect, Not the Write
- What Counts as the Same Operation?
- Where You Put the Acknowledgement Decides Everything
- Exactly-Once Is a Scope, Not a Guarantee
- Deduplication: Bounded Memory Against an Unbounded Stream
- 80/18
Overload, Deadlines and Caching
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.
- Backpressure Is a Signal That Has to Travel — and Reach Someone Who Can Slow Down
- Rejecting Work on Purpose — and Rejecting It Cheaply Enough to Help
- Decide at the Door Whether the Capacity Exists
- One Retry per Tier Is Not One Retry — It Multiplies
- Cap Retries as a Fraction of Traffic, Not as a Count per Request
- Without Jitter, Every Client That Failed Together Retries Together
- Containment Is Decided by What Is Shared, Not by Where the Service Boundaries Are
- Bulkheads: Buying Independence by Giving Up Utilisation
- A Deadline Is Divided Across the Call Chain, Not Repeated at Every Hop
- Pass the Remaining Budget Down, Not a Fresh One
- The Caller Is Gone — Stopping Is Usually Right and Sometimes Unsafe
- Send a Second Request After p95 and Take Whichever Answers First
- Fan Out to 100 and the Component’s Tail Becomes the System’s Median
- A Cache Across Machines Is a Replica With No Replication Protocol
- Invalidation Is a Messaging Problem, Which Is Why Cache Bugs Are Hard
- One Key Expires and Five Hundred Instances Miss at the Same Millisecond
- You Cannot Enumerate the Caches, So TTL Is the Bound and Invalidation Is the Optimisation
- Sharding Does Not Help a Single Key
- 90/19
Storage, Compute and Geography
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.
- A Distributed Database Is a Stack, Not a Box
- Distributed File Systems: Chunks, a Metadata Service, and Where the Copies Go
- Object Storage: A Flat Namespace With a Hash Behind It
- Acknowledged, Durable, Replicated: Three Different Things
- Recovered State Is a Checkpoint Plus the Log After It
- A Consistent Cut, Without Stopping the World
- Splitting a Computation Across Machines
- MapReduce: The Model That Made the Trade-offs Visible
- The Shuffle Is the Job
- Move the Computation to the Data
- Who Runs What, and What Happens When a Worker Goes Quiet
- One Slow Task Sets the Pace for Everything
- A Region Boundary Is a Consistency Decision
- The One Number You Cannot Optimise
- Three Ways to Accept a Write in More Than One Place
- Active-Passive: Simple to Reason About, Rarely Tested
- Active-Active: Every Conflict Scenario Becomes Real
- EU and US Are Partitioned. Can Both Keep Accepting Writes?
- When the Data Is Not Allowed to Leave
- 100/22
Operating Distributed Systems
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.
- Detect, Contain, Recover, Reconcile, Verify
- Graceful Degradation: Which Dependency Is Actually Critical
- The Steady-State Hypothesis and the Abort Condition
- Cascading Failure: When the Response to Failure Causes More Failure
- Dependency Blast Radius: What Dies If This Node Dies
- Chaos Engineering Is Not Randomly Breaking Production
- Fault Injection: The Catalogue, and Which Faults Are Hard
- Distributed Debugging: The Question Ladder
- Three Nodes, Three Logs, and You Cannot Sort by Timestamp
- Microservices Are a Distribution Decision, Not a Scaling Technique
- The Distributed Monolith: All of the Cost, None of the Autonomy
- Four Questions That Test a Proposed Boundary
- The Shared Database: An Honest Trade, Not a Prohibition
- Exactly One Component Owns Each Piece of State
- Source of Truth: The Question Every Inconsistency Incident Is Really Asking
- Materialized Views: A Read Model That Lags
- Reconciliation Is a Component, Not a Cleanup Script
- An Agent System Is a Distributed System
- The Model Retries Because It Cannot See the Result
- Agents Do Not Negotiate. Processes Contend for State.
- Resuming a Workflow That Died Halfway Through
- The Failures That Produce No Errors