Key Concepts & Terminology

Master essential vocabulary and concepts used in every system design discussion.

Last generated

Lesson 2 of 18 available16 practice questions

SPACED REPETITION Β· 16 practice questions

Make this lesson stick.

Try 3 questions now. No account needed. Sample answers aren't saved.

One box, twenty times the traffic

Your team runs a ticketing site. Everything lives on one machine: the web app and PostgreSQL share a 16-core server. On a normal evening it serves 300 requests per second, and each request costs about 8 ms of CPU: 5 ms in the app and 3 ms in the database. That is 2.4 cores of work, so the box is 15% busy and every page feels instant.

On Friday at 10:00 a stadium tour goes on sale. Marketing expects 20Γ— the normal peak for the first ten minutes: 6,000 requests per second.

Predict first: what goes wrong, and in what order?

Check your answer

6,000 requests Γ— 8 ms = 48 cores of work every second, on a machine with 16. The box can finish at most 16 / 0.008 = 2,000 requests per second, so every second another 4,000 requests join the queue. Ten seconds in, 40,000 requests are waiting, and a new arrival sits about 20 seconds behind them. Clients time out and retry, which adds load to a machine that is already drowning.

Fixing the arithmetic raises three more questions:

  • Rent a 64-core machine and the math works (48 of 64 cores, 75% busy). But it is still one machine. If it dies at 10:02, the sale stops during its most valuable ten minutes.
  • At 6,000 requests per second, "fast" needs a number. Is a 3-second page acceptable for 1 fan in 100?
  • Two fans click seat 14F at the same instant. Who gets it, and could both?

Those questions are the four ideas of this lesson. Every system design interview comes back to them, whatever the product:

# The interviewer asks Concept The move you'll learn
1 "What happens at 20Γ— traffic?" Scalability Scale up to buy time, scale out a stateless tier, then relieve the data tier
2 "How fast is it for the slowest users?" Latency, throughput, percentiles Size with Little's law, set targets on p99, keep headroom below the knee
3 "What if that machine dies at 10:02?" Availability Remove single points of failure, do the nines math, make recovery fast
4 "Can two fans get seat 14F?" Consistency, CAP, PACELC Choose a consistency level per piece of data and name its price

This lesson assumes you have read Foundations of System Design, which follows one request from DNS to the database and back and covers the latency numbers used in estimates. Here you learn the concepts every later lesson leans on. Each idea ends with the sentence or two you would actually say in an interview, and the follow-up question you should expect.

Before any move: turn adjectives into numbers

"Fast", "scalable", "highly available" and "consistent" are adjectives. An interviewer cannot check an adjective, and neither can you. Each concept in this lesson starts to do work the moment you attach a number and a scope.

Vague Precise The follow-up you'll get
"It should be fast." "Seat-map reads: p99 under 300 ms during the 6,000 req/s on-sale peak, measured at the load balancer." "What is your p99 today?"
"It should scale." "Handle 20Γ— the normal peak for ten minutes, then shrink back." "Which tier breaks first?"
"It should always be up." "99.95% of purchase requests succeed each month." "What is your single point of failure?"
"The data should be consistent." "A seat is sold at most once; the seat map may lag by up to 2 s." "What happens during a network partition?"

The precise versions are not longer for show. Each number decides something. A 300 ms p99 caps how many network hops one request can afford. 99.95% leaves the equivalent of 21.6 minutes of full outage per month. "Sold at most once" forces one kind of write, and "may lag by 2 s" permits a cheaper kind of read. Requirements Gathering shows how to get these numbers out of an interviewer; this lesson teaches what they mean.

πŸ’‘ Habit worth building: whenever you hear yourself say an adjective, finish the sentence with a number and a scope. "Fast: p99 under 200 ms for the feed endpoint."

Idea 1: Scalability β€” scale up, then scale out

Pressure: the load is growing past what one machine can do.

Slow is not the same as "doesn't scale"

Before you add machines, find out which problem you have. Two symptoms look alike on a dashboard and need opposite fixes:

Symptom Diagnosis Fix
Slow at 10 req/s, and just as slow at 1,000 Performance problem: each request does too much work Fix the work: the query, the N+1 loop, the algorithm
Fast at 10 req/s, slow only near peak Capacity problem: requests wait in a queue Add capacity, or remove the bottleneck

Scalability is how cheaply you can add capacity. Ideally, twice the machines give twice the throughput at the same latency. A system can scale well and still be slow (every request makes three slow database calls, on any number of servers). It can also be fast and scale badly (one finely tuned machine that falls over at 2,001 requests per second). Adding servers to a performance problem just serves more slow requests in parallel.

Why does latency climb suddenly near peak instead of growing steadily? Queueing. Take one worker that needs 10 ms per request on average, with requests arriving at random and service times varying at random (the simplest queueing model, called M/M/1). Its mean response time is 10 ms / (1 - utilization):

Utilization Arrivals per second Mean response time
50% 50 20 ms
70% 70 33 ms
80% 80 50 ms
90% 90 100 ms
95% 95 200 ms
99% 99 1,000 ms

The work per request never changed. Going from 50% to 99% busy made responses 50Γ— slower, and every extra millisecond was spent waiting in line. Real systems with many workers bend later and more sharply, but they bend. This is why capacity plans aim for peak utilization well below 100% (60–70% is a common rule of thumb). It is also why "slow only at peak" is a capacity problem, not a performance one.

Scale up: the bigger box

Vertical scaling (scaling up) means a bigger machine: more cores, more memory, faster disks. It needs no code changes and keeps the data next to the logic, with no network in between. Cloud providers rent single machines with hundreds of cores and terabytes of memory, and within one instance family the price usually grows roughly in step with size. For a database, scaling up is often the right first move.

What you give up:

  • A ceiling. Eventually there is no bigger machine, and the largest sizes are specialised and scarce.
  • One failure domain. A bigger box is still one box. Its crash, its kernel patch and its resize are all your outages.
  • Downtime to grow. Resizing usually means a restart, or a failover to a new machine.

Scale out: more boxes, and where the state goes

Horizontal scaling (scaling out) means more machines sharing the work behind a load balancer. The ceiling is much higher, and one failure removes a fraction of the capacity instead of all of it. The price is new machinery (a load balancer, health checks, deploys across a fleet) and one hard requirement: a request must not depend on landing on a particular machine.

That requirement is what stateless means. A stateless server keeps nothing in its own memory or disk that the next request needs. The state has not disappeared; it has moved somewhere every server can reach.

STATEFUL: the session lives in one server's RAM

  client ──► load balancer ──► app-1   [session: alice, 2 seats held]
                          └──► app-2   [has never heard of alice]

STATELESS: the session lives outside the servers

  client ──► load balancer ──► app-1 ──┐
                          β”œβ”€β”€β–Ί app-2 ──┼──► shared session store
                          └──► app-3 β”€β”€β”˜

There are three common ways to handle sessions once you scale out:

Option How it works What you give up
Sticky sessions The load balancer pins each user to one server (usually with a cookie); the session stays in that server's RAM A crash or scale-in logs that server's users out; load gets uneven; deploys must wait for users to drain
Shared session store Sessions live in Redis or a database; any server looks them up by session ID One more network hop per request (well under a millisecond inside a data centre), and the store must itself be highly available
Signed token The client carries a signed token (for example a JWT) holding the user ID and an expiry; servers verify the signature, no lookup Hard to revoke before it expires, so keep lifetimes short; the token rides on every request, so keep it small

⚠️ Whichever you choose, get the user's identity from the session token, never from a user ID the client sends. GET /cart?user_id=42 lets anyone read cart 42.

Two precise points interviewers like to probe:

  • A stateful process is not a stateful server. A checkout has state ("card charged, tickets not yet issued"). If each step is written to a database or a workflow engine, any worker can resume it, and the servers stay stateless.
  • Some tiers are stateful on purpose. A WebSocket gateway holds open connections. A game server holds a live match in memory because it updates it 30 times a second. You scale such tiers by partitioning with affinity: route everything for match 812 to the server that owns it, and decide what happens when that server dies (snapshot the state, replicate it, or accept losing it). Statelessness makes scaling easiest, but it is not the only way.

The bottleneck moves to the data

Back to the on-sale. Put the database on its own machine and scale out the app. Each app server has 8 cores, and the app's share is 5 ms of CPU per request, so one server tops out at 8 / 0.005 = 1,600 req/s. Plan for 70% at peak (1,120 req/s each): 6,000 / 1,120 = 5.4, so 6 servers, plus 1 spare so that losing a server at peak does not overload the rest. Seven small stateless servers, as far as CPU is concerned; Idea 2 checks whether they also have enough threads.

Now look at the database. It still receives every request: 6,000 Γ— 3 ms = 18 cores of work, on a 16-core machine. The same overload has moved one tier down. Suppose 90% of requests are seat-map reads. Read replicas or a cache can absorb those 5,400 reads per second, which leaves 600 writes per second, or 1.8 cores, on the primary. Reads scale by copying data. Writes are harder, because every copy must agree. Databases & Storage covers replicas, sharding and caching in depth.

Your turn: a ten-times bigger artist announces a tour: 60,000 req/s at peak, still 10% writes. You run the app tier in three availability zones and must survive losing a whole zone at peak. How many app servers do you need, and what breaks next?

Check your answer

At 1,120 req/s per server, 60,000 req/s needs 54 servers. With 18 in each of three zones, losing a zone leaves 36 servers, each asked for 1,667 req/s, which is more than the 1,600 maximum. The two surviving zones must carry the whole peak, so each zone needs 27 servers: 81 in total (by CPU; Idea 2 adds a check on threads). At this size, "one spare" becomes "one spare zone".

What breaks next is the write path: 6,000 writes per second Γ— 3 ms = 18 cores on a 16-core primary. Your options are a bigger primary (scale up again), splitting the writes across several primaries (for example by event), or admitting buyers at a rate the primary can handle, with a waiting room. Each one has a cost; say which you pick and why.

Say it in the interview: "The app tier is stateless, with sessions in Redis, so I scale it out behind the load balancer: about 6 servers at 70% CPU plus a spare, each with enough worker threads for its share of the requests in flight. That moves the bottleneck to the database, so next I'd separate reads, which can go to replicas, from writes, which must go to the primary."

Likely follow-up: "Why not sticky sessions?" Because a server's crash or a scale-in logs out everyone pinned to it, load becomes uneven, and every deploy turns into a slow drain.

Idea 2: Latency, throughput and the tail

Pressure: "fast" has to hold for everyone, including the unlucky.

Three different words

  • Latency (or response time) is the time from sending a request to receiving the whole response, as the caller sees it. It includes network time, queueing and processing. The processing part alone is the service time. When latency grows but service time does not, requests are waiting in line.
  • Throughput is completed work per second: requests, messages or bytes.
  • Bandwidth is the capacity of a link, usually in bits per second, and throughput can never exceed it. A 1 Gbit/s link moves at most 125 MB/s, so it can deliver at most 250 responses of 500 KB per second, however many CPUs sit behind it.

Latency and throughput are related, but they are separate levers:

Change Latency per request Capacity (max throughput)
Add identical servers behind the load balancer Unchanged at low load; lower near peak (less queueing) Up
Batch writes (wait a few ms to fill each batch) Up Up
Cache a slow lookup Down Up
Run the same servers at 95% busy instead of 50% Up, sharply (the knee) Unchanged: you use more of it and keep nothing for bursts

The throughput you actually deliver is utilization Γ— capacity. It rises with demand until it reaches capacity, then stays flat while latency grows without bound.

Little's law: the formula to know by heart

For any stable system (one server, a connection pool, a queue, a whole data centre), averaged over time:

L = Ξ» Γ— W β€” the average number of items inside = arrival rate Γ— average time each one spends inside.

It makes no assumption about how arrivals or service times are distributed, which is why it survives the back of an envelope. Three uses:

  1. Concurrency. The on-sale's 6,000 req/s with a 50 ms average latency means 6,000 Γ— 0.05 = 300 requests in flight on average. Idea 1 sized the app tier by CPU: seven servers. If each runs 40 worker threads and a request holds a thread for its whole life, those seven hold only 280 requests at once. Threads run out before CPU does: you would need 8 servers before any headroom, and 12 with 70% headroom and a spare. Each request spends only 5 ms of its 50 ms on CPU and the rest waiting for the database, so the cheaper fix is more threads per server (or asynchronous I/O): with 100 threads each, the same seven servers hold 700 requests. Size a tier by whichever limit it hits first.
  2. The ceiling of a pool. A database connection pool of 20 connections, where each query holds a connection for 4 ms, can serve at most 20 / 0.004 = 5,000 queries per second. If a slow query makes each connection busy for 40 ms, the ceiling drops to 500 queries per second, and every other request queues behind it. This is how a slow dependency takes down a healthy service: it does not fail, it just holds on to your concurrency.
  3. Waiting time. Rearranged, W = L / Ξ». A backlog of 12,000 messages drained at 400 per second means a new message waits about 30 seconds.

Your turn: your API's database pool has 50 connections, and each request runs one query. A normal query holds a connection for 5 ms. A bad deploy changes 10% of the queries so that they hold a connection for 200 ms. The API needs 3,000 queries per second. What happens, and to whom?

Check your answer

The average hold time becomes 0.9 Γ— 5 ms + 0.1 Γ— 200 ms = 24.5 ms. By Little's law the pool can now complete at most 50 / 0.0245 β‰ˆ 2,040 queries per second, against 3,000 needed (before the deploy the ceiling was 10,000). The pool saturates and a queue builds in front of it. Everyone waits, including the 90% of requests whose own query still takes 5 ms. The fixes: a timeout on the slow query, a separate small pool for it (a bulkhead) so it cannot starve the rest, and then fixing or rolling back the query.

Percentiles: the average is hiding something

Run this simulation of a service where 97% of requests take the usual path and 3% hit a slow one (a garbage-collection pause, a cold cache, a lock wait):

import random, statistics

random.seed(7)

def one_request():
    if random.random() < 0.03:                  # 3%: slow path (GC pause, cold cache, lock wait)
        return random.uniform(400, 1200)
    return random.lognormvariate(3.6, 0.35)     # the usual path, median about 37 ms

def pct(samples, p):
    s = sorted(samples)
    return s[min(len(s) - 1, int(p / 100 * len(s)))]

lat = [one_request() for _ in range(100_000)]
print(f"mean  {statistics.mean(lat):5.0f} ms")
for p in (50, 95, 99, 99.9):
    print(f"p{p:<4} {pct(lat, p):5.0f} ms")
mean     61 ms
p50      37 ms
p95      75 ms
p99     922 ms
p99.9  1172 ms

The p99 is the latency that 99% of requests beat. Here the mean (61 ms) looks healthy and the median (37 ms) looks great, yet 1 request in 100 takes almost a second. A few thousand slow values averaged with ninety-seven thousand fast ones disappear into a comfortable number. That is why latency targets are written on percentiles: "p99 under 300 ms", never "average under 300 ms". Google's SRE book makes the same argument.

Two traps with percentiles:

  • They do not average. A fleet's p99 is not the mean of each server's p99, weighted or not; it depends on the whole latency distribution on every server. Merge the raw histograms first, then take the percentile.
  • Your most active users live in the tail. A page that makes 40 requests hits at least one of your slowest 1% on about a third of its loads.

Tail at scale: fan-out multiplies the tail

Predict first: a request fans out to 100 servers in parallel and must wait for all of them. Each server is slow on 1% of requests. Roughly what share of requests is slow: 1%, 10%, or more than half?

Check your answer

More than half. The request is as slow as its slowest server. If each server is slow independently with probability 1%, the chance that at least one is slow is 1 - 0.99**N:

Fan-out N Requests that hit at least one slow server
1 1%
10 9.6%
100 63%

Jeffrey Dean and Luiz AndrΓ© Barroso described exactly this in The Tail at Scale (Communications of the ACM, 2013), with the same 63%. Their remedies are now standard interview moves:

  • Hedged requests. Send the request to one replica. If it has not answered by about its 95th-percentile latency, send a copy to a second replica and use whichever answers first. Waiting until p95 limits the extra load to about 5% of requests. In the paper's Bigtable benchmark (1,000 keys spread across 100 servers), hedging after 10 ms cut the 99.9th-percentile latency from 1,800 ms to 74 ms for 2% more requests. Hedge only reads and other operations that are safe to repeat.
  • Less fan-out, where the product allows it.
  • Timeouts and partial results, so one straggler cannot hold the whole page.

Say it in the interview: "I'll set the target on p99, say 300 ms at the load balancer, because the average hides the slowest 1%. This call fans out to 30 shards, so the tail compounds; I'd hedge to a second replica after the p95 latency and put a timeout on each shard."

Likely follow-up: "What does hedging cost?" Extra load roughly equal to the share of requests that cross the hedge threshold (about 5% at p95), and it is only safe for idempotent operations.

Idea 3: Availability and the nines

Pressure: machines, networks and deploys fail, and the sale must go on.

What availability means

Availability is the share of time, or of requests, in which the service does its job correctly. Google's SRE book gives two ways to measure it (Embracing Risk):

  • Time-based: uptime / (uptime + downtime).
  • Request-based: successful requests / total requests. Most distributed services use this one, because they are rarely "all up" or "all down".

Three terms that are often mixed up:

  • SLI (indicator): the measurement, such as "share of purchase requests that returned a success in under 1 s".
  • SLO (objective): the target for that measurement, such as "99.95% over 30 days".
  • SLA (agreement): a contract with consequences, such as refunds, if the target is missed.

Keep the internal SLO stricter than the SLA, so you notice trouble before customers are owed money. AWS, for example, says S3 Standard is designed for 99.99% availability and backs it with a 99.9% SLA (storage classes).

Availability Downtime per 30-day month Downtime per year
99% 7.2 hours 3.65 days
99.9% 43.2 minutes 8.76 hours
99.95% 21.6 minutes 4.38 hours
99.99% 4.3 minutes 52.6 minutes
99.999% 26 seconds 5.3 minutes

Each extra nine is ten times less downtime, and around four nines humans stop being fast enough. 4.3 minutes a month is less than it takes to get paged, open a laptop and read a dashboard. From 99.99% up, detection and failover must be automatic.

The complement of the SLO is an error budget. 99.9% leaves 0.1% of requests free to fail. You can spend that budget on risky deploys, migrations and experiments, and when it runs out, you slow down.

Serial parts multiply their availability; redundant copies multiply their unavailability

Components that a request passes through in series must all work, so their availabilities multiply. For redundant copies where any one is enough, all of them must fail at once, so their unavailabilities multiply:

  • In series: A = A₁ Γ— Aβ‚‚ Γ— A₃ …
  • Redundant, any one of n is enough: A = 1 βˆ’ (1 βˆ’ A₁)ⁿ

Walk the ticketing path on a normal evening, when one app server can carry all 300 req/s: a load balancer at 99.99%, app servers at 99.5% each, and a database at 99.9%. That last figure is roughly what you get from a machine that fails about every 90 days and takes 2 hours to restore by hand: 2,160 h / (2,160 h + 2 h) = 99.91%.

Design Availability Downtime per 30-day month
One of each: load balancer β†’ 1 app server β†’ 1 database 99.39% 263 min
Three app servers, if any one can carry the load 99.89% 47.5 min
… plus a database standby with automated failover in 60 s 99.989% 4.7 min

(The Redis session store from Idea 1 is one more component in series. It is left out here to keep the table small, but count it in an interview: it needs a replica and failover too.)

Two lessons hide in that table. First, after every fix a different component becomes the weakest link: the app server, then the database, and now the load balancer, which accounts for 4.3 of the remaining 4.7 minutes. Second, the database improvement came from repair time, not from better hardware. Availability = MTBF / (MTBF + MTTR), where MTBF is the mean time between failures and MTTR the mean time to recover. Cutting MTTR from 2 hours to 60 seconds took the database from 99.91% to 99.999%.

⚠️ Where the math stops. The redundancy formula assumes copies fail independently. A bad deploy to all three app servers, a zone outage, an expired certificate or a shared configuration service fails every copy at once. Spread copies across zones and roll deploys out gradually. Otherwise three redundant servers fail together as easily as one.

⚠️ When you need k of n. The redundancy formula assumes any one copy can carry the load. During the on-sale you need 6 of the 7 app servers, so the tier is up only while at most one of them is down: 0.995⁷ + 7 Γ— 0.995⁢ Γ— 0.005 = 99.948%, about 22 minutes of downtime a month. That alone uses up the 99.95% target from the start of the lesson. One more spare (6 of 8) gives 99.9993%, under 20 seconds a month. At peak, spares buy availability as well as headroom.

Single points of failure and failover

A single point of failure (SPOF) is a component whose failure stops the service. Find them by walking the request path and asking "what if this dies?" at every box, including the load balancer, DNS, the session store, the database primary and the configuration service.

Redundancy is having spare copies. In active-active redundancy every copy serves traffic, like app servers behind a load balancer. In active-passive redundancy a standby waits and takes over on failure, the usual set-up for a database primary. Failover is that takeover: detecting the failure (health checks, heartbeats) and switching. Detection time plus switching time is your MTTR.

The load balancer needs the same treatment. A self-managed pair in one data centre shares a virtual IP: the standby watches the active one's heartbeats and claims the IP when they stop (keepalived with VRRP is a common tool). The two sit side by side, not in a chain, and either one can reach every app server. A VRRP address can only move within one network segment, so such a pair survives the loss of one load-balancer machine, not the loss of the whole site or zone. For that, use a managed cloud load balancer that spans zones, or one pair per zone behind DNS or anycast.

                              clients
                                 β”‚
                     virtual IP 203.0.113.10
                 β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
           [ LB active ] ◄──── heartbeats ────► [ LB standby ]
                 β”‚                               β”‚   the standby claims the IP
                 β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜   if heartbeats stop
               β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
           [ app-1 ]         [ app-2 ]         [ app-3 ]    at normal load, any one can serve
               β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                          [ DB primary ] ── async replication ──► [ DB standby ]
                                                                   promoted on failure

When you promise failover, name two numbers:

  • RTO (recovery time objective): how long until service is back.
  • RPO (recovery point objective): how much recent data you may lose, measured in time.

An asynchronous standby gives a short RTO but an RPO above zero. PostgreSQL's documentation notes that streaming replication is asynchronous by default, so transactions committed on the primary but not yet replicated are lost if it crashes, in proportion to the replication delay (PostgreSQL docs). A synchronous standby gives an RPO of zero, but every commit waits for it. And if a required synchronous standby crashes, the same docs warn that commits "may never be completed", which is why you name several standbys and require only some of them (ANY 1 (s1, s2)). That is the choice between consistency and availability, in miniature, inside one database. Idea 4 generalises it.

Two more pairs of words that sound alike:

  • Durability is not availability. Durability means acknowledged data is never lost; availability means you can use it right now. AWS designs S3 Standard for 99.999999999% durability and 99.99% availability: your photo is almost never lost, but reaching it can fail for a few minutes.
  • Replication is not backup. A replica applies your accidental DELETE FROM orders within milliseconds. Replicas protect you from machine failure; backups (point-in-time copies kept elsewhere) protect you from bad writes and bugs.

When a dependency gets sick: timeouts, retries, circuit breakers

Many outages start as a slow dependency, not a dead one (remember the connection pool). The standard toolkit:

  1. Timeouts on every network call, set from the dependency's measured percentiles (a little above its p99.9), so a stuck call gives its thread back.
  2. Retries with exponential backoff and jitter, and only for operations that are safe to repeat: idempotent ones, or ones protected by an idempotency key (Networking Basics covers idempotency). The random delay stops thousands of clients from retrying in lockstep.
  3. Retry at one layer, within a budget. Retries multiply. Google's SRE book (Addressing Cascading Failures) gives the example of three layers that each make up to 4 attempts: one user action becomes as many as 4 Γ— 4 Γ— 4 = 64 attempts on the database, exactly when the database is struggling. Cap retries per request and per server.
  4. Circuit breakers. After repeated failures, stop calling the dependency for a cool-down period and fail fast, or serve a fallback. Then let a limited number of trial calls through: Martin Fowler's description allows one, and libraries such as Resilience4j make it configurable (their default is 10). If the trials succeed, the circuit closes again; if they fail, it reopens.
  5. Load shedding and graceful degradation. When overloaded, reject extra work early (HTTP 429 or 503) instead of letting everything time out. Switch off non-essential features, such as recommendations, to keep the essential one, checkout, alive.
            failures reach the threshold
  CLOSED ──────────────────────────────► OPEN     fail fast; no calls go through
    β–²                                     β”‚
    β”‚ trial calls succeed                 β”‚ cool-down elapses
    β”‚                                     β–Ό
    └───────────────────────────────── HALF-OPEN  a limited number of trial calls go through
                                          β”‚
                                          └── a trial call fails ──► back to OPEN

Here is a retry helper with capped exponential backoff and "full jitter" (each delay is random between zero and the current backoff ceiling), one of the variants that AWS's backoff and jitter analysis found far better than backoff without jitter:

import random, time

def call_with_retries(op, attempts=3, base=0.1, cap=2.0):
    """Call op(). Retry timeouts with capped exponential backoff and full jitter.
    Only use this for operations that are safe to repeat."""
    for attempt in range(attempts):
        try:
            return op()
        except TimeoutError:
            if attempt == attempts - 1:
                raise                                  # out of attempts: fail, don't hang
            delay = min(cap, base * 2 ** attempt)      # ceiling doubles: 0.1 s, 0.2 s, 0.4 s ... capped at 2 s
            time.sleep(random.uniform(0, delay))       # full jitter spreads clients apart

calls = {"n": 0}
def flaky():
    calls["n"] += 1
    if calls["n"] < 3:
        raise TimeoutError
    return "ok"

print(call_with_retries(flaky), "after", calls["n"], "attempts")   # ok after 3 attempts

Your turn: your database primary fails about once a month. Automated failover works, but detection waits for ten missed heartbeats, so each failover takes 5 minutes. What is the database's availability, and what is the cheapest way to add a nine?

Check your answer

MTBF is 720 hours and MTTR is 5 minutes: 720 / (720 + 0.083) = 99.988%, about 5 minutes of downtime a month. Detecting in 30 seconds instead gives 99.9988%, about 30 seconds a month: ten times less downtime without a single new machine. The cost is more false alarms, since a shorter timeout mistakes slow moments for failures, and each needless failover can lose the unreplicated tail of writes.

Say it in the interview: "Components in series multiply, so for 99.95% end to end every hop has to beat that. The app tier runs in three zones behind a managed load balancer that spans them. The database primary is the weakest link, so I'd add a standby with automated failover. With an asynchronous standby we could lose the last second or so of writes on failover; if that's not acceptable, a synchronous standby costs a few milliseconds per commit."

Likely follow-up: "How do you stop both databases from acting as primary after a failover?" With majority-based leader election and fencing, covered in Advanced Topics & Final Prep.

Idea 4: Consistency, CAP and PACELC

Pressure: copies of the same data disagree for a while, and some data can't afford that.

The stale read

A fan buys seat 14F. The write commits on the primary. 40 ms later the app loads "My tickets" from a read replica that is 90 ms behind:

time β†’         0 ms                  40 ms                          90 ms
fan's app    POST /buy 14F ──►      GET /my-tickets ──┐
primary      commits 14F β†’ fan                        β”‚
replica                              not applied yet β—„β”˜ β†’ "no tickets"
                                                                     applies 14F β†’ fan

Nothing is broken. Replication to read replicas is usually asynchronous (PostgreSQL's default), and the replica is doing its job 90 ms late. MongoDB avoids this by default by reading from the primary, and its documentation warns that every other read preference may return stale data. Whether a stale read is a bug depends on what the data is, which is why "consistency" needs a precise name.

Consistency is a ladder, not a switch

Model The promise What you'd see without it
Linearizable (strong) Every operation appears to happen at a single instant between its request and its response. Once a write completes, every later read, by anyone, sees it or something newer. Fan A's purchase of 14F has completed, yet fan B, loading the page a second later, still sees 14F as free
Bounded staleness Reads may be stale, but never by more than a stated time (or number of versions) The seat map still shows a seat sold ten minutes ago
Causal Effects never appear before their causes, for anyone A reply shows up before the comment it answers
Read-your-writes You always see your own writes You buy a ticket and "My tickets" is empty
Monotonic reads You never see time go backwards A refresh shows 2 tickets, then 1, then 2
Eventual If writes stop, all replicas converge That is all: any read may be stale, with no stated bound

The order is rough, not a strict ladder. Bounded staleness limits how old a read can be, while causal consistency limits the order in which things appear. Read-your-writes and monotonic reads sit side by side, since neither implies the other; causal consistency gives both. Read-your-writes and monotonic reads are called session guarantees: promises to one client, which is often all a product needs. Cheap ways to get read-your-writes: read the user's own data from the primary for a few seconds after they write, return the new value in the write's response, or carry a version token and read only from a replica that has caught up to it. Azure Cosmos DB offers this ladder as settings, from strongest to weakest: strong, bounded staleness, session, consistent prefix and eventual (Cosmos DB docs).

⚠️ Three meanings of "consistency". ACID's C means a transaction preserves the database's rules (no negative balances, valid foreign keys). CAP's C means linearizability: reads see the latest write across replicas. And "the cache is consistent" usually means "not stale". A fully ACID database with asynchronous read replicas still serves stale reads from those replicas. Say which one you mean.

BASE ("basically available, soft state, eventually consistent") is a slogan for the available, eventually consistent end of this ladder, coined as a contrast to ACID. It is not a precise model. In an interview, name the actual guarantee instead.

Quorums: turning the dial per operation

Dynamo-style stores such as Cassandra copy each key to N replicas and let every operation choose how many replicas must answer: W for a write, R for a read. If R + W > N, every set of replicas a read asks overlaps every set a write reached, so a read always includes at least one replica holding the latest acknowledged write.

N = 3 replicas of "seat:14F"      write with W = 2, then read with R = 2

             after the write     the read asks A and C
replica A    v2                  β†’ v2   ◄─ the overlap: the newest version wins
replica B    v2
replica C    v1 (behind)         β†’ v1

With R = 1 and W = 1, reads and writes are fastest, and a read that lands on C returns v1. Cassandra calls these per-operation choices consistency levels (ONE, QUORUM, LOCAL_QUORUM, ALL) and states the same R + W > RF rule (Cassandra docs). DynamoDB reads are eventually consistent by default. A strongly consistent read (ConsistentRead=true) costs twice as much read capacity and is not available on global secondary indexes (DynamoDB docs).

⚠️ Overlap is not the whole story. Cassandra resolves conflicting writes by last write wins on timestamps, so clock skew can silently discard a write, and two concurrent writers still race. Quorums make stale reads rare; they do not make "read, check, then write" safe. For "sell 14F at most once" you need a conditional write on a single authority: UPDATE seats SET owner = 42 WHERE id = '14F' AND owner IS NULL on one primary (exactly one buyer sees "1 row updated"), a DynamoDB condition expression, or a lightweight transaction in Cassandra.

CAP, stated precisely

A network partition means some nodes cannot exchange messages: a cut link, a bad firewall rule, a failed switch. The nodes on each side are alive; they just cannot hear each other.

The CAP theorem, conjectured by Eric Brewer and proved by Seth Gilbert and Nancy Lynch in 2002, says: in a network that can lose messages, no system can provide a read/write store that is both consistent (linearizable) and available (every request to a node that has not crashed eventually gets a non-error response). Gilbert and Lynch restate it plainly in Perspectives on the CAP Theorem.

So while a partition lasts, a replica that cannot reach the others has two options:

  • Refuse or wait: keep C, and give up A for the clients it serves.
  • Answer from what it has: keep A, and risk stale or conflicting data.

What CAP does not say:

  • It does not offer "pick two of three". Brewer's own twelve-years-later retrospective calls that framing misleading. In a distributed system, partitions are not optional, so "CA" only means you haven't said what happens during one. When there is no partition, CAP allows both C and A.
  • Its A is not your SLO. CAP-available means every live node answers, however late. A system that rejects writes on the minority side of a rare partition can still meet 99.99%.
  • Its C is not ACID's C.
  • It says nothing about latency. That is PACELC's job.

Many real systems are neither strictly "CP" nor strictly "AP". A primary with asynchronous read replicas is not linearizable even without a partition, and its cut-off side cannot accept writes either. So label operations, not products: "purchases are CP, seat-map reads are AP".

PACELC: the trade-off you pay every day

Daniel Abadi's PACELC (IEEE Computer, 2012) completes the sentence: if there is a Partition, choose Availability or Consistency; Else, in normal operation, choose Latency or Consistency.

The "else" half matters more day to day. Partitions are rare, but latency is paid on every request, because consistency needs coordination and coordination costs round trips:

  • A synchronous replica makes every commit wait for it.
  • A quorum read must wait for two replicas instead of the first one to answer.
  • A strongly consistent write across regions waits at least one inter-region round trip, since a majority of regions must agree. Azure Cosmos DB documents its multi-region strong write latency as about two round trips between the two farthest regions plus 10 ms at p99, and blocks strong consistency by default for regions more than 8,000 km apart.

Abadi's paper classes Dynamo, Cassandra and Riak as PA/EL in their default configurations, and fully consistent systems such as VoltDB/H-Store and Megastore as PC/EC. Tunable stores move per operation. Cassandra at ONE behaves PA/EL; at QUORUM it pays latency for consistency. etcd serves linearizable reads by default, through its Raft consensus, and lets a client ask for "serializable" reads, which may be stale but skip the consensus round trip (etcd guarantees).

Choosing per piece of data

Data Consistency Why What it costs
Seat ownership Linearizable: a conditional write on one primary "Sold at most once" Every purchase goes to the primary; during a partition, the side without it cannot sell
Payment Strong, plus an idempotency key Never charge twice Latency; every retry must reuse the key
Seat map shown to browsers Bounded staleness, at most 2 s A stale "available" only costs one failed click Lag monitoring: a replica more than 2 s behind is pulled from rotation; some clicks fail at purchase time
"My tickets" page Read-your-writes A fan who can't see a purchase will buy again The buyer's reads go to the primary for a few seconds
"Interested" counter on the event page Eventual Nobody audits it Counts jump around

Failure modes

  • Asynchronous failover loses writes. Promote an asynchronous replica and the unreplicated tail is gone. If the old primary comes back, its extra writes conflict with the new history.
  • Split brain. A network cut makes the standby believe the primary died, and both accept writes. Prevent it with majority-based leader election (a third zone or a witness breaks ties) and fencing: the old primary must be cut off, by an expired lease or by storage that rejects its stale token, before the new one acts.
  • Last write wins with skewed clocks silently drops whichever write carries the older timestamp, even if it really happened later.
  • Time travel after failover. A client that read version 7 from the old primary can read version 6 from a replica of the new one (monotonic reads broken) unless it carries a version token.
  • The minority side stops. Split a three-node consensus group into a pair and a single node: the pair keeps or elects a leader and carries on, while the single node can neither commit writes nor serve linearizable reads. That is CP working as designed; say so out loud.

Your turn: a shopping site runs in two regions, and users mostly stay in the nearest one. "Add to cart" must keep working even when the regions are cut off from each other. Placing an order must happen exactly once. Which consistency do you choose for each, and what do you give up?

Check your answer

The cart is AP: each region accepts cart writes locally, and when the regions disagree you merge the versions (the union of items). Amazon's Dynamo paper (2007) describes this choice for its shopping cart: an "add to cart" is never lost, but a deleted item can reappear after a merge. That cost is acceptable for a cart.

Placing the order is CP: one leader (or a consensus group) records it, guarded by an idempotency key, so a retry cannot create a second order. The price is that a user cut off from the order's leader cannot check out until the partition heals, and users far from the leader pay a cross-region round trip at checkout.

Say it in the interview: "Seat purchases need linearizability, so they're a conditional write on the seat's primary. During a partition the side without the primary rejects purchases; that's a deliberate CP choice. The seat map reads from replicas, and any replica more than 2 s behind is pulled from rotation, so staleness stays under 2 s. The buyer's own tickets page reads from the primary for a few seconds after a purchase, which gives read-your-writes."

Likely follow-up: "What does that cost across regions?" With one primary region, remote buyers pay one inter-region round trip per purchase, while reads stay local.

Words that change meaning mid-interview

Interviewers listen for precision, and these pairs are where candidates slip. Each row is a place where one word means two things.

Word Precise meaning Often confused with
Partition (network) Nodes are alive but cannot exchange messages Partition (data): a shard, one slice of the data set
Consistency (CAP) Linearizable: reads see the latest completed write ACID's C (rules and invariants hold); "the cache isn't stale"
Availability (SLO) The share of time or requests served successfully CAP's A (every live node answers); cache hit rate, a different metric (load absorbed, not success rate)
Latency Time until the caller has the response, queueing included Service time (processing only); network round-trip time
Throughput Work completed per second Bandwidth: the link's capacity, in bits per second
Performance How fast one request is at low load Scalability: how cheaply capacity grows as you add resources
Replication Live copies of the same data Backup (a point-in-time copy that survives a bad write); sharding (different data on each node)
Durability Acknowledged data is not lost Availability: you can reach it right now
Fault tolerance Keeps working correctly through a component failure High availability (little downtime, maybe a short failover); graceful degradation (keeps running with fewer features)
Stateless server Keeps nothing the next request needs "No state anywhere": the state moved to a store or a token
Queue Each message goes to one of the competing consumers (sharing work) Pub/sub: every subscriber gets every message (broadcast); see Messaging & Queues

Final round: no label on the problem

Real prompts don't say which idea they test. Turn the adjectives into numbers, find the pressure, then pick the move and name its cost.

Challenge 1: the sneaker drop

500 pairs of a limited-edition sneaker go on sale at 12:00:00. You expect 200,000 people to press Buy within the first minute. Never oversell; the site must not fall over.

  1. Which data needs which consistency?
  2. Each purchase updates one stock row under a lock, which takes 2 ms. How many purchases per second can that row absorb?
  3. What do you do with 200,000 clicks?
Check your answer
  1. The stock count needs a linearizable update: a conditional decrement such as UPDATE items SET stock = stock - 1 WHERE id = 7 AND stock > 0 on one primary; "0 rows updated" means sold out. The product page and a "sold out" flag can be eventually consistent, a few seconds behind.
  2. One row serialises its writers: 1 / 0.002 s = 500 updates per second. All 500 pairs could be gone in about a second. The counter is not the real problem.
  3. The problem is the other 199,500 requests. Put a waiting room in front of checkout that admits buyers at a rate checkout can handle. Once stock reaches zero, publish "sold out" through a cache or CDN so the rest never reach the database (load shedding with a cheap cached answer). Payments carry idempotency keys, and a reservation expires if the buyer does not pay within a few minutes. The cost: some people see "available" and then "sold out", and fairness depends on the queue's order.

Challenge 2: the internal tool

An HR tool for a 300-person company. The interviewer asks: "How would you make it scale and be highly available?"

Check your answer

Size it to the requirement first. 300 users means a few requests per second at most, so one small server has plenty of headroom. Office-hours use makes 99.9% a generous target. A sensible design is two small app instances behind a managed load balancer, plus a managed PostgreSQL with automated failover and point-in-time backups: no sharding, no multi-region, no caching layer. Spend the effort where the risk is: backups (replication is not a backup) and access control, because HR data is sensitive. Then say what would change your answer, such as 100Γ— the users or an external SLA. Restraint that you can justify is a senior signal.

Challenge 3: feature flags in three regions

Every request in every service reads 20–50 feature flags, which adds up to millions of reads per second across three regions. Flags change a few times a day through an admin UI. A region may be cut off from the others for several minutes. A flag change should reach every service within a minute.

Check your answer

Reads are PA/EL. Each service keeps the flags in memory and refreshes them every 30 seconds or so from a replica in its own region. During a partition it keeps serving the last values it saw, and it ships with safe defaults for a cold start before any value has loaded. Writes are PC and rare: they go to one leader region, and an admin can retry if that region is unreachable. The price: a change takes up to the refresh interval to take effect, a cut-off region may run on old values until the partition heals, and an emergency kill switch is not instant. The one-line summary for the interviewer: "Reads are local and available; writes are consistent and rare."

Cheat sheet: pressure β†’ move β†’ cost

Pressure Move What it costs
One box is near capacity Scale up A ceiling, one failure domain, downtime to resize
Load keeps growing Scale out a stateless tier behind a load balancer A load balancer, fleet deploys; state must move out of the servers
Sessions live in server RAM A shared session store, or signed tokens An extra hop and an HA store, or revocation that is hard
The tier is stateful by nature Partition with affinity (route by room or match ID) Uneven load; state lost on a crash unless snapshotted
Slow even at low load Fix the work: query, N+1 loop, algorithm Engineering time; more servers won't help
Slow only at peak Add capacity; keep utilization at 60–70% Paying for idle headroom
Sizing threads, pools, servers Little's law, L = Ξ»W Uses averages; add headroom for bursts
Tail latency amplified by fan-out Fewer calls, hedged requests after p95, timeouts A few % extra load; hedge only safe operations
An availability target Remove SPOFs, redundancy across zones, automated failover Money, failover testing, split-brain risk
A slow or failing dependency Timeouts, retries with jitter at one layer, circuit breaker, load shedding Some requests fail fast or degrade
Must never double-sell Linearizable conditional write on one leader or a consensus group Latency; during a partition, the side without the leader cannot sell
Stale is fine for seconds Replica reads (eventual; monitor lag if a bound is required) Users may see old data; add read-your-writes where it matters
Global users, low-latency reads Local replicas (PA/EL) Staleness; conflicts if several regions accept writes

Before moving on, explain the ticketing site aloud as if to an interviewer, in three minutes and without notes: how each tier scales and which one breaks first; the p99 target, and how many servers the CPU and Little's law each call for; the availability target and which component limits it; and which data gets strong consistency, which gets a weaker guarantee, and what each side of a partition does. If you can also say what each choice costs, the design lessons ahead will feel like applications of these four ideas.

Next: Networking Basics for HTTP, idempotency, DNS, CDNs and TCP vs UDP; Core Building Blocks for load balancers and the stateless service tier; and Databases & Storage for replication, sharding and caching.