System Design (High-Level Design) Module 3 – Performance and reliability fundamentals
How distributed systems fail
Classify crash, omission, timing and Byzantine faults, spot partial, gray and metastable failures, and handle the unknown outcome of a call that timed out.
What you will learn
- Classify crash, omission, timing and Byzantine faults
- Recognise partial, gray and metastable failures from their symptoms
- Design for the failure modes that actually occur in production, starting with unknown outcomes
Before you start
On this page
On one computer, a program mostly either works or stops. A distributed system can do something worse: fail partially. Some machines work, some have crashed, some messages arrive and some are lost, and a call can end without any answer at all. That last case is the hardest to design for, because a call that timed out may have succeeded, failed, or still be running. This lesson names the kinds of faults that distributed systems are designed against (crash, omission, timing and Byzantine), shows how a timeout leaves an unknown outcome and how idempotency keys handle it, and describes two failure patterns that health checks and retries make worse rather than better: gray failures and metastable failures.
Partial failure, and the unknown outcome
Amazon’s Builders’ Library takes one request and reply between a client and a server and lists eight steps that can each go wrong independently: sending the request, delivering it, the server validating it, the server updating its state, the server sending the reply, delivering the reply, the client validating it, and the client updating its own state. Client, server and network fail independently of each other, and most of these failures look identical from the client’s side: nothing comes back.
One remote call, and the places it can fail
Text description of the diagram
The diagram follows one call from top to bottom.
- The client sends a request to charge order 7.
- The network carries it, and may lose or delay it.
- The server may crash before charging the card, or after charging it but before replying.
- The network carries the reply, and may lose or delay it.
- The client receives the reply or, if nothing arrives in time, a timeout.
A timeout looks the same whether the request was lost, the server crashed before acting, or the card was charged and the reply was lost. The client cannot tell which, so it has to handle an unknown outcome.
The article’s conclusion is the rule this lesson is built around: when a call times out, its result is unknown. The request may never have arrived, the server may have crashed before acting, or it may have acted and lost its reply. The client has to handle all three at once. This is not a rare corner: the same article warns that bugs in how a service handles these cases can stay hidden in production for months before a rare combination of failures sets them off.
Four fault models
Designs and protocols state which faults they tolerate, in the terms of a system model. Martin Kleppmann’s Cambridge lecture notes describe a system model in three parts, network, nodes and timing, and the four classic kinds of fault map onto them. From the easiest to tolerate to the hardest:
- Crash. A node stops, and may restart with only what it saved to disk. From outside it looks like silence, then a node that has forgotten its memory; a process killed for using too much memory is the everyday case. Defences: replicas, durable writes and restarts.
- Omission. A message is lost on the way in or on the way out, for example by a congested switch. From outside it looks like silence. Defences: timeouts, retries and deduplication.
- Timing. A correct answer arrives too late, after a garbage-collection pause or behind a slow disk: silence, then a late reply. Defences: timeouts, hedging, and protocols that stay correct when delays have no bound.
- Byzantine, or arbitrary. A node does anything at all: wrong answers, lies, different stories to different peers, because of corrupted memory, a compromised machine or a participant who cheats. Its answers look valid. Defences: signatures, checksums and protocols that tolerate a minority of liars.
The notes distinguish crash-stop nodes, which never come back, from crash-recovery nodes, which may restart having lost their memory but kept their disk; most real services are crash-recovery. On the network side, a “fair-loss” link may lose, duplicate or reorder messages, but a message retried often enough eventually gets through, which is why retries plus deduplication can turn it into a reliable link. For timing, the notes explain why real systems are only partially synchronous: networks and nodes are usually predictable, then occasionally pause for garbage collection, page faults, congestion or rerouting, and no timeout can tell a slow node from a dead one.
Byzantine faults are the general case, named after Lamport, Shostak and Pease’s paper on generals who must agree on a plan while some of them are traitors. Their result is that with ordinary messages, agreement is possible only if more than two thirds of the generals are loyal; with unforgeable signed messages it is possible for any number. Kleppmann notes that a bug shared by every node cannot be outvoted by such a protocol, because every node has it, so the term is usually kept for deliberate deviation. Services inside one organisation assume crash-recovery faults; blockchains and other systems shared between parties who do not trust each other must tolerate Byzantine ones.
A timeout does not tell you what happened
The program below pays for 10,000 orders through a fake payment server. The network loses 2% of requests before they reach the server and 3% of replies after the server has charged the card; the client sees the same timeout in both cases. Three clients face the same network luck:
"""A payment call that times out has an unknown outcome: the charge may or may not have happened.
A fake payment server charges a card when a request reaches it. The network loses 2% of requests on the way in
and 3% of replies on the way out; in both cases the client only sees a timeout. Three client strategies pay
for the same 10,000 orders with the same network luck. The rates are assumptions for the sketch.
"""
import random
class PaymentServer:
def __init__(self, deduplicate):
self.deduplicate = deduplicate
self.charges = [] # every charge made: the order it was for
self.done = {} # idempotency key -> charge id, when deduplicating
def charge(self, order, key):
if self.deduplicate and key in self.done:
return self.done[key] # a retry: the same answer, no new charge
self.charges.append(order)
charge_id = f"ch-{len(self.charges)}"
if self.deduplicate:
self.done[key] = charge_id
return charge_id
def pay(server, order, network, retries):
"""Try up to 1 + retries times; return True when a reply confirmed the payment."""
key = f"order-{order}" # the client picks one key per order and reuses it on every retry
for _ in range(1 + retries):
if network.random() < 0.02: # the request is lost: the server never sees it
continue
server.charge(order, key)
if network.random() < 0.03: # the charge happened, but the reply is lost
continue
return True
return False
print(f"{'strategy':28}{'charges':>8}{'charged 2+ times':>18}{'unconfirmed':>13}{'charged, but told no':>22}")
for name, retries, deduplicate in [
("no retry", 0, False),
("retry up to 3 times", 3, False),
("retry with idempotency key", 3, True),
]:
server, network = PaymentServer(deduplicate), random.Random(9)
confirmed = {order: pay(server, order, network, retries) for order in range(10_000)}
per_order = {}
for order in server.charges:
per_order[order] = per_order.get(order, 0) + 1
twice = sum(1 for n in per_order.values() if n > 1)
unconfirmed = sum(1 for ok in confirmed.values() if not ok)
charged_unconfirmed = sum(1 for order, ok in confirmed.items() if not ok and order in per_order)
print(f"{name:28}{len(server.charges):>8,}{twice:>18,}{unconfirmed:>13,}{charged_unconfirmed:>22,}") Output
strategy charges charged 2+ times unconfirmed charged, but told no no retry 9,802 0 503 305 retry up to 3 times 10,325 308 0 0 retry with idempotency key 10,000 0 0 0
Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 unknown_outcome.py
Runs on this device, in your browser. The first run downloads Python (about 13.5 MB), which is kept for the next runs.
Your run, in this browser
- No retry leaves 503 orders unconfirmed. Of those, 305 were in fact charged: their replies were lost. Those customers paid and were told the payment failed.
- Retry up to 3 times confirms every order, but charges 308 customers twice or more, because a retry after a lost reply makes a second charge.
- Retry with an idempotency key makes exactly 10,000 charges. The client picks one key per order and sends it with every attempt; the server records the result under the key, and a repeat gets the stored result instead of a new charge.
An idempotency key turns “do this” into “do this once”, which is what makes retries safe. When even the retries run out, the key still lets a later process look the payment up and settle it, which is called reconciliation: the client, or a background job, asks the server what happened to order 7 rather than guessing. The next lessons of this module build retries, timeouts and backoff on top of this.
UUID Generator Generate random UUIDs to use as idempotency keys, one per logical operation.Gray failures: healthy to the checker, broken for users
A gray failure is a component that is failing for the people who use it while the system’s own failure detectors see it as healthy. Huang and colleagues, writing from experience with large cloud systems, argue that this differential observability is the defining trait of gray failures, and that subtle faults of this kind, rather than clean crashes, cause many of the worst outages: a crash is detected and handled, a half-working component is not. Their example is a server whose request-handling module is stuck while its heartbeat module still runs, so every check that relies on heartbeats reports it as healthy.
// A gray failure: the health checks stay green while users fail. A checkout service has a shallow health check
// (is the process up?) and a deep one (can it reach its database?). From minute 4, the path to one of its three
// pricing shards breaks, so checkouts for customers on that shard fail. The workload is an assumption.
function generator(seed) {
let a = seed >>> 0;
return () => {
a = (a + 0x6d2b79f5) >>> 0;
let t = Math.imul(a ^ (a >>> 15), a | 1);
t ^= t + Math.imul(t ^ (t >>> 7), t | 61);
return ((t ^ (t >>> 14)) >>> 0) / 4294967296;
};
}
const random = generator(17);
const service = {
brokenShard: null, // which pricing shard is unreachable, if any
shallowHealth: () => "200 ok", // the process answers
deepHealth: () => "200 ok", // the database answers SELECT 1
checkout(customer) {
const shard = customer % 3; // customers are spread over three pricing shards
if (shard === this.brokenShard) return "503 pricing unavailable";
return random() < 0.001 ? "500 error" : "200 paid";
},
};
console.log("minute shallow check deep check checkouts succeeded alarm on success rate < 99%");
for (let minute = 1; minute <= 8; minute++) {
if (minute === 4) service.brokenShard = 2;
let ok = 0;
const requests = 600;
for (let i = 0; i < requests; i++) {
const customer = Math.floor(random() * 90000);
if (service.checkout(customer).startsWith("200")) ok += 1;
}
const rate = ok / requests;
const alarm = rate < 0.99 ? "FIRING" : "-";
console.log(
`${String(minute).padStart(6)} ${service.shallowHealth().padEnd(13)} ${service.deepHealth().padEnd(10)} ${String(requests).padStart(9)} ${(100 * rate).toFixed(1).padStart(8)}% ${alarm}`,
);
} Output
minute shallow check deep check checkouts succeeded alarm on success rate < 99%
1 200 ok 200 ok 600 100.0% -
2 200 ok 200 ok 600 100.0% -
3 200 ok 200 ok 600 99.7% -
4 200 ok 200 ok 600 67.3% FIRING
5 200 ok 200 ok 600 65.5% FIRING
6 200 ok 200 ok 600 68.8% FIRING
7 200 ok 200 ok 600 67.3% FIRING
8 200 ok 200 ok 600 69.5% FIRING
Recorded with Node.js 24.21.0 on macOS 26 arm64. To run it yourself: mise exec node@24.21.0 -- node gray_failure.mjs
Runs on this device, in your browser. The first run downloads JavaScript (about 0.6 MB), which is kept for the next runs.
Your run, in this browser
From minute 4, one of the service’s three pricing shards is unreachable, so about a third of checkouts fail. The shallow health check only asks whether the process answers, and the deep one only asks whether the database answers, so both stay at “200 ok” while the success rate drops to between 65% and 70%. Only the measurement of what users actually get, the share of successful checkouts, raises the alarm. The paper’s advice is to move from a single failure detector, such as a heartbeat, to several signals of health, and to measure from where the applications stand; in practice that means alerting on user-facing success rates, the SLIs of the previous lesson, and health checks that exercise the paths users depend on.
Metastable failures: when the fix keeps the failure going
Bronson, Aghayev, Charapko and Zhu describe a pattern behind widespread outages at large internet companies, some lasting hours: a metastable failure. A system runs in a vulnerable state, efficient and healthy, but with a hidden feedback loop. A trigger pushes it into a bad state, and a sustaining effect keeps it there even after the trigger is gone. Leaving the state takes a strong push, such as cutting the load sharply or changing the policy that feeds the loop.
Their first example is retries. A web application normally sends 280 queries a second to a database that answers quickly below 300 queries a second and much more slowly above it, and each user request retries a query that has not answered within a second. A 10-second network outage triggers a surge of queued requests and retries; the database slows, so queries time out and are retried, and the load stays at 560 queries a second, twice the real demand. The network has recovered, yet the system stays down until the load drops or the retry policy changes. The paper notes that outages like this are often blamed on the trigger, when the real root cause is the sustaining effect, and that the features creating the loop, retries, failover and caching, are usually ones added for efficiency or reliability.
The defences are the subject of the rest of this module: retries with backoff and a budget, so that retries cannot double the load; circuit breakers that stop calling a failing dependency; load shedding, so that an overloaded service refuses work quickly instead of slowly failing all of it; and capacity tests that look for the load at which the system becomes vulnerable.
Design for the failures that happen
- Put a timeout on every remote call, and treat a timeout as an unknown outcome, not a failure.
- Make every operation that changes state idempotent, with a key chosen by the caller, so retries are safe.
- Reconcile what timed out: look outcomes up by key rather than guessing.
- Alert on what users get (success rates and latency at the edge), not only on health checks.
- Make health checks exercise real dependencies, and keep the action they trigger proportionate: a check that removes every server when a shared dependency fails turns a partial failure into a total one.
- Give retries, failovers and other recovery actions a budget, so they cannot become the sustaining effect.
- Assume crash-recovery faults inside your organisation, and design for Byzantine faults only where participants do not trust each other.
Practice
Exercise · Medium · Python
A payment server that charges once per idempotency key
Complete PaymentServer.charge(key, amount) in payments.py, so that a client may safely retry a payment whose reply it never received.
- The first request with a new idempotency key makes a charge: it appends amount to self.charges and returns a new charge id, "ch-1" for the first charge the server ever makes, "ch-2" for the second, and so on. - A request that repeats a key with the same amount is a retry: it makes no new charge and returns the id of the original charge. - A request that repeats a key with a different amount is a client bug, not a retry: raise ValueError and make no charge.
For example, charge("order-7", 500) twice returns "ch-1" both times and leaves charges as [500].
Starter code · payments.py
class PaymentServer:
"""A payment server that charges at most once per idempotency key."""
def __init__(self):
self.charges = [] # every amount actually charged, in order
self.seen = {} # idempotency key -> (charge id, amount)
def charge(self, key, amount):
"""Charge `amount` once for `key` and return the charge id; a retry returns the same id."""
# Replace this line with your code.
return "" The sample tests · test_payments.py
from payments import PaymentServer
def test_first_charge():
"""a new key makes a charge and returns a new id"""
server = PaymentServer()
assert server.charge("order-7", 500) == "ch-1"
assert server.charge("order-8", 120) == "ch-2"
assert server.charges == [500, 120]
def test_retry_charges_once():
"""a repeated key with the same amount returns the first id and charges nothing"""
server = PaymentServer()
first = server.charge("order-7", 500)
assert server.charge("order-7", 500) == first
assert server.charge("order-7", 500) == first
assert server.charges == [500]
def test_mismatched_amount():
"""a repeated key with another amount is refused and charges nothing"""
server = PaymentServer()
server.charge("order-7", 500)
try:
server.charge("order-7", 900)
refused = False
except ValueError:
refused = True
assert refused, "a different amount for the same key must raise ValueError"
assert server.charges == [500]
def test_ids_count_real_charges():
"""ids count the charges made, so retries do not use up numbers"""
server = PaymentServer()
server.charge("a", 1)
server.charge("a", 1)
assert server.charge("b", 2) == "ch-2" A hint
Keep a dictionary from each key to what it produced, for example self.seen[key] = (charge_id, amount). Look the key up before charging: if it is there, compare the amounts and either return the stored id or raise; if it is not, make the charge, store the result and return the new id. The id can come from the number of charges made so far.
Results of the sample tests
| Test | Result | Details |
|---|
What your code printed
The sample tests run on this device, in your browser (Pyodide): nothing is sent to mysmartcopilot.com. The first run downloads Python (about 13.5 MB), which is kept for the next runs. A check in your browser is feedback for you, not proof that the code is right for every input.
Check yourself
8 questions about this lesson. Every answer and why it is right is on the page, behind “Show the answer”. Your score stays in this browser.
Interview questions
Warm-up (fresher to mid level): what is a partial failure, and why does it make distributed systems harder than one machine? A partial failure is one in which some parts of a system fail while others keep working: one node crashes, one network link drops messages, one dependency slows down. On a single machine a failure usually stops the whole program, which is easy to notice. In a distributed system the working parts must carry on around the broken ones, and a caller often cannot tell which happened: a timed-out call may have failed, succeeded with a lost reply, or still be running. So every remote call needs a timeout, and every operation that changes state needs to be safe to retry.
Key takeaways
- A remote call can fail at eight points between client and server; a timeout means its outcome is unknown.
- Crash, omission, timing and Byzantine faults are increasingly hard to tolerate; internal services assume crash-recovery, multi-party systems Byzantine.
- Idempotency keys make retries safe: without them, retries after lost replies charge customers twice; without retries, customers are charged and told it failed.
- Gray failures pass health checks while users fail: alert on user-facing success rates.
- Metastable failures persist after their trigger because a feedback loop, often retries, sustains them; budget every recovery action.
References
- Challenges with distributed systems (Amazon Builders' Library) (Amazon Web Services)
- Distributed Systems, lecture notes: system models (Martin Kleppmann) (University of Cambridge, Department of Computer Science and Technology)
- The Byzantine Generals Problem (Lamport, Shostak and Pease) (ACM Transactions on Programming Languages and Systems (via Microsoft Research))
- Gray Failure: The Achilles' Heel of Cloud-Scale Systems (HotOS 2017 (Microsoft Research))
- Metastable Failures in Distributed Systems (HotOS 2021 (ACM SIGOPS))
Related tools
Report a problem with this lesson
Kept only in this browser. Your Learn progress