Advanced Topics & Final Prep
Consistent hashing, microservices and service discovery, consensus and leases, sagas and the outbox, observability, and a final-prep checklist
SPACED REPETITION Β· 16 practice questions
Make this lesson stick.
Try 3 questions now. No account needed. Sample answers aren't saved.
or sign in to practice all 16Many machines, one system
Your checkout used to be one database transaction:
BEGIN;
INSERT INTO orders ...;
UPDATE stock SET available = available - 1 WHERE sku = ...;
INSERT INTO payments ...;
COMMIT;
Either all three rows changed or none did. Then the company grew. Orders, Stock and Payments became separate services with separate databases, and the card charge moved to an external provider. The same checkout is now several network calls and several commits.
Suppose you run 300 checkouts per second at peak, and one checkout in a thousand dies between two of those steps: a deploy, a timeout, a crashed pod. That is 0.3 per second, or about 1,080 checkouts an hour with stock held but no payment, or money taken but no order.
No single line of that design is wrong. It just lost a guarantee that one machine used to give you for free. Most "advanced" system design topics are the replacements you build when one box becomes many, and this lesson teaches six of them:
| # | One machine gave you⦠| Many machines need⦠| Move |
|---|---|---|---|
| 1 | Every key in one place | A placement rule that survives nodes joining and leaving | Consistent hashing |
| 2 | Function calls between modules | Calls that can fail, and a way to find the callee | Service boundaries and discovery |
| 3 | One process in charge | Exactly one leader, even when the network misbehaves | Consensus, leases and fencing |
| 4 | One ACID transaction | All steps done, or all compensated, across services | Sagas |
| 5 | Saving and notifying in one place | A commit and its event that can never disagree | Outbox, CDC, event sourcing, CQRS |
| 6 | One log file and a debugger | A view of one request across many services | Observability |
The lesson ends with a final-prep checklist for the whole course. Timed mock interviews have their own lesson, Mock Interviews & Communication.
It builds on the earlier lessons: estimation from Foundations of System Design, CAP and percentiles from Key Concepts & Terminology, replication and sharding from Databases & Storage, and delivery guarantees from Messaging & Queues.
Before any move: prove you need more than one box
Every move in this lesson costs latency, operational work or consistency. Interviewers notice when a candidate pays that price for no reason. An internal tool for 50 employees needs one web app and one PostgreSQL instance with backups. Put Kafka, a Redis cluster and six services in front of it, and you have shown that you can't size a problem.
Ask four questions before you reach for anything below:
| Ask | If the answer is no⦠|
|---|---|
| Does the data or traffic outgrow one machine plus a replica? | You don't need to partition, so you don't need consistent hashing |
| Do several teams need to release independently, or do parts scale very differently? | A modular monolith is the stronger answer |
| Must exactly one actor do something (one primary, one run of a job)? | You don't need leader election |
| Does one business action write to more than one database or external system? | A local transaction still works, so you don't need a saga |
π‘ What to say: "At 150 requests per second this fits on one database with a replica, so I'll keep it simple, and I'll tell you the numbers at which I'd split it." Then name those numbers. That shows judgement. Twelve boxes with no numbers show the opposite.
Move 1: Put keys on a ring
Pressure: a cache or storage tier where nodes join and leave, and every client must find a key's node without asking anyone.
Adding a node empties the cache
Your product cache has 4 nodes, and clients pick a node with hash(key) % 4. It serves 200,000 reads per second at a 95% hit rate, so the database behind it sees 10,000 reads per second. Before a big sale you add a fifth node and change the formula to hash(key) % 5.
Predict first: what fraction of keys now map to a different node, and what does the database see right after the change?
Check your answer
About 80% of keys move. A key stays put only when hash % 4 equals hash % 5, which holds for 4 out of every 20 hash values. Every moved key misses on its new node, so the hit rate falls to about 0.95 Γ 0.20 = 19%. The database goes from 10,000 to about 162,000 reads per second, 16 times its normal load, at the start of your busiest day.
In general, going from N to N + 1 nodes moves N/(N + 1) of the keys: 75% for 3 β 4, and 91% for 10 β 11.
The ring
Hash the node names too, and put nodes and keys on the same circle of hash values. A key belongs to the first node clockwise from its position. Unrolled into a line that wraps around at the end:
before: 0 ββAββββk1βββββBββββk2βββk3ββββCββββk4βββDβββΊ end, then back to 0
k1 β B k2 β C k3 β C k4 β D
add E: 0 ββAββββk1βββββBββββk2βEβk3ββββCββββk4βββDβββΊ
k1 β B k2 β E k3 β C k4 β D
Only k2 moved. It sat between B and E, and it moved from C to E.
Removing a node hands its keys to the next node clockwise. On average, adding one node to N moves 1/(N + 1) of the keys: exactly the new node's fair share.
import bisect
import hashlib
def h(s: str) -> int:
return int.from_bytes(hashlib.md5(s.encode()).digest()[:8], "big")
class Ring:
def __init__(self, nodes, vnodes=100):
self.vnodes = vnodes
self.points = [] # sorted hash positions
self.owner = {} # position -> physical node
for node in nodes:
self.add(node)
def add(self, node, weight=1):
for i in range(self.vnodes * weight): # a bigger box gets more points
p = h(f"{node}#{i}")
self.owner[p] = node
bisect.insort(self.points, p)
def node_for(self, key):
i = bisect.bisect(self.points, h(key)) % len(self.points) # next point clockwise
return self.owner[self.points[i]]
keys = [f"user:{i}" for i in range(100_000)]
nodes = ["cache-a", "cache-b", "cache-c", "cache-d"]
def moved(before, after):
return sum(b != a for b, a in zip(before, after)) / len(keys)
mod_before = [nodes[h(k) % 4] for k in keys]
mod_after = [(nodes + ["cache-e"])[h(k) % 5] for k in keys]
print(f"modulo 4 -> 5: {moved(mod_before, mod_after):.0%} of keys move")
ring = Ring(nodes)
ring_before = [ring.node_for(k) for k in keys]
ring.add("cache-e")
ring_after = [ring.node_for(k) for k in keys]
print(f"ring 4 -> 5: {moved(ring_before, ring_after):.0%} of keys move")
print("every moved key went to:", {a for b, a in zip(ring_before, ring_after) if b != a})
modulo 4 -> 5: 80% of keys move
ring 4 -> 5: 19% of keys move
every moved key went to: {'cache-e'}
With the ring, the sale starts at a hit rate of about 0.95 Γ 0.80 = 76%, so the database sees about 48,000 reads per second. That is still 4.8 times normal. Consistent hashing shrinks the spike; it doesn't remove it. Add capacity before the peak, one node at a time, and let each node warm up.
Why it works
A key's owner depends only on the key's position and on the nearest node clockwise. Adding a node changes the nearest node only for keys in the arc just before the new node, so no other key can move. With modulo, every key's owner depends on N itself, so changing N reshuffles almost everything.
Virtual nodes: many small arcs instead of one big one
With one point per node, arc lengths are random, and so is each node's share. Measured with the code above on 100,000 keys and 10 nodes:
| Points per node | Busiest node vs fair share | Idlest node vs fair share |
|---|---|---|
| 1 | 2.26Γ | 0.06Γ |
| 10 | 1.46Γ | 0.65Γ |
| 100 | 1.19Γ | 0.90Γ |
| 1,000 | 1.05Γ | 0.96Γ |
Virtual nodes give each physical node many points on the ring. That buys three things:
- Balance. Many small arcs average out, as the table shows.
- Spread-out failover. With one point per node, a dead node's keys all land on one neighbour. In a 5-node test with one point each, that neighbour already owned 41% of the keys because of an unlucky arc, and it inherited all of the dead node's keys too. With 100 points per node, the dead node's keys spread over all four survivors.
- Weights. Give a node with twice the RAM twice the points, and it owns about twice the keys. That is the
weightparameter in the code.
The cost is a bigger membership table and, in a replicated store, more peers per node. Cassandra calls each point a token. Its production guide suggests 16 tokens for clusters that grow and shrink often and 4 for clusters heading past 30 nodes, and warns that more tokens means sharing data with more peers, which lowers availability. Cassandra gets away with so few because it places tokens with an allocation algorithm (allocate_tokens_for_local_replication_factor), not at random as in the table above.
Replicas on the ring
To keep R copies of a key, walk clockwise and take the next R distinct physical nodes, skipping points that belong to a node you already picked. Amazon's Dynamo paper (2007) calls this list the key's preference list. Without the skip, two of a key's "three replicas" could be the same machine.
It isn't the only answer: fixed slots
Many systems get the same property another way. They split the key space into a fixed number of slots, many more than there are nodes, and keep a small table of which node owns which slot. A key never changes slot, and scaling moves whole slots between nodes.
- Redis Cluster has 16,384 hash slots:
slot = CRC16(key) mod 16384. Hash tags such as{user42}:cartput related keys in one slot (cluster spec). - Kafka's Java client picks a partition with
murmur2(key) mod partitions(clients built on librdkafka default to a CRC32 hash instead, so mixed-language producers must agree on one). Moving partitions between brokers doesn't change which partition a key goes to. Adding partitions changes the formula and remaps keys, which can break per-key ordering during the change.
Rendezvous hashing is a third option: score every node with hash(key + node) and pick the highest score. It moves the same minimal set of keys and needs no ring, but a lookup costs O(N).
Failure modes
- Hot keys stay hot. The ring balances keys, not traffic. One very popular key still lands on one node. A read-hot key, like a product page on sale, needs more copies: replicas, or a short in-process cache on the app servers. Splitting it into sub-keys would make every read gather all the pieces. Splitting (salting) is the fix for a write-hot key, such as a counter, whose pieces you sum on read. Both are covered in Databases & Storage.
- Cold after a change. Even the ring's move of about 20% sends a burst of misses to the database, as the sale example showed.
- Clients that disagree about membership. If half your app servers know about node E and half don't, the same key is read and written in two places: extra misses and stale values. Membership must come from one source, such as a configuration service or the cluster's own gossip, and every client must follow it.
- Nodes that flap. A node that drops out and rejoins every minute moves its keys twice each time. Remove a node after a grace period, not on its first missed heartbeat.
π‘ What to say: "I'll place keys with consistent hashing and about 100 virtual nodes per server, so adding a server moves only the new server's share of keys, roughly 1/(N + 1), instead of nearly all of them." Likely follow-up: "What about a hot key?" Answer: "The ring doesn't help there. If it's read-hot, I'd replicate it and cache it in process; if it's write-hot, I'd split it into sub-keys and sum them on read."
Your turn: a 10-node cache adds one node. How much of the cache goes cold with modulo placement, and how much with a ring with virtual nodes? Then the requirement changes: the new node has three times the RAM of the others. What do you change?
Check your answer
Modulo moves 10/11 β 91% of keys. The ring moves only the new node's share: 1/11 β 9% with equal weights.
With three times the RAM, give the new node three times the virtual nodes. Its share is its weight over the total weight: 3 out of 10 + 3 = 13, so it should own 3/13 β 23% of the keys, and only those keys move. Modulo can't express weights without changing the formula, and changing the formula moves nearly everything again.
Move 2: Split services where it pays, then let them find each other
Pressure: eight teams wait for one release train; image resizing needs 20 times the CPU of everything else; a bug in reporting takes checkout down with it.
Those are the three good reasons to split a service out: independent releases, different scaling, and failure isolation. "Microservices are modern" is not one of them. Each service should also own its data. Two services that share tables are really one service with two deploy pipelines.
What a split costs
Predict first: checkout calls 5 services one after another. Each is available 99.9% of the time, and they fail independently. What is checkout's availability, and how many minutes of failure does that allow per 30 days?
Check your answer
0.999β΅ β 99.5%. A 30-day month has 43,200 minutes. At 99.9% that allows 43 minutes of failure. At 99.5% it allows about 216 minutes, five times as much. Every synchronous dependency multiplies its own unavailability into yours.
Three more costs:
- Tail latency grows with fan-out. If each backend call is slow 1% of the time and a page calls 10 backends in parallel, 1 β 0.99ΒΉβ° β 9.6% of pages wait for at least one slow call. At 100 backends it is 63%. Jeff Dean and Luiz Barroso describe this effect in The Tail at Scale (2013). Percentiles are covered in Key Concepts & Terminology.
- No joins or transactions across services. A report that joined three tables now calls three APIs, and a transaction becomes a saga (Move 4).
- Operations. Every service needs a pipeline, dashboards, alerts, an on-call owner and a versioned API.
That is why a modular monolith is a common and defensible start: one deployable with strict module boundaries, split later along the seams that actually hurt. Stack Overflow has described serving its whole Q&A network in 2016 from one multi-tenant application on nine web servers, with separate applications for things like Careers and its APIs.
Keep a slow dependency from taking you down
A network call can fail in a way a function call can't: it can return nothing at all. Four standard defences:
| Defence | What it does | The trap it avoids |
|---|---|---|
| Timeout on every call | Frees the caller after a deadline | One slow dependency ties up every thread |
| Retry with exponential backoff and jitter, only for idempotent calls, under a retry budget | Rides out a short blip | Synchronized retry waves; double charges |
| Circuit breaker | After many failures, fails fast for a while, then lets a few trial calls through | Hammering a service that is already down |
| Bulkhead | Gives each dependency its own thread or connection pool | A slow recommendations call starving checkout |
β οΈ Retry storms. If each of three layers makes up to 3 attempts, one user request can become 3 Γ 3 Γ 3 = 27 calls at the bottom layer, just when that layer is already struggling. Retry at one layer only, cap retries at a small share of traffic (a retry budget), and pass the remaining deadline downstream so inner calls stop once the caller has given up.
How does Orders find Payments?
Instances come and go with every deploy, scale-up and crash, so hard-coded addresses go stale within a day. Service discovery has two halves: a registry that knows the healthy instances, and a way for callers to use it.
Client-side discovery
Orders ββ(1) "who serves payments?"βββΊ Registry
Orders βββββ 10.0.1.7, 10.0.2.9 βββββββββ
Orders ββ(2) request ββββββββββββββββββΊ 10.0.2.9 (payments-2) Orders picks the instance
Server-side discovery
Orders βββΊ payments.internal (load balancer or proxy) βββ¬βββΊ 10.0.1.7 (payments-1)
β² ββββΊ 10.0.2.9 (payments-2)
βββ the registry keeps the proxy's list of healthy instances current
| Client-side | Server-side | |
|---|---|---|
| How | The caller asks the registry, picks an instance and balances the load itself | The caller sends to one stable address; a load balancer or proxy picks the instance |
| Examples | Netflix Eureka with a client library; gRPC client-side balancing | Cloud load balancers; a Kubernetes Service |
| Gains | No extra hop; per-request balancing | Simple callers in any language |
| Costs | A discovery library in every language you use | An extra hop, and the proxy must itself be highly available |
Registration is done either by the instances themselves (register at start-up, then send heartbeats) or by the platform. Kubernetes does it for you. Only pods that pass their readiness probe are marked ready in a Service's EndpointSlices, and kube-proxy routes only to ready endpoints; a failing liveness probe restarts the container instead. A normal Service gets a stable virtual IP and a DNS name, and kube-proxy forwards each connection to a ready pod. A headless Service (clusterIP: None) returns the pod IPs themselves, for callers that balance on their own.
β οΈ Connections are not requests. L4 balancing, which is what a Kubernetes ClusterIP and most TCP load balancers do, picks a backend once per connection. Clients that keep a long-lived connection (HTTP/2, gRPC, database pools) keep using the pods they connected to first. Balancing each request needs an L7 proxy or client-side balancing. L4 vs L7 is covered in Core Building Blocks.
A service mesh (for example Istio with Envoy sidecars) moves discovery, per-request balancing, retries, mutual TLS and telemetry into a proxy next to every instance. You pay for an extra proxy on each side of every call, CPU and memory for every sidecar, and a control plane to run.
The registry must not be a single point of failure, so it is replicated. Some registries sit on a consensus store: Kubernetes keeps its state in etcd, and Consul's servers run Raft (Move 3). Others, like Eureka, replicate between peers without consensus and prefer a stale answer to no answer. When the registry is unreachable, callers should keep using their last known list instead of failing every call.
π‘ What to say: "Services find each other through Kubernetes Services, which only route to pods that pass their readiness checks. Every call has a timeout. Only idempotent calls are retried, with backoff and a budget, and the payment provider sits behind a circuit breaker." Likely follow-up: "What if the registry goes down?" Answer: "Callers keep their cached list, so existing traffic flows; new instances aren't discovered until it's back. That's why the registry itself is replicated, for example on a three- or five-member etcd cluster."
Your turn: checkout's p99 target is 300 ms. It calls Pricing (p99 40 ms), then Inventory (p99 60 ms), then Payments (p99 150 ms), one after another. A product manager asks for one more synchronous call, to a new Recommendations service (p99 120 ms), "to show related items on the confirmation page". What do you do?
Check your answer
Adding p99s is only a rough estimate, because the slow tails of independent calls rarely line up, but it shows the problem: 40 + 60 + 150 = 250 ms already, and another 120 ms makes 370 ms, over the target.
More importantly, a recommendation isn't needed to take the money. Keep it off the critical path. Render the confirmation, then let the page load recommendations separately. Or call Recommendations in parallel with a short timeout (say 50 ms) and show nothing if it misses. A failure in Recommendations must never fail a checkout: that is a timeout plus a fallback (and, if the call shares a thread pool with checkout, a bulkhead).
Move 3: Agree on one leader
Pressure: "exactly one": one primary database accepting writes, one scheduler firing the nightly job, one owner per partition.
Why "just pick one" is hard
A node can't tell a dead leader from a slow one or from a broken network link. All it sees is silence, and a timeout is only a guess. Guess too early and you get split brain: two nodes both acting as leader, both accepting writes. Guess too late and nothing happens for a long time.
Predict first: two scheduler replicas share a lock in a key-value store. The lock expires after 10 seconds and the holder renews it every 3 seconds. Replica A holds it and is halfway through writing the nightly payout file when its whole process freezes for 15 seconds, for example in a stop-the-world garbage-collection pause. What happens?
Check your answer
The lock expires at most 10 seconds into the pause (the last renewal came up to 3 seconds before the freeze), and replica B takes it and starts the payout run. At least five seconds later A wakes up. From A's point of view no time has passed, so it still believes it holds the lock, and it keeps writing. Now two writers are producing payouts.
A renewal thread doesn't save A, because the whole process was frozen. Checking the lock before each write doesn't save it either, because the pause can fall between the check and the write. The fix, fencing tokens, comes later in this move.
Majorities: why 3 or 5, not 2 or 4
Consensus protocols decide by majority. Any two majorities of the same group share at least one member. Since each member votes once per election, two candidates can never both win it, and every future majority includes at least one member that knows about every committed decision.
| Members | Majority | Failures tolerated |
|---|---|---|
| 1 | 1 | 0 |
| 2 | 2 | 0 |
| 3 | 2 | 1 |
| 4 | 3 | 1 |
| 5 | 3 | 2 |
| 7 | 4 | 3 |
An even size adds a machine without adding tolerance: four members survive one failure, just like three.
Placement matters as much as the count. Three members in three availability zones survive the loss of any one zone. The same three members in two zones, placed 2 + 1, do not: lose the zone with two, and the survivor is alone. With only two data centers, no layout survives losing the larger one, which is why two-site setups add a small third site whose only job is to vote.
A 99.99% target is usually met inside one region by surviving the loss of a zone. DynamoDB, for example, commits to 99.99% for a table in one region and to 99.999% for global tables that span regions (SLA). Multi-region consensus is for requirements such as "survive losing a whole region without losing a write".
Raft in five sentences
Raft is the consensus protocol behind etcd and Consul. At interview depth:
- Time is divided into numbered terms, and each term has at most one leader.
- A follower that hears no heartbeat within a randomized election timeout becomes a candidate, starts a new term and asks for votes. Each node votes at most once per term, and only for a candidate whose log is at least as up to date as its own.
- A candidate with a majority of votes wins. Randomized timeouts make split votes rare: the paper suggests 150β300 ms, and etcd defaults to a 1,000 ms election timeout with 100 ms heartbeats (tuning guide).
- The leader appends each write to its log and replicates it. An entry is committed once a majority stores it, and only then is the client told "done".
- Every majority overlaps the one that stored a committed entry, and voters reject candidates with older logs, so a new leader always has every committed entry.
ZooKeeper uses its own protocol, Zab, built on the same majority logic.
Trace: a partition
Five members, A to E. A leads term 7. The network splits A and B away from C, D and E.
| Time | Minority side: A and B | Majority side: C, D and E |
|---|---|---|
| 0 s | A still thinks it leads, and clients send it writes | Heartbeats from A stop arriving |
| about 1 s | A appends the writes and B stores them. That is 2 of 5, so nothing commits, and those clients time out | C's election timeout fires. C starts term 8 and wins with votes from D and E |
| 1β20 s | Still nothing commits. B can't win an election with 2 votes | C's writes reach 3 of 5 members and commit |
| heal | A and B see term 8 and follow C. Their uncommitted entries are overwritten by C's log | Unchanged |
The minority side is unavailable for writes during the partition. That is the consistency-over-availability choice of CAP made concrete (Key Concepts & Terminology). There is one subtle trap: A could still answer reads from its stale state, unless reads are confirmed with a majority too. etcd does this for its default, linearizable reads.
The price: every write waits for a majority
A commit needs acknowledgments from a majority, so commit latency is roughly the round trip to the member that completes the fastest majority. Three members in three zones of one region: a millisecond or two per write. Members on different continents: tens to hundreds of milliseconds per write, all the time, not only during failures. That is why most systems keep a consensus group inside one region and replicate between regions asynchronously.
β οΈ Isolation is not consensus. "PostgreSQL with serializable isolation" describes how concurrent transactions on one primary behave. It says nothing about replicas. A primary with asynchronous replicas can still lose acknowledged writes when it fails over. CP-like behaviour needs synchronous or quorum replication and a failover process that can never produce two primaries.
Use it, don't build it
You will almost never implement Raft yourself. You use a system that has it built in:
- Coordination stores (etcd, ZooKeeper, Consul) for leader election, locks, configuration and service registries. Kubernetes keeps all cluster state in etcd and elects its controllers with Lease objects.
- Databases and brokers. Kafka keeps its metadata in its own Raft-based quorum, KRaft; version 4.0 dropped ZooKeeper support entirely. Spanner-style databases replicate every shard with consensus.
- Failover managers. Patroni keeps a PostgreSQL cluster's leader lock in etcd, Consul, ZooKeeper or Kubernetes with a TTL (30 seconds by default, per its configuration docs). A primary that can't renew the lock is demoted instead of being left running as a second primary.
Leases and fencing tokens
A lease is a lock with an expiry, so a crashed holder can't block everyone forever. The expiry causes the problem you predicted: a paused holder doesn't know its lease has expired. Martin Kleppmann's How to do distributed locking gives the cure. Every grant comes with a fencing token, a number that only ever increases, and the resource being protected rejects any token older than the newest one it has seen.
| Step | Event | Highest token the storage has seen | Result |
|---|---|---|---|
| 1 | A gets the lease with token 33, then freezes | none | |
| 2 | The lease expires; B gets it with token 34 | none | |
| 3 | B writes with token 34 | 34 | Accepted |
| 4 | A wakes up and writes with token 33 | 34 | Rejected, because 33 < 34 |
class FencedStore:
"""The protected resource checks tokens itself."""
def __init__(self):
self.highest = 0
self.rows = {}
def write(self, key, value, token):
if token < self.highest: # an older leader woke up
raise PermissionError(f"stale token {token}, already saw {self.highest}")
self.highest = token
self.rows[key] = value
store = FencedStore()
store.write("payouts/2026-09-24", "file from B", token=34)
try:
store.write("payouts/2026-09-24", "file from A", token=33)
except PermissionError as err:
print("rejected:", err)
print(store.rows)
It prints rejected: stale token 33, already saw 34, and only B's file remains.
Tokens come from the consensus store: etcd's revision numbers and ZooKeeper's transaction ids only ever increase. The lease alone can't protect you; only a check at the resource can. When the resource can't check tokens (an email provider, a bank's API), make the operation itself idempotent, for example keyed by the payout date, so a duplicate run does no harm.
Two more failure modes:
- Clocks. A lease holder should time its lease on a local monotonic clock and stop working well before the expiry, because wall clocks jump and machines disagree about the time. Even so, only fencing makes it safe against pauses.
- Election storms. A timeout that is too short for your network causes needless elections, and nothing commits while one runs. etcd's tuning guide ties the timeouts to the round-trip time between members.
π‘ What to say: "Only one scheduler may run, so the replicas elect a leader through an etcd lease. A lease can outlive a paused process, so every write carries a fencing token that the database checks, and the job is idempotent per run date anyway." Likely follow-up: "Why not a Redis lock?" Answer: "A single Redis lock gives me no fencing token, and a failover can lose the lock. That's fine for efficiency, to avoid duplicate work, but not for correctness."
Your turn: your company has two data centers and wants the payments database to fail over automatically, with no split brain, if either data center is lost. The team proposes a 4-member etcd cluster, 2 members in each data center, to decide which side runs the primary. Does it meet the requirement?
Check your answer
No. Four members need 3 votes. Lose either data center and 2 remain: no majority, no election, no automatic failover. A 2 + 1 layout doesn't help either, because losing the side with two members leaves one of three.
Add a third site, even a single small cloud VM that holds no data, so the members are placed 2 + 2 + 1. Now either data center can fail and 3 of 5 remain. The third site only has to vote, but it must be reachable from both data centers over network paths that don't fail together.
Move 4: Replace the cross-service transaction with a saga
Pressure: one business action writes to several services' databases or calls an external API. That is the checkout from the start of this lesson.
Why not two-phase commit?
In two-phase commit (2PC), a coordinator first asks every participant to prepare: make the change durable but not yet visible, keep its locks, and vote yes or no. If every vote is yes, it tells them all to commit.
Coordinator Orders Stock Payments
βββ PREPARE (to all three) βββββΊ β β β write, lock rows, vote yes
ββββ three YES votes βββββββββββ
βββ COMMIT (to all three) ββββββΊ β β β make visible, release locks
If the coordinator dies between the two phases, every participant keeps its locks
and cannot decide alone. It stays "in doubt" until the coordinator returns.
2PC gives real atomicity, and it is the right tool inside one database system. Google's Spanner, for example, runs 2PC between shards that are themselves replicated with Paxos, so the coordinator can't simply disappear. Across independent services it breaks down:
- It blocks. If the coordinator fails after the votes, participants that hold locks can neither commit nor abort on their own.
- Locks are held across network round trips, and every participant must be up for anything to commit. That is the serial availability math from Move 2.
- Every participant must support it. Your card provider's HTTP API, an email service or Kafka won't join your transaction. Kafka's own transactions are atomic only across Kafka partitions.
The saga
A saga is a sequence of local transactions, each committed on its own, where each step has a compensating action that undoes it in business terms. Hector Garcia-Molina and Kenneth Salem named the idea in 1987. If step k fails, you run the compensations for steps k β 1 back down to step 1.
Compensation is not rollback. You can't un-send an email or un-charge a card; you send a correction or issue a refund. So the order of the steps matters:
- Reversible steps first, the ones that are cheapest to undo: reserve stock, redeem a coupon.
- The pivot: the last step that is allowed to say "no" for a business reason. Usually that is the charge.
- After the pivot, only steps that can't fail for business reasons, so you can retry them until they succeed: create the shipment, send the receipt.
Here is an orchestrator for that checkout, with a declined card:
class Service:
"""A fake downstream service that remembers idempotency keys."""
def __init__(self, name, says_no_to=(), loses_reply_to=()):
self.name = name
self.says_no_to = set(says_no_to)
self.loses_reply_to = set(loses_reply_to)
self.results = {} # idempotency key -> first result
def call(self, op, key):
if key in self.results: # a retry: same answer, no new work
return self.results[key]
if op in self.says_no_to:
raise ValueError(f"{self.name}.{op}: declined")
self.results[key] = f"{self.name}.{op} done"
print(" ", self.results[key])
if op in self.loses_reply_to: # the work happened, the reply is lost
self.loses_reply_to.discard(op)
raise TimeoutError(f"{self.name}.{op} timed out")
return self.results[key]
def call_with_retry(svc, op, key, attempts=3):
for _ in range(attempts):
try:
return svc.call(op, key)
except TimeoutError as err:
print(" retrying with the same key after:", err)
raise RuntimeError(f"{svc.name}.{op} keeps timing out: page a human")
stock = Service("stock", loses_reply_to={"release"})
coupons = Service("coupons")
payments = Service("payments", says_no_to={"charge"}) # the card is declined
shipping = Service("shipping")
SAGA = [ # (service, action, compensation)
(stock, "reserve", "release"),
(coupons, "redeem", "restore"),
(payments, "charge", "refund"), # pivot: the last step allowed to say no
(shipping, "create_shipment", None), # after the pivot: retry until it succeeds
]
def run_saga(order_id):
done = [] # in production: a durable saga-state row
for svc, action, undo in SAGA:
try:
call_with_retry(svc, action, key=f"{order_id}:{action}")
done.append((svc, undo))
except ValueError as err: # a business "no": compensate in reverse
print(" ", err, "-> compensating")
for prev_svc, prev_undo in reversed(done):
call_with_retry(prev_svc, prev_undo, key=f"{order_id}:{prev_undo}")
return "CANCELLED"
return "CONFIRMED"
print(run_saga("order-42"))
stock.reserve done
coupons.redeem done
payments.charge: declined -> compensating
coupons.restore done
stock.release done
retrying with the same key after: stock.release timed out
CANCELLED
Compensations run in reverse order: the coupon first, then the stock. Look at release. The stock service did the work, but the reply was lost. The retry used the same key, so the service returned its first answer instead of releasing the stock twice.
Why it works. Every step is either compensatable, or comes after the pivot and can be retried until it succeeds. So a saga ends in one of two states: everything done, or everything compensated. That holds only if every call is idempotent, because timeouts force retries, and if the saga's progress is stored durably, so a crashed orchestrator can resume from the last recorded step.
Orchestration or choreography
| Orchestration | Choreography | |
|---|---|---|
| How | A coordinator (a saga state machine) calls each step and each compensation | Each service reacts to events and emits its own: StockReserved β Payments charges β PaymentFailed β Stock releases |
| Gains | The whole flow is in one place: easy to read, time out and monitor | No central component; services stay loosely coupled |
| Costs | One more service to run, and it knows every step | The flow exists only as the sum of subscriptions: hard to trace, easy to create cycles |
| Good fit | Many steps, compensations, timeouts | Two or three steps with simple reactions |
Workflow engines such as Temporal and AWS Step Functions store an orchestrator's state for you.
What a saga gives up: isolation
A saga is atomic through compensation, but it has no isolation. Other requests see the intermediate states. A customer can see an order as "placed" that is cancelled two seconds later, and a stock count can drop and come back. The standard countermeasures:
- Semantic locks: mark records as
PENDING, and make other flows respect that status. - Commutative updates: "add 1 back to stock" is correct whatever happened in between; "set stock to 7" is not.
- Version checks: before an update that depends on earlier state, re-read it and compare versions.
Paying safely: the idempotency key, done right
The charge is where a careless retry costs real money. A safe pattern:
- Before calling the provider, store a row
(idempotency_key, order_id, status = PENDING)in your own database. - Call the provider with the same key. Stripe, for example, saves the result of the first request for a key, whether it succeeded or failed, and returns that result for every retry with the key. It may prune a key once it is 24 hours old (docs), so recovery must run well within a day.
- Record the outcome,
SUCCEEDEDorFAILED, against the key. - After a crash, a recovery job finds
PENDINGrows and repeats step 2 with the same key. The provider returns the original result instead of charging again.
The row written first tells recovery that a charge may be in flight. The key sent to the provider makes repeating the call harmless. If you store a result only after the call succeeds, a crash in between makes the retry look like a brand-new payment.
Failure modes
- A compensation fails. Retry it, since it is idempotent. If it keeps failing, park the saga in a
NEEDS_ATTENTIONstate and alert. A saga that silently stops halfway is the worst outcome. - The orchestrator crashes. It resumes from its stored state and re-sends the current step with the same key.
- Timeouts are ambiguous. No reply means "maybe it happened". Retry with the same idempotency key, never with a new one.
- The pivot's result is unknown. Find out by querying the provider, or retrying with the same key, before you compensate anything.
π‘ What to say: "Checkout spans three services and a card provider, so I'll use an orchestrated saga: reserve stock, charge, then create the shipment. Every step has an idempotency key and a compensation, and shipping comes after the charge because it can be retried until it succeeds. The cost is isolation: for a moment an order can look placed and then be cancelled." Likely follow-up: "Why not 2PC?" Answer with the three reasons above: it blocks, it holds locks across the network, and the card provider can't take part.
Your turn: finance now wants customers charged only when the parcel ships, up to 3 days after checkout. What changes in the saga?
Check your answer
Split the charge in two. At checkout, authorize the card: this reserves the funds and is the step that can be declined, so it becomes the pivot. Its compensation is voiding the authorization. When the parcel ships, capture the authorized amount.
The saga is now long-running, measured in days, so its state must live in a durable store or a workflow engine with timers, not in memory. Stock also stays reserved for days, so other flows must respect the PENDING reservation. Two new failure cases need a plan: the authorization can expire before shipment (card networks and providers limit how long it lasts), and a capture can still fail. Re-authorize, or contact the customer and cancel with compensations.
Move 5: Commit once, publish from the log
Pressure: "when X is saved, tell the other services." The order is committed, and Loyalty, Email and Analytics must hear about it.
The dual-write trap
Predict first: the process can crash between any two lines. Which version is safe?
# Version 1 # Version 2
db.commit(order) broker.publish(order_placed)
broker.publish(order_placed) db.commit(order)
Check your answer
Neither. In version 1, a crash between the lines leaves an order that nobody hears about: no loyalty points, no email, a gap in analytics. In version 2, the event goes out and then the commit fails (a constraint violation, a deadlock, a crash), so every consumer acts on an order that doesn't exist.
Retrying the second write narrows the window but can't close it, because the process can still crash between the lines. And you can't make a database and a broker commit together without a distributed transaction that brokers such as Kafka don't offer.
The transactional outbox
Write the event into the same database, in the same transaction, as the business row. A separate relay publishes it afterwards.
Order service PostgreSQL Kafka
BEGIN ββββββββββββββββββββββββ
INSERT order ββββββββΊ β orders β
INSERT outbox row βββΊ β outbox (not yet sent)β βββ relay reads unsent rows
COMMIT ββββββββββββββββββββββββ and publishes them ββββΊ topic "orders"
then marks them sent
Now the order and its event can't disagree: both rows commit, or neither does. The relay is either a poller (select unsent rows, publish them, mark them sent) or change data capture (CDC): a tool such as Debezium reads the database's own write-ahead log and publishes the new outbox rows. Debezium ships an outbox event router for exactly this.
Here it is end to end, with SQLite standing in for both services' databases, and a relay that crashes at the worst moment:
import json
import sqlite3
import uuid
orders_db = sqlite3.connect(":memory:")
orders_db.executescript("""
CREATE TABLE orders (id TEXT PRIMARY KEY, total_cents INTEGER NOT NULL);
CREATE TABLE outbox (seq INTEGER PRIMARY KEY AUTOINCREMENT,
event_id TEXT NOT NULL, payload TEXT NOT NULL,
sent INTEGER NOT NULL DEFAULT 0);
""")
def place_order(order_id, total_cents):
with orders_db: # ONE local transaction: both rows or neither
orders_db.execute("INSERT INTO orders VALUES (?, ?)", (order_id, total_cents))
event = {"type": "OrderPlaced", "order_id": order_id, "total_cents": total_cents}
orders_db.execute("INSERT INTO outbox (event_id, payload) VALUES (?, ?)",
(str(uuid.uuid4()), json.dumps(event)))
topic = [] # stand-in for a Kafka topic
def relay_once(crash_after_publish=False):
rows = orders_db.execute(
"SELECT seq, event_id, payload FROM outbox WHERE sent = 0 ORDER BY seq").fetchall()
for seq, event_id, payload in rows:
topic.append((event_id, payload)) # publish
if crash_after_publish:
return # dies before marking the row sent
with orders_db:
orders_db.execute("UPDATE outbox SET sent = 1 WHERE seq = ?", (seq,))
loyalty_db = sqlite3.connect(":memory:")
loyalty_db.executescript("""
CREATE TABLE points (order_id TEXT PRIMARY KEY, points INTEGER NOT NULL);
CREATE TABLE inbox (event_id TEXT PRIMARY KEY);
""")
def consume(event_id, payload):
event = json.loads(payload)
with loyalty_db: # dedupe and effect in ONE transaction
if loyalty_db.execute("SELECT 1 FROM inbox WHERE event_id = ?", (event_id,)).fetchone():
return # already processed: a redelivery
loyalty_db.execute("INSERT INTO inbox VALUES (?)", (event_id,))
loyalty_db.execute("INSERT INTO points VALUES (?, ?)",
(event["order_id"], event["total_cents"] // 100))
place_order("o-1", 4999)
relay_once(crash_after_publish=True) # relay dies after publishing
relay_once() # restarted relay publishes o-1 again
print("messages on the topic:", len(topic))
for event_id, payload in topic:
consume(event_id, payload)
print("points rows:", loyalty_db.execute("SELECT * FROM points").fetchall())
messages on the topic: 2
points rows: [('o-1', 49)]
The relay died after publishing but before marking the row as sent, so after the restart it published the same event again. That is the price of the outbox: delivery is at least once, never exactly once. So consumers must be idempotent. Record each processed event id in the same transaction as the effect it causes (an inbox table), and skip ids you have already seen. Loyalty got two messages and awarded the points once. More on this in Messaging & Queues.
Two more rules:
- Ordering. A CDC relay reads the log, so it sees transactions in commit order. A poller publishes in
seqorder, which is insert order: under concurrent transactions a row with a lowerseqcan commit after one with a higherseq. For a single entity that still matches commit order, as long as that entity's writes are serialized (for example by its row lock). So poll by thesentflag, as the code does, never byseq > last_seen, or a row that commits late is skipped for good. Use the entity id (here, the order id) as the message key, so all events for one order land on one partition, in order. - Housekeeping. Delete or archive sent rows so the relay's query stays fast, and alert on the age of the oldest unsent row. A stuck relay delays events without losing them, and "delayed by 40 minutes" is still an incident.
Side effects outside your database
Email, SMS and payment calls can't join your transaction at all. Suppose the email consumer sends an email and then records that it sent it:
- Send, then record: a crash in between makes the retry send it again (at least once).
- Record, then send: a crash in between means it is never sent (at most once).
No ordering of two separate writes gives you exactly once. Pick the failure you can live with. A rare duplicate receipt is usually fine; a lost password-reset email is not. If the provider accepts an idempotency key, send one derived from the event id, and the provider can drop the duplicate for you.
Event sourcing: the log is the database
The outbox keeps events next to the state. Event sourcing goes further: the events are the source of truth, and current state is computed by replaying them.
WalletOpened +500
PurchaseMade β150
RefundReceived +100
balance = 500 β 150 + 100 = 450
You gain a complete audit history, the ability to rebuild state as of any moment, and new read models built later by replaying old events. You pay with snapshots (replaying millions of events on every read is too slow), event versioning (events written years ago must still be readable), and a separate read model for most queries. Event sourcing is a storage model, not a product. An append-only table in PostgreSQL works; Kafka is a way to move events, not a requirement.
CQRS: separate models for writing and reading
Command Query Responsibility Segregation (CQRS) keeps the write model (normalized, validated, transactional) apart from one or more read models (denormalized and shaped for each screen, perhaps in a search engine). The read models are fed by events or by CDC. It fits when reads and writes have very different shapes: rich search over a catalogue that changes a few times a second, or dashboards over an event-sourced ledger.
The cost is lag. A read model is eventually consistent: usually milliseconds to seconds behind, and longer during incidents. So a user who saves and immediately reloads may not see their own change. Three fixes: return the new state in the command's response; read the author's own recent writes from the write model; or let the client wait until the read model has reached the version its command produced.
Failure modes
- Relay down: events pile up in the outbox, delayed but safe. Alert on the oldest unsent row.
- Consumer crashes after acting but before committing its offset: the event is delivered again, and the inbox check makes that harmless.
- A buggy projection corrupts a read model: fix the code and rebuild the read model by replaying the log. That is one of the payoffs of keeping events.
π‘ What to say: "I won't write to the database and to Kafka separately, because a crash between them loses or invents events. The order and its event go into one PostgreSQL transaction through an outbox table. Debezium publishes it keyed by order id, and consumers dedupe by event id." Likely follow-up: "So is that exactly once?" Answer: "Delivery is at least once. The effect happens once, because each consumer dedupes in the same transaction as the effect."
Your turn: the requirement changes. Legal says the Loyalty service must be able to prove, a year later, why every customer has the points they have. What do you store, and what does it cost?
Check your answer
Make loyalty event-sourced. Keep every PointsAwarded, PointsRedeemed and PointsExpired event, each with the id of the order or campaign that caused it, in an append-only store, and derive balances from those events. The audit question becomes a replay.
The costs: a snapshot of each balance, so reads don't replay a year of events; versioned event schemas; and a separate read model for queries such as "top customers this month". The OrderPlaced events the Orders outbox already publishes are the input for awards, so the two patterns fit together; redemptions and expiries are Loyalty's own events.
Move 6: See inside the system
Pressure: "p99 doubled after Tuesday's deploy, and checkout touches nine services. Which one did it?" In a monolith you would open one log file. Here you need to follow one request across many machines.
Three signals, three jobs
| Signal | What it is | Answers | Cost |
|---|---|---|---|
| Metrics | Numbers aggregated over time | "Is something wrong, and since when?" Cheap to keep and to alert on | Low per series, but it multiplies with every label |
| Logs | One structured record per event | "What exactly happened to this request?" | High at volume |
| Traces | One request's tree of timed spans across services | "Where did the time go?" | Too big to keep every request, so you sample |
OpenTelemetry is the vendor-neutral standard for producing all three.
A trace, step by step
Every request gets a trace id at the edge. Each service records a span (name, start, duration, parent span) and passes the context on to the services it calls, in a header. The W3C Trace Context standard names that header traceparent. It holds a version, a 32-hex-digit trace id, the 16-hex-digit id of the calling span and some flags, for example 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01 (spec).
Here is one slow checkout's trace as a waterfall:
span start duration
checkout POST /checkout 0 ms 1,180 ms ββββββββββββββββββββββββββββββ
ββ cart GET /cart 5 ms 40 ms β
ββ pricing POST /price 45 ms 60 ms ββ
ββ promotions POST /apply 105 ms 930 ms βββββββββββββββββββββββ
β ββ 48 Γ SELECT promo_rules 110 ms 900 ms ββββββββββββββββββββββ
ββ payments POST /charge 1,040 ms 120 ms βββ
It takes seconds to read. A trace from before the deploy showed promotions at about 30 ms with a single query. Now one query runs 48 times inside it: an N+1 query that the deploy introduced. Metrics told you checkout got slow; nine services' logs wouldn't easily tell you which one; the trace points straight at it. Put the trace id in every log line as well, so you can jump from the slow span to that request's logs.
β οΈ Asynchronous hops break traces unless you carry the context across them. Put traceparent into message headers, including in the outbox row from Move 5, so the relay publishes it with the event.
What to measure: RED, USE and the golden signals
- For every service, measure RED: rate (requests per second), errors, and duration as percentiles such as p50 and p99, never only averages.
- For every resource (CPU, disk, connection pool, queue), measure USE: utilization, saturation (work waiting in a queue), and errors.
- Google's SRE book names four golden signals that cover much the same ground: latency, traffic, errors and saturation.
β οΈ Cardinality. Every unique combination of label values is a separate time series. A latency histogram with 50 endpoints Γ 5 status classes Γ 4 regions Γ 12 buckets is 12,000 series. Add a user_id label with 10 million users and it becomes 120 billion. Prometheus's naming guide says it plainly: don't put user ids or other unbounded values in labels. Per-user detail belongs in logs and traces.
SLOs and error budgets: page on what users feel
An SLI measures user experience, for example "the fraction of checkout requests that succeed in under 800 ms". An SLO is the target for it, for example "99.9% over 30 days". The error budget is what's left over: 0.1% of requests. For a total outage that is 30 Γ 24 Γ 60 Γ 0.001 = 43.2 minutes a month.
Alert on how fast the budget is burning, not on CPU. The burn rate is the error rate divided by the budgeted error rate. 1% errors against a 0.1% budget is a burn rate of 10, which empties a 30-day budget in 3 days. Google's SRE workbook suggests paging when a 1-hour window burns at a rate of 14.4 or more (2% of the month's budget in one hour, confirmed over a 5-minute window) and opening a ticket for slow burns. A pod at 90% CPU that users don't notice is not a page.
Sampling: you can't keep every trace
Assume 20,000 requests per second, 12 spans per request and 500 bytes per span. That is 120 MB/s, or about 10.4 TB of traces a day.
- Head sampling decides at the edge, for example "keep 1%": about 104 GB a day. It is cheap and simple, but it keeps a random 1%, so it misses most rare slow requests.
- Tail sampling decides after the trace has finished, for example "keep every error, every trace over 1 second, and 1% of the rest". It keeps the interesting traces, but a collector must buffer all the spans of every trace until it decides.
Failure modes
- Monitoring in the same failure domain as the thing it watches: the region goes down and takes its dashboards with it. Run at least one external probe from somewhere else.
- Alert fatigue: pages that nobody acts on teach people to ignore pages. Every page should be about user impact and need a human now.
- Averages hide tails: a p99 regression can leave the mean almost unchanged.
π‘ What to say: "I'd define SLOs on checkout success and latency and page on burn rate. Every service exports RED metrics, and OpenTelemetry tracing with tail sampling keeps all errors and slow requests, with the trace id in every log line." Likely follow-up: "How do you roll out changes safely?" Answer: "Canary each deploy on a small share of traffic, and compare its SLIs with the baseline before going wider."
Your turn: your payments API's SLO is 99.9% of requests succeeding over 30 days. An incident makes 5% of requests fail. What is the burn rate, how long until the budget is gone, and should this page someone at 3 a.m.?
Check your answer
Burn rate = 5% / 0.1% = 50. The 30-day budget lasts 30 / 50 = 0.6 days, about 14.4 hours, and each hour of the incident uses 50/720 β 7% of the month's budget. That is far above the 1-hour paging threshold, a burn rate of 14.4, so yes: page.
Final round: no label on the problem
Real prompts don't say which move they want. Find the pressure, pick the move, and say what it costs.
1. The triple digest
A weekly-digest job runs as a cron entry on each of 3 app servers, so users get 3 emails. Make it one email per user, even if a server dies mid-run.
Check your answer
Two layers. One runner: elect a leader through a Kubernetes Lease or an etcd lease, or move the job to a scheduler that runs it once. Harmless duplicates: a lease can still overlap during a pause, so make each send idempotent. Before sending, skip users already listed in a digests_sent table for this week; after sending, insert (user_id, week) under a unique constraint. If the provider accepts an idempotency key, derive it from the same pair.
If the leader dies halfway, the next leader reruns the job and skips everyone already done. Cost: an extra write per email, and a crash between sending and recording can still cause a rare duplicate, which is acceptable for a digest.
2. More partitions, same order
A Kafka topic keyed by account_id has 12 partitions, and consumers rely on per-account order. Throughput now needs 24 partitions. What goes wrong if you simply add 12 partitions, and what do you do instead?
Check your answer
The partition is a hash of the key mod the partition count (murmur2 in the Java client). A key keeps its partition only if its hash mod 24 is below 12, so about half of the accounts move to a new partition. New events for a moved account can then be consumed before older events still waiting in its old partition: out of order.
Instead, create a new 24-partition topic and cut over: producers switch to the new topic, and consumers finish draining the old topic before they read the new one. Or plan partitions for growth from day one, which is the fixed-slots idea from Move 1. Cost: a planned migration, or extra partitions carried from the start.
3. Two accounts, two shards
A bank transfer debits an account on shard 3 and credits an account on shard 7 of your own distributed SQL database. Nobody may ever see money missing or doubled. Saga or not?
Check your answer
Not a saga. A saga exposes intermediate states (debited but not yet credited), which the requirement forbids. Both rows live in one database system that you control, so use its distributed transactions. Distributed SQL databases such as Spanner or CockroachDB offer ACID transactions across shards; Spanner, for example, runs two-phase commit between consensus-replicated shards, so the coordinator can't vanish. Cost: higher write latency and contention between transfers on the same accounts. If you sharded PostgreSQL yourself, the choice is harder: either accept a saga's visible intermediate state, which this requirement forbids, or run two-phase commit across your own shards (PostgreSQL supports PREPARE TRANSACTION) and own its blocking failure mode; a distributed SQL database does the latter for you. Sagas are for steps that live in systems that can't share a transaction.
4. The ghost listings
Search sometimes shows products that were deleted minutes ago. The product service deletes the row in PostgreSQL, then calls the search service to remove the document. About one call in 5,000 times out and isn't retried.
Check your answer
It's a dual write, failing just as Move 5 predicts. Write a ProductDeleted event to an outbox in the same transaction as the delete (or let CDC capture the delete itself), and let a relay feed the search indexer, which applies events idempotently by product id.
Cost: search becomes eventually consistent through the relay, usually well under a second behind, and you now watch relay lag. Adding retries to the synchronous call would shrink the gap but not close it: a crash between the delete and the call still leaves a ghost.
5. The 2 a.m. slowdown
Every night around 2 a.m., the orders API's p99 triples for 20 minutes, while its CPU graphs look normal. You have per-service RED metrics and tracing with tail sampling.
Check your answer
Pull the slow traces the tail sampler kept from 2:00β2:20 and see which span grows. Suppose it's the database span. Then ask what else runs at 2 a.m.: a backup, a batch job or index maintenance saturating the disk. Check USE on the database's disk and connection pool, not CPU. The fix follows the cause: run the batch job against a replica, throttle it, or move it away from traffic.
The general move: when resource graphs look fine, let traces show where the time goes, then use USE on that component to find why.
Cheat sheet: pressure β move β cost
| Pressure in the requirements | Move | What you give up |
|---|---|---|
| Nodes join and leave; clients must find a key's node | Consistent hashing with virtual nodes, or fixed slots plus a map | A membership source everyone agrees on; hot keys still need their own fix |
| Teams blocked on one release; parts that scale very differently | Split along business capabilities; each service owns its data | Network failures, tail latency, no cross-service joins or transactions |
| Instances come and go | Service discovery: a registry plus readiness checks, client-side or server-side | A registry that must be highly available; stale lists |
| A slow or failing dependency | Timeouts, retries with backoff and a budget, circuit breakers, bulkheads | Tuning work; degraded answers instead of errors |
| Exactly one actor | Leader election on a 3- or 5-member consensus store, plus fencing tokens | The minority side can't write; a majority round trip on every write |
| One action across several services | Saga: reversible steps, then the pivot, then retry-until-done steps | Isolation: intermediate states are visible |
| Save and notify | Transactional outbox, relay or CDC, idempotent consumers | At-least-once delivery; a relay to operate |
| Audit history, rebuildable views | Event sourcing, with CQRS read models | Snapshots, event versioning, read-model lag |
| Many services, one slow request | Traces, RED metrics, SLO burn-rate alerts | Sampling choices, instrumentation work, cardinality discipline |
Final prep: the checklist
Tick an item only if you can do it out loud, on a blank page, without notes.
| Can you⦠| Review |
|---|---|
| Turn DAU into the peak QPS, storage and bandwidth numbers that decide something, with units straight? | Foundations of System Design |
| State CAP and PACELC precisely, and turn a number of nines into downtime? | Key Concepts & Terminology |
| Choose between WebSockets, SSE and long polling, and say what a CDN caches and for how long? | Networking Basics |
| Say when a load balancer should work at L4 or L7, and how you would generate unique ids? | Core Building Blocks |
| Choose a database from the access pattern, pick a shard key, and pick a caching strategy? | Databases & Storage |
| Explain at-least-once delivery, idempotent consumers and per-partition ordering? | Messaging & Queues |
| Run the four phases, Scope β Sketch β Deep dive β Wrap-up, in a 45-minute slot, closing scope by minute 10, the sketch by 20 and one deep dive by 35? | Interview Framework & Strategy |
| Scope: ask the questions that change the design, and turn the answers into the two or three numbers that decide something? | Requirements Gathering |
| Sketch an API, a data model and a diagram with one write and one read traced end to end, then deep-dive the riskiest part with its failure modes? | Design Process Steps |
| Attack a prompt you have never seen with a repeatable playbook? | Realistic Design Examples |
| Design a rate limiter and a URL shortener from memory, numbers included? | Design URL Shortener & Rate Limiter, URL Shortener (e.g. TinyURL) |
| Handle fan-out and celebrities, video upload and delivery, and chat delivery and receipts? | Design Social & Streaming Systems, Twitter/Instagram News Feed, YouTube / Netflix Streaming, Chat Application (WhatsApp) |
| Use each of the six moves in this lesson, and name its cost? | The cheat sheet above |
Practice: change one requirement
Memorized designs fall apart the moment an interviewer changes one number, so practise the change. Take a design you know and rerun it with one mutation at a time:
- writes grow 100Γ, or reads flip from 100:1 to 1:1;
- one user has 50 million followers, or one product gets 30% of all traffic;
- events for the same entity must be processed strictly in order;
- users spread over three continents, and each needs a 100 ms p99;
- data about EU users must stay in the EU;
- the budget drops by 80%: what do you cut first?
For each one, say which components stay, which change, and what the new design gives up. After each rep, write down one sentence: the hardest trade-off you made. Ten minutes on a mutation is better practice than an hour of rereading.
When you don't know something
You will blank on a detail. Say what you know, what you don't, and how much it matters to the decision: "I don't remember Redis Cluster's exact failover timing. The design only needs it to be seconds, not minutes, and I'd confirm that in the docs." Then keep going. A precise "I don't know" is worth more than a confident guess that the interviewer can see through.
When to start mock interviews
Now. Once you can take a prompt you haven't seen through Scope, Sketch, Deep dive and Wrap-up, a timed mock teaches you more than another read-through. Use this checklist between mocks to decide what to review. Mock Interviews & Communication shows how to run one end to end and how to grade yourself.
Before moving on
Pick one move from this lesson and explain it aloud, as if an interviewer had just asked "what happens when�". Cover the pressure that calls for it, how it works in two sentences, what breaks without it, what it costs, and the follow-up question you would expect next. If you can do that for all six moves, you are ready for the mocks.