Reference · Distributed Systems / Failure Models

Distributed Systems / Failure Models — Cheat Sheet

Pairs with the distributed-systems lessons: Lesson 7: Retries and Partial Failure through Lesson 21: The Traffic Path.

Glossary

TermMeaning
Partial failureA network call fails ambiguously — the caller can't tell whether the request was never received, was received and is still processing, or was processed and only the response was lost.
TimeoutA deliberate, bounded time limit a caller places on waiting for a response, after which it treats the call as failed even though the true outcome downstream is still unknown.
Retry storm / retry amplificationMany callers retrying a failing or slow dependency add load on top of the load that caused the failure in the first place, deepening the outage rather than riding it out.
Exponential backoffEach successive retry waits longer than the last (e.g. doubling), instead of retrying at a fixed interval.
JitterRandomizing the exact wait time within a range so retries from many callers don't stay synchronized and re-collide on the same instant.
Circuit breakerA guard around a dependency with three states — closed (calls go through normally), open (calls fail immediately, no network attempt, after failures cross a threshold), half-open (after a cooldown, a limited number of calls are let through to test whether the dependency has recovered).
CAP theoremUnder a network partition, an operation that needs agreement across nodes must choose between answering with possibly-stale data (favoring availability) or refusing to answer until agreement is possible (favoring consistency) — one line, not the full model.

Glossary II — messages, logs, and atomicity

TermMeaning
Dual writeWriting to two systems (a database and a broker) with no shared transaction — the crash window between them loses the event or creates a phantom one. The bug the outbox pattern exists to fix.
Outbox (transactional outbox)Writing the event as a row in the same database transaction as the business write, then having a relay publish it. Makes the DB transaction the atomicity boundary for both systems' writes.
RelayThe process that reads unpublished outbox rows, publishes them to the broker, and marks them done. Crash between publish and mark → duplicate publish → consumers must be idempotent.
SagaMulti-service business process as retryable, idempotent local transactions, each with a compensating transaction that undoes it if a later step fails. Eventual consistency by running the undo, instead of atomicity by not committing.
Compensating transactionThe undo of a saga step (refund, restock, cancel). It is itself a distributed call: it can fail and must be retried with backoff until it succeeds.
LogAn append-only sequence of records with monotonically increasing offsets. Read by consumers at their own pace; nothing is deleted on delivery — replay is possible.
OffsetA consumer's position in a log. Committing it after processing gives at-least-once; before processing, at-most-once. The storage location of the "already did it" flag.
PartitionOne shard of a log: an independent log with its own offsets. Order is preserved within a partition, not across partitions — partition by key to keep per-entity order.
Consumer groupConsumers sharing a log's partitions; each partition is read by exactly one member at a time, so scaling readers preserves per-partition order.
At-least-once / at-most-once / effectively-onceDelivery guarantees: at-least-once loses nothing but may duplicate (needs idempotent consumers); at-most-once duplicates nothing but may lose; effectively-once = at-least-once + dedup.
Consumer lagThe gap between the latest offset and a consumer's position — Little's law in disguise, an unbounded backlog converting a throughput gap into delayed failure.

Glossary III — partitioning and scaling

TermMeaning
Partitioning / shardingSplitting one logical dataset into subsets, each on a different node. Replication makes copies (fault tolerance); partitioning makes subsets (capacity). They combine: each partition is replicated.
Range partitioningKey space split into ranges (by id, name, time). Good for range scans; hotspots on the active range (recent timestamps, new signups).
Hash partitioninghash(key) decides the node. Even spread by construction; kills range scans; hot keys still saturate their one node.
RebalancingMoving partitions when nodes join or leave. The metric that matters: data moved per node added — should be ~1/N, not everything.
Consistent hashingKeys and nodes on a ring; a key belongs to the next node clockwise. Adding a node moves only the arc between it and its predecessor (~1/N of keys). Virtual nodes smooth the distribution.
HotspotOne key or range getting disproportionate load — a load problem, not a size problem; no shard count fixes it (cache or shed instead).
Fan-out / scatter-gatherA query no single shard can answer hits all shards and merges. Latency = slowest shard (tail at scale).
Cross-shard transactionA transaction touching two shards; needs 2PC or sagas. A sign the partition key doesn't match the access pattern.

Glossary IV — agreement and consistency

TermMeaning
QuorumA subset of nodes large enough to decide: write quorum W, read quorum R, with R + W > N so they intersect. Majority = N/2 + 1 is the special case.
ConsensusGetting nodes to agree on a value (who's the leader, which writes are committed) despite failures and partitions. Raft is the protocol to know.
Leader electionTerms + randomized timeouts + majority of votes. At most one leader per term, because two majorities would have to overlap.
Commit on majorityA write is durable only when a majority acknowledges it — the Raft rule that makes the durability window (Lesson 15) disappear.
Split brainTwo sides both accepting writes after a partition. Unrecoverable divergence; consensus prevents it by refusing to commit on the minority side.
LinearizabilityEvery operation takes effect atomically at one instant, respecting real time: a completed write is visible to subsequent reads. What single-node DBs give by default.
Sequential consistencySome global order exists in which each process's ops appear in program order, but real-time ordering of concurrent ops is dropped — weaker and easier to implement.

Decision rule

The default policy for any network call

Every network call needs a bounded timeout and a backoff-plus-jitter retry policy. Every dependency you retry against eventually needs a circuit breaker, so a struggling dependency gets protected instead of buried.

Atomicity across systems

One database + one event that must be atomic → outbox (event row in the same transaction, relay publishes). A multi-service process that must converge → saga (idempotent steps + compensations). Both assume at-least-once delivery and require idempotent consumers; both are eventually consistent.

Agreement, in one line

R + W > N gives you intersection; waiting for full quorums gives you linearizable reads; below a majority of reachable nodes, a quorum system refuses writes — the availability cliff that prevents split brain.

Primary sources