System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Building blocks: communication, reliability and observability


Two questions appear in almost every design: how do services talk to each other, and how does the client find out something changed. Once they are answered, two more follow in the wrap-up: what happens when those calls fail, and how would you know.

This lesson covers all four. It starts with communication patterns. Then the reliability patterns — the vocabulary that makes a design sound production-grade, each one existing because of a specific failure that happens constantly in distributed systems. It ends with observability and deployment, which answer the wrap-up's questions about knowing the system works and changing it safely.

Service-to-service: REST, RPC, GraphQL

REST (Representational State Transfer) models the system as resources addressed by Uniform Resource Locators (URLs) and manipulated with HTTP verbs. Ubiquitous, debuggable with a browser, and cacheable by HTTP infrastructure. Its weakness is over- and under-fetching: an endpoint returns a fixed shape, so a mobile client gets fields it does not need, or must make three calls to assemble one screen.

RPC (Remote Procedure Call), typically gRPC, models the system as functions you call. Binary encoding, a schema (protocol buffers) generating client and server code, built-in streaming, and noticeably lower latency and payload size than JSON over HTTP. Its weaknesses: browsers cannot speak it without a translating proxy, and a binary wire format is harder to inspect. This is the right default for internal service-to-service traffic in a performance-sensitive system.

GraphQL lets the client specify exactly which fields it wants in one request. It solves over-fetching and the "three calls per screen" problem, which is why it is common in mobile backends. Its costs are real: HTTP caching no longer works the same way because every query is a POST to one endpoint; a naive resolver implementation issues one database query per item (the N+1 problem); and a client can construct an arbitrarily expensive query, so you need query cost limits.

Client updates: the four options

The question "how does the client learn about a new message?" recurs in Sections 12, 14, 17 and 19. There are four answers.

Short polling. The client asks "anything new?" every N seconds.

  • Trivial to build; works everywhere; stateless server.
  • Wasteful and slow. With 1 million connected users polling every 5 seconds, that is 200,000 requests per second, the overwhelming majority returning "nothing". And the average notification is delayed by half the interval — 2.5 seconds.

Long polling. The client sends a request; the server holds it open until there is data or a timeout (say 30 seconds), then responds; the client immediately reconnects.

  • Near-real-time delivery over ordinary HTTP, through any proxy or firewall.
  • Each waiting client occupies a server-side connection, so the server must handle many idle connections cheaply. Reconnection churn every timeout adds overhead.

Server-sent events (SSE). One long-lived HTTP connection over which the server pushes a stream of messages. Built-in reconnection and event identifiers for resuming.

  • Simple, one-directional, HTTP-native, works with standard infrastructure.
  • Server to client only. The client sends anything back over a normal request.

WebSockets. A single TCP connection upgraded to a full-duplex protocol; either side sends at any time.

  • Lowest latency, bidirectional, low per-message overhead.
  • The server becomes stateful: a specific user's connection lives on a specific machine, so you need a way to find which machine (The stateful problem, Section 14), and a plan for when it dies.
  • Cost per connection is real. At roughly 10–50 KB of memory per idle connection, 1 million connections is 10–50 GB across the fleet, and a tuned server holds on the order of 100,000 connections — so about 10 machines for a million users, plus headroom for reconnection storms.
PatternLatencyServer costDirectionChoose when
Short pollingHalf the intervalVery high at scaleClient pullsUpdates are rare and infrequent, or the client is a script
Long pollingNear real-timeModerateServer pushes on one requestYou need push but cannot run WebSockets
Server-sent eventsReal-timeModerateServer → client onlyFeeds, notifications, live dashboards
WebSocketsReal-timeHigh (stateful)Both waysChat, collaborative editing, gaming
Short pollingClient asks every N seconds.goodtrivial to buildworks everywherecostswasted requests when nothing changedlatency up to N secondslow-frequency updates, small user baseLong pollingServer holds the request open until there isnews.goodnear-real-timeplain HTTPcostsa held connection per clienttimeouts and reconnect churnchat, when WebSocket is unavailableServer-sent eventsOne long-lived stream, server to client only.goodautomatic reconnectcheap on the servercostsone direction onlyconnection limits over HTTP/1.1feeds, notifications, live scoresWebSocketOne connection, both directions, stays open.goodlowest latencyfull duplexcostsstateful — complicates load balancingown heartbeat and reconnectchat, collaborative editing, tradingPick the cheapest option that meets the latency requirement — and say the requirement out loud first, because it is what makes the choice defensible.
The jump in cost is between polling and holding a connection — everything after that is a duplex question.

Reliability patterns: timeouts

Every one of those calls can fail, hang or arrive twice. The reliability patterns are what stop one such failure from spreading.

Each pattern answers one failureReliabilitypatternsTimeout: hung calleeBackoff: retry stormBreaker: cascadeBulkhead: pool starvedIdempotency: dup writeShedding: overload
Naming the failure each pattern defends against is what makes the vocabulary sound production-grade rather than recited.

Every network call gets an explicit timeout. Without one, a hung dependency holds a thread or connection until the operating system gives up, which can be minutes.

The arithmetic: a service with 200 worker threads calling a dependency that hangs for 30 seconds sustains 6.6 requests per second before every worker is blocked. The service is then down, even though the only actual failure was elsewhere.

Set timeouts from the latency distribution: roughly the 99th percentile plus headroom. And make them shorter as you go deeper — if the client gives up after 2 seconds, a backend call with a 5-second timeout is doing work nobody will receive.

Retries with exponential backoff and jitter

A retry converts a transient failure into a success. It also multiplies load at the worst possible moment.

  • Only retry idempotent operations. Retrying a non-idempotent write creates duplicates.
  • Exponential backoff: wait 100 ms, then 200, 400, 800 — giving the dependency room to recover.
  • Jitter: randomise each wait. Without it, every client that failed at the same instant retries at the same instant, producing a synchronised wave that re-breaks the recovering service.
  • Cap the attempts, and budget them across the stack. Three retries at each of three layers is 27 attempts for one user request — an amplification that turns a small failure into an outage.

Circuit breakers

A circuit breaker stops calling a dependency that is clearly broken, so callers fail fast instead of accumulating timeouts.

Three states: closed (calls pass through, failures counted), open (calls fail immediately without a network attempt, after the failure rate crosses a threshold — say 50% over 20 requests), and half-open (after a cooldown, a few trial requests are allowed; if they succeed the breaker closes, if not it reopens).

The gain is throughput: a caller failing in 1 ms instead of waiting 2 seconds keeps its threads free for requests it can serve.

Bulkheads

Isolate resources so one failure cannot consume everything, named after the compartments in a ship's hull. Give each downstream dependency its own connection pool or thread pool. When the recommendation service hangs, it exhausts its own pool of 20 connections rather than the shared pool of 200, and checkout keeps working.

The analogy has a limit worth stating: ship bulkheads are about containing water that has already entered, and they do not help if the compartments are connected above the waterline. Software bulkheads similarly fail if the "isolated" pools share one thread pool or one CPU-bound process underneath.

Idempotency keys

The client generates a unique key per logical operation and sends it with the request. The server records the key with the result, and a repeat of the same key returns the stored result instead of performing the operation again.

This is what makes retries safe on a write path, and it is the mechanism behind every "exactly-once" claim in this course (Message queues and event streaming; Sections 12, 23, 24, 28). The storage detail matters: the key must be recorded in the same transaction as the effect, or a crash between the two reintroduces the duplicate.

Graceful degradation and load shedding

When the system cannot serve everything, decide in advance what to drop.

  • Serve stale cached data rather than an error when the database is unavailable.
  • Turn off expensive non-essential features under load — personalised ranking falls back to chronological, recommendations return a static list.
  • Shed load deliberately: reject a fraction of requests quickly so the rest succeed. Rejecting 20% of traffic at the edge is far better than serving 100% of it 30 times slower and timing out.
  • Put the switches behind feature flags so a human can flip them during an incident.

Observability: the three signals, and what each is for

The wrap-up step (Step 4: wrapping up) asks how you would know the system is working and how you would change it safely. These are the answers.

Four questions, four signalsMetrics: is it broken?Logs: what happened?Traces: where is the time?SLO: broken enough to act?
Averages hide the one percent of users having the worst hour, which is why you alert on p99 and not the mean.

Metrics. Numbers aggregated over time — request rate, error rate, latency percentiles, queue depth, CPU. Cheap to store because they are pre-aggregated, so you can keep them for years, and they are what alerts fire on. They cannot tell you why: a metric says 3% of requests failed, not which ones or what the error was. Section 22 designs a metrics system.

Logs. A record per event, with detail. They answer "what happened to this specific request", which metrics cannot. They are expensive: at 10,000 requests per second and 1 KB per log line, that is 10 MB/s, 864 GB a day, before any application logging. This is why sampling and log levels exist.

Traces. One request followed across every service it touched, recorded as a tree of spans with timings. Traces answer "where did the 2 seconds go" in a system where the answer spans five services. Usually sampled — 1% of requests, plus 100% of errors — because full tracing is prohibitively expensive.

The rule of thumb: metrics tell you something is wrong, traces tell you where, logs tell you what. A design that names all three, with that division of labour, reads as operationally literate.

Percentiles, not averages

Report p50, p95, p99, and p99.9 — never only the mean. A service with a 50 ms mean and a 2-second p99 is failing 1% of requests badly, and at 15,000 requests per second that is 150 users every second having a bad experience while the average looks excellent.

Two related points worth making unprompted: a user request that fans out to 10 backend calls experiences the p99 of the slowest of them, so tail latency amplifies with fan-out; and averaging percentiles across servers is mathematically meaningless — aggregate the underlying distribution instead.

Health checks from the service's side

The load balancer's health checks only work if the service answers them honestly. Liveness: is the process functioning at all? Failure means restart it. Readiness: can it serve traffic right now? Failure means remove it from the load balancer but leave it running.

Conflating them is a common production bug: a service that is warming its cache reports unhealthy, gets restarted, warms again, and never becomes ready.

Keep health endpoints shallow. A health check that queries the database means one slow database marks the entire fleet unhealthy and takes the system down — a failure amplifier rather than a safety net.

Deployment strategies

StrategyMechanismCostRollback
RollingReplace instances a few at a timeCheap; both versions run togetherRoll forward or reverse the rollout
Blue-greenTwo full environments; switch traffic at onceDouble the infrastructure during the switchInstant — switch back
CanarySend a small share of traffic to the new version, increase if healthyNeeds good metrics and traffic splittingRoute the small share away

Canary is the default answer for a large system. Route 1% of traffic to the new version, watch error rate and latency for ten minutes against the old version as a control, then 5%, 25%, 100%, with an automatic abort if error rate rises. It limits the blast radius of a bad release to a fraction of users instead of all of them.

Two things to mention alongside: both versions run simultaneously in rolling and canary deployments, so database migrations must be backward-compatible — add a column, deploy code that writes both, backfill, then remove the old one, never all at once. And feature flags separate deploying code from enabling behaviour, which turns a rollback into a configuration change taking seconds rather than a redeploy taking minutes.

Service level objectives

An SLI (service level indicator) is a measurement — "the fraction of requests served under 200 ms". An SLO (service level objective) is the target — "99.9% of requests under 200 ms over 30 days". The gap between the target and 100% is the error budget: at 99.9% over 30 days you may fail about 43 minutes' worth of requests. Spent budget means stop shipping features and fix reliability; unspent budget means you can take more risk. Saying this in a wrap-up shows you have thought about reliability as a quantity rather than an aspiration.