Distributed Systems / Reliability / Backend Engineering
Distributed Systems Are About Partial Failure
Every remote call has three outcomes, not two. Idempotency, retry budgets, deadlines, quorums, and fencing tokens are what you build once you accept that.
On this page
- A remote call has three outcomes
- Make retries safe before making them automatic
- Retry with a budget, not a loop
- Timeouts are a design decision
- Overload is a failure mode, not a capacity number
- Replication, and what "consistent" means
- Do not coordinate with clocks
- Observe the request, not the machine
- A checklist for every remote call
- Further reading
Leslie Lamport once described a distributed system as one in which the failure of a computer you didn't even know existed can render your own computer unusable. The line is a joke, but it names the property that separates distributed systems from everything else: parts of the system fail while other parts keep running, and the running parts often cannot tell.
A single process either works or crashes. A system of processes connected by a network can be half-working in ways that are hard to observe: a replica that accepts writes but cannot replicate them, a dependency that answers health checks but times out on real traffic, a request that was executed but never acknowledged. Most of the techniques that make distributed systems dependable are responses to that one fact.
A remote call has three outcomes
A local function call returns or throws. A remote call can succeed, fail, or end with no information at all. When a request times out, the caller does not know which of these happened:
- the request never reached the server;
- the server received it and failed before doing the work;
- the server did the work, and the response was lost or arrived after the caller gave up.
The third case is the one that causes incidents. If the operation was "charge this card" or "send this email", retrying blindly repeats the side effect, and not retrying may drop a request that never ran.
Treat "unknown" as a first-class result in code. A timeout should not be logged as "payment failed"; it should be recorded as "payment outcome unknown" and resolved by a mechanism that can find out, either an idempotent retry or a status lookup.
Make retries safe before making them automatic
An operation is idempotent when performing it twice has the same effect as performing it once. Reads usually are. Writes usually are not, unless you design them to be.
The common design is an idempotency key: the client generates a unique key per logical operation and sends it with every attempt. The server stores the key with the outcome and, on a repeat, returns the stored outcome instead of executing again.
type StoredResult = { status: number; body: unknown };
async function handleCharge(req: ChargeRequest): Promise<StoredResult> {
const key = req.headers['idempotency-key'];
if (!key) return { status: 400, body: { error: 'idempotency key required' } };
return db.transaction(async tx => {
// Claim the key. The unique constraint on (client_id, key) means a
// concurrent duplicate blocks here and then sees the committed row.
const existing = await tx.idempotency.findForUpdate(req.clientId, key);
if (existing) {
if (existing.requestHash !== hash(req.body)) {
return { status: 422, body: { error: 'key reused with a different request' } };
}
return existing.result; // replay the original outcome
}
const charge = await tx.charges.insert(toCharge(req.body));
const result = { status: 201, body: charge };
await tx.idempotency.insert({
clientId: req.clientId,
key,
requestHash: hash(req.body),
result,
});
return result;
});
}A few details matter more than the happy path:
- Scope the key to the client or tenant, so two clients cannot collide.
- Store a hash of the request. Reusing a key with a different payload is a client bug and should be rejected, not silently answered with an unrelated result.
- Put the key and the side effect in the same transaction when they share a database. If the side effect is external, such as a call to a payment provider, pass the key through so the provider can deduplicate too, and record an in-progress state so concurrent retries wait instead of racing.
- Decide on retention. Keys need to outlive the client's retry window. They do not need to live forever.
Retry with a budget, not a loop
Once operations are safe to repeat, retries become a tool for absorbing transient failures. They are also the most common way to turn a small outage into a large one.
Three rules keep retries useful:
- Back off exponentially and add jitter. Synchronized retries from many clients arrive as waves. Randomizing the delay spreads them out.
- Respect a deadline. The caller has a time budget for the whole operation, and retries spend from it. A retry that cannot finish before the deadline should not start.
- Retry at one layer. If the client, the gateway, and the service each retry three times, one failing call becomes 27 attempts against the dependency that is already struggling.
async function withRetry<T>(
attempt: (signal: AbortSignal) => Promise<T>,
{ deadline, baseMs = 100, maxMs = 2_000, isRetryable }: RetryOptions
): Promise<T> {
for (let n = 0; ; n += 1) {
const remaining = deadline - Date.now();
if (remaining <= 0) throw new DeadlineExceeded();
try {
return await attempt(AbortSignal.timeout(remaining));
} catch (error) {
if (!isRetryable(error)) throw error;
// "Full jitter": a random delay between 0 and the exponential cap.
const cap = Math.min(maxMs, baseMs * 2 ** n);
const delay = Math.random() * cap;
if (Date.now() + delay >= deadline) throw error;
await sleep(delay);
}
}
}A retry budget adds a fleet-level limit: allow retries only while they stay below some fraction of total requests to a dependency. When the dependency is healthy, nearly every retry fits in the budget. When it is failing broadly, the budget runs out and callers fail fast instead of multiplying load.
isRetryable deserves real thought. Connection resets and explicit "try again" responses usually are retryable. Validation errors never are. A timeout on a non-idempotent call is not retryable unless an idempotency key makes it so.
Timeouts are a design decision
Every remote call needs a timeout, and the library default is almost never the right one. A useful timeout comes from the dependency's observed latency distribution and from what the caller can afford: long enough that healthy-but-slow requests succeed, short enough that a stuck dependency does not hold threads, connections, and memory hostage.
Prefer deadlines over per-hop timeouts. A deadline is an absolute time by which the whole request must finish, propagated with the request to every downstream call. Each service compares the remaining time with how long its work usually takes and refuses work it cannot complete. Without propagation, a downstream service keeps working on a request the user abandoned seconds ago.
Overload is a failure mode, not a capacity number
Systems rarely fail at exactly their capacity limit. They fail when demand rises past it and the system's own reactions make things worse: queues grow, latency grows, clients time out and retry, and the extra retries add more load. The system can stay in that degraded state even after the original trigger has gone away. Bronson and colleagues called these metastable failures, and the defining feature is that the cause and the sustaining effect are different things.
The defenses are about refusing work early:
- Bound every queue. An unbounded queue converts overload into latency and memory exhaustion instead of fast, visible errors.
- Shed load at the edge. Reject excess requests before they consume downstream resources, and reject cheaply.
- Prioritize. When you must drop work, drop retries and background jobs before interactive requests.
- Return backpressure. A clear "overloaded, retry after N seconds" response lets well-behaved clients slow down.
Replication, and what "consistent" means
Replication keeps data available when a machine fails and puts it closer to readers. It also means there is more than one copy, and the copies are not always identical.
In a leader-based system, writes go to one leader and replicate to followers. Reading from the leader gives you the latest value. Reading from a follower may give you an older one, depending on replication lag. In a leaderless system, a client writes to several replicas and reads from several, and quorums decide when an operation counts as complete.
The quorum condition is simple arithmetic. With N replicas, a write that waits for W acknowledgements and a read that queries R replicas are guaranteed to overlap if W + R > N. Overlap is not the whole story. Concurrent writes, failed partial writes, and the rules for picking the newest version still decide what the reader actually sees.
It helps to choose consistency per operation rather than per database:
| Guarantee | What the user observes | Typical need | | --- | --- | --- | | Linearizable | Behaves like a single copy; reads see the latest completed write | Uniqueness checks, locks, balances | | Read-your-writes | A user always sees their own updates | Profile edits, settings, comments | | Monotonic reads | Never sees data go back in time across requests | Feeds, dashboards, timelines | | Eventual | Replicas converge if writes stop | Counters, analytics, caches |
Stronger guarantees cost latency and availability during network partitions. Weaker guarantees cost application complexity, because the code has to tolerate stale or out-of-order data. Neither choice is free, which is why it is worth making deliberately for each feature.
Do not coordinate with clocks
Wall clocks on different machines disagree, and NTP corrections can make a clock jump backward. Ordering events across machines by timestamp is therefore unreliable, and conflict resolution based on "last write wins" can silently discard a write that was actually later.
When order matters, use mechanisms designed for it. A single leader can assign sequence numbers. Logical clocks, introduced in Lamport's 1978 paper on time and ordering, capture causality without trusting wall time.
The same caution applies to locks. A process that holds a lease can pause, for a garbage-collection stop or a slow disk, long enough for the lease to expire and another process to acquire it. When the first process resumes, it still believes it holds the lock. The standard defense is a fencing token: the lock service hands out a number that increases with every grant, and the protected resource rejects any request carrying a number lower than one it has already seen.
// Storage side: reject writes from a holder whose lease was superseded.
async function writeWithFence(file: string, data: Buffer, token: number) {
const current = await meta.get(file);
if (current && token < current.highestToken) {
throw new StaleLeaseError(file, token, current.highestToken);
}
await meta.set(file, { highestToken: token });
await blobs.put(file, data);
}The check has to happen at the resource, not in the client. A client cannot reliably know that it has been paused.
Observe the request, not the machine
Per-host metrics tell you that something is wrong. They rarely tell you where a specific request went. Propagate a trace identifier with every call, record each attempt as its own span, including retries and the reason for each, and log which replica or partition served it. When a user reports that their update disappeared, the trace should show whether the write was acknowledged, by which quorum, and which replica served the read that missed it.
A checklist for every remote call
Most distributed-systems bugs come from a small number of unanswered questions. Asking them for each new dependency catches a surprising share of them:
| Question | Why it matters | | --- | --- | | What happens if this call times out after the work was done? | Decides whether you need idempotency or a status lookup. | | Is it safe to retry, and who retries? | Prevents duplicate side effects and retry amplification. | | What is the deadline, and is it propagated? | Stops work on abandoned requests. | | What does the caller do when the dependency is down? | Forces an explicit degraded mode instead of an outage. | | Which consistency does this read need? | Avoids paying for strong reads that do not need it, and missing ones that do. | | How would I find this request in the logs? | Makes the failure diagnosable after the fact. |
None of these require exotic infrastructure. They require accepting that the network is part of the program, and that "it didn't answer" is a result the program has to handle.
For a closer look at the communication patterns these calls run over, see How Systems Communicate.
Further reading
- Martin Kleppmann, Designing Data-Intensive Applications: replication, consistency models, and the fencing-token argument in depth.
- Leslie Lamport, "Time, Clocks, and the Ordering of Events in a Distributed System" (1978).
- Nathan Bronson et al., "Metastable Failures in Distributed Systems" (HotOS 2021).
- The "fallacies of distributed computing", originally collected at Sun Microsystems, still a good list of assumptions to check.