Pairs with the distributed-systems lessons: Lesson 7: Retries and Partial Failure through Lesson 21: The Traffic Path.
| Term | Meaning |
|---|---|
| Partial failure | A 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. |
| Timeout | A 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 amplification | Many 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 backoff | Each successive retry waits longer than the last (e.g. doubling), instead of retrying at a fixed interval. |
| Jitter | Randomizing the exact wait time within a range so retries from many callers don't stay synchronized and re-collide on the same instant. |
| Circuit breaker | A 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 theorem | Under 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. |
| Term | Meaning |
|---|---|
| Dual write | Writing 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. |
| Relay | The 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. |
| Saga | Multi-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 transaction | The 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. |
| Log | An append-only sequence of records with monotonically increasing offsets. Read by consumers at their own pace; nothing is deleted on delivery — replay is possible. |
| Offset | A 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. |
| Partition | One 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 group | Consumers 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-once | Delivery 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 lag | The 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. |
| Term | Meaning |
|---|---|
| Partitioning / sharding | Splitting 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 partitioning | Key space split into ranges (by id, name, time). Good for range scans; hotspots on the active range (recent timestamps, new signups). |
| Hash partitioning | hash(key) decides the node. Even spread by construction; kills range scans; hot keys still saturate their one node. |
| Rebalancing | Moving partitions when nodes join or leave. The metric that matters: data moved per node added — should be ~1/N, not everything. |
| Consistent hashing | Keys 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. |
| Hotspot | One 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-gather | A query no single shard can answer hits all shards and merges. Latency = slowest shard (tail at scale). |
| Cross-shard transaction | A transaction touching two shards; needs 2PC or sagas. A sign the partition key doesn't match the access pattern. |
| Term | Meaning |
|---|---|
| Quorum | A 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. |
| Consensus | Getting 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 election | Terms + randomized timeouts + majority of votes. At most one leader per term, because two majorities would have to overlap. |
| Commit on majority | A write is durable only when a majority acknowledges it — the Raft rule that makes the durability window (Lesson 15) disappear. |
| Split brain | Two sides both accepting writes after a partition. Unrecoverable divergence; consensus prevents it by refusing to commit on the minority side. |
| Linearizability | Every 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 consistency | Some 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. |
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.
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.
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.