Software Engineering WikiSE Wiki

System design

Latency and capacity numbers, scaling, caching, load balancing, resilience, data and API decisions, and the failure mode each one introduces.

Reviewed MarkdownEdit

On this page

Numbers to reason with#

OperationOrder of magnitude
L1 cache reference1 ns
Main memory reference100 ns
SSD random read100 µs
Datacentre round trip500 µs
Disk seek (spinning)10 ms
Sydney to US west coast round trip~150 ms
Read 1 MB sequentially from memory~10 µs
Read 1 MB sequentially from SSD~200 µs

Anything crossing a region pays the speed of light in fibre, and no tuning removes it. Design for one round trip per user action, not five.

QuantityReference point
1 request/second sustained2.6 million per month
1,000 rpsA single well-tuned service instance can often do this
99.9% availability43 minutes of downtime per month
99.99% availability4.3 minutes per month, which needs automated failover
1 TB at 1 Gbps~2.2 hours to transfer

Latency, throughput, saturation#

Throughput is work per unit time; latency is time per unit work; they trade against each other through queueing. As utilisation approaches 100%, queueing delay rises without bound. This is why a system at 80% CPU feels fine and the same system at 95% feels broken.

Measure percentiles, never means. An average of 100 ms with a p99 of 4 s means one request in a hundred is unacceptable, and every page composed of 50 calls hits it.

SignalWhy
Rate, errors, duration (RED)Per service, from the caller’s perspective
Utilisation, saturation, errors (USE)Per resource: CPU, memory, disk, network
Queue depthLeading indicator; latency is the lagging one

Scaling#

ApproachBuysCosts
VerticalSimplicity; no distributed reasoningA ceiling, and a single failure domain
HorizontalHeadroom and redundancyStatelessness, coordination, data partitioning
Read replicasRead capacityReplication lag and stale reads
ShardingWrite capacityCross-shard queries, rebalancing, hot keys
QueueingAbsorbs burstsLatency, ordering, and at-least-once delivery
CachingLatency and loadStaleness and invalidation

Scale the stateless tier first; it is the easy half. Almost every real limit is the database, and the fix is usually removing work rather than adding replicas.

Shard on a key with even distribution and no cross-shard queries in the hot path: user ID or tenant ID, rarely time. Time-based sharding puts every current write on one shard.

Caching#

PatternHow it worksFailure mode
Cache-asideApp reads cache, on miss loads and populatesThundering herd on expiry
Read-throughCache loads from the store itselfSame, hidden inside the cache
Write-throughWrite to cache and store togetherSlower writes, consistent reads
Write-behindWrite to cache, flush asynchronouslyData loss if the cache dies
Refresh-aheadRefresh before expiryWasted work on cold keys
key = "user:v2:7"       # version in the key: deploy invalidates without flushing
ttl = 300 + rand(0, 60) # jitter, or every key expires in the same second

Build three defences in from the start: jittered TTLs, a single-flight lock so one miss triggers one load, and a negative cache for “not found” so a missing key cannot be used to hammer the database.

Invalidation is the hard part. Prefer short TTLs and versioned keys over event-driven invalidation, which is correct in theory and wrong in the one code path nobody updated.

Load balancing#

LayerSeesCan do
L4 (TCP)Addresses and portsFast, protocol-agnostic, no retries or routing by path
L7 (HTTP)Method, path, headersRouting, retries, rewriting, per-route timeouts, TLS termination
AlgorithmUse
Round robinHomogeneous backends, uniform requests
Least connectionsVariable request duration
Least time / EWMAHeterogeneous backends; best default for HTTP
Consistent hashingCache affinity; minimises reshuffling when a node leaves
Random two choicesNearly as good as least-connections, far cheaper to coordinate

Health checks must exercise the dependency path that matters without cascading: a readiness check that queries the database takes every instance out at once when the database blips. Check liveness shallowly and readiness slightly deeper, and shed load rather than failing health checks under pressure.

Resilience#

MechanismPreventsWatch out for
TimeoutUnbounded waitsMust be shorter than the caller’s timeout
Retry with jittered backoffTransient failureRetry storms; only retry idempotent work
Circuit breakerHammering a dead dependencyHalf-open probes need to be cheap
BulkheadOne slow dependency exhausting the poolSizing each pool
Rate limitOverload and abusePer-tenant, not just global
Load sheddingTotal collapseShed cheaply, at the edge, with 429/503
Idempotency keyDuplicate side effects on retryStorage and expiry of the keys

Timeouts must shrink as you go deeper: if the edge allows 5 s, the service should allow 3 s and the database 1 s. Equal timeouts at every layer mean the whole chain waits for the slowest thing before anybody gives up.

Exponential backoff without jitter synchronises every client into a retry wave. sleep(random(0, min(cap, base * 2^attempt))) is the version that actually works.

Data#

QuestionPick
Relationships, transactions, ad-hoc queriesRelational, until proven otherwise
Known access pattern, huge volume, simple keysKey-value or wide-column
Full-text search and rankingA search engine, fed from the store of record
Time series with retention and rollupsA purpose-built TSDB
Events consumed by many independent readersA log (Kafka, Kinesis)

One store of record; everything else is a derived projection that can be rebuilt. The moment two systems both claim to be authoritative, reconciliation becomes a permanent tax.

Replication is asynchronous unless you paid for it not to be: a read straight after a write can miss it. Route read-after-write to the primary, or carry a version and wait for it.

CAP in practice: during a partition you choose between serving stale data and refusing requests. Decide per endpoint: a product page can be stale, a payment cannot.

APIs#

GET    /v1/orders?status=open&limit=50&cursor=eyJ...    200
POST   /v1/orders                                        201 + Location
GET    /v1/orders/{id}                                   200 | 404
PATCH  /v1/orders/{id}                                   200 | 409
DELETE /v1/orders/{id}                                   204
DecisionGuidance
VersioningIn the path (/v1): visible in logs, routable, obvious
PaginationCursor, not offset: stable under concurrent writes, and cheap at depth
ErrorsA consistent envelope with a machine-readable code and a human message
Partial failureReport per-item status in bulk endpoints; do not fail the batch
IdempotencyAccept an Idempotency-Key on POST and store the result
ConcurrencyETag plus If-Match, returning 409 on conflict
Rate limitingx-ratelimit-* headers and retry-after on 429
Long operationsReturn 202 with a status URL rather than holding the connection

Breaking changes need a new version; additive changes do not. A field that becomes required, an enum that gains a value clients must handle, or a default that changes are all breaking even when the schema still validates.

Asynchronous work#

Queues turn a synchronous dependency into a durable one, at the cost of eventual consistency and a new failure surface.

ConcernAnswer
DeliveryAt-least-once in practice; make consumers idempotent
OrderingOnly within a partition or key; design so global order is not needed
Poison messagesDead-letter queue with a retry budget and an alert on depth
BacklogAlert on age of the oldest message, not just depth
Fan-outTopic per event type; one queue per consumer group
Exactly-onceAchievable only as “at-least-once plus deduplication at the sink”

Multi-tenancy#

IsolationModel
Shared everythingTenant column on every row and every query; cheapest, riskiest
Shared database, schema per tenantBetter blast radius, harder migrations
Database per tenantClean isolation, operational overhead grows linearly
Cluster per tenantRegulatory or very large tenants only

Whatever the model, enforce the tenant boundary in one place, such as a row-level policy or a repository layer, never by remembering to add WHERE tenant_id = ? in each query.

Noisy neighbours are the default failure: per-tenant quotas and rate limits are not optional once more than one tenant matters.

API patterns in detail#

The table above says what to choose; the shapes below are what clients actually see. Consistency across endpoints matters more than any single choice, because a client library is written once against the pattern and then assumed everywhere.

{
  "error": {
    "code": "order_not_found",
    "message": "order 7f3a does not exist",
    "request_id": "01J9Z...",
    "details": [{ "field": "customer_id", "issue": "required" }]
  }
}

code is a stable string clients branch on; message is for humans and may change. request_id ties the response to the server-side trace. Use the same envelope for every 4xx and 5xx, and never leak stack traces or SQL through it. RFC 9457 (application/problem+json) is the standard form if you want one.

Cursor pagination encodes the position of the last item returned, usually the sort key plus a tiebreaker, opaque to the client:

GET /v1/orders?limit=50                     -> { "items": [...], "next_cursor": "eyJjcmVhdGVkIjoiMjAyNi0wMS0xNVQxMDoyMzo...In0" }
GET /v1/orders?limit=50&cursor=eyJjcmVh...  -> next page; the query is WHERE (created_at, id) < ($1, $2) ORDER BY created_at DESC, id DESC LIMIT 50

Offset pagination (?page=40) scans and discards 2,000 rows to return page 40 and skips or repeats items when rows are inserted between requests. A cursor is an index seek and is stable. Return next_cursor: null at the end rather than an empty page, so clients stop one round trip earlier.

Bulk endpoints accept a bounded array (document the limit, reject above it with 413 or 422) and report per-item results, because a batch of 500 with one bad item should not fail the 499:

{ "results": [ { "id": "a1", "status": 201 }, { "id": "a2", "status": 409, "error": { "code": "duplicate" } } ] }

Long-running work returns 202 Accepted with a Location pointing at a job resource the client polls (GET /v1/jobs/{id} returning { "state": "running" | "succeeded" | "failed", "result_url": ... }), with Retry-After on the 202 to set the poll interval. Webhooks invert that: the server calls the client, which needs a signature (HMAC of the body with a shared secret, timestamped to stop replays), retries with backoff on non-2xx, and a client that treats every delivery as possibly duplicated.

Field naming, time (RFC 3339 in UTC, always with the offset), money (integer minor units plus a currency code, never a float), identifiers (opaque strings, not sequential integers that leak volume) and null versus absent (PATCH must distinguish “clear this field” from “do not touch it”) are the decisions that cost most to change later. Write them down once in an API guideline and lint against it.

Idempotency#

An operation is idempotent when doing it twice has the same effect as doing it once. GET, PUT and DELETE are idempotent by definition; POST is not, and every retry, redelivery and double-click turns a non-idempotent POST into a duplicate charge, order or email. The fix is a key the client generates once per logical operation and sends with every attempt.

POST /v1/payments
Idempotency-Key: 4f0c1a2e-...       # UUID chosen by the client, reused on retry

server:
  1. look up key in the idempotency store (Redis or a table with a unique index)
  2. found and finished  -> return the stored status and body, do nothing else
  3. found and in flight -> 409 Conflict (or wait briefly), so two concurrent attempts cannot both run
  4. not found           -> insert key as in-flight, do the work, store (status, body), return

The store must be written atomically with the side effect, or a crash between “charge succeeded” and “record result” leaves a key that replays the charge. In a relational database, do both in one transaction with the key in a table that has a unique constraint; the constraint is the lock. Scope keys per client or tenant so two clients cannot collide, expire them after 24 hours or so, and reject a reuse with a different request body as 422, because that is a client bug rather than a retry.

Consumers of a queue need the same idea: a message ID or a business key (order_id plus event_type) checked against a processed-set before acting, or an operation designed so repetition is harmless (SET balance = 100 rather than SET balance = balance + 10, INSERT ... ON CONFLICT DO NOTHING, an upsert keyed by the event’s own ID). Idempotent consumers are what make at-least-once delivery safe; nothing else does.

Retries and timeouts#

A retry is a bet that the failure was transient. It pays off for connection resets, 503s and timeouts against a dependency that is healthy but briefly busy; it makes everything worse for a dependency that is overloaded, because every client now sends two or three requests instead of one. Three rules keep retries safe: retry only idempotent operations, cap the total retries a caller may add to the system, and do it in one layer.

delay(attempt) = random(0, min(cap, base * 2^attempt))    # full jitter; base 100 ms, cap 10 s
retry budget: at most 10% extra requests per minute across the service, then fail fast
retry on: connection error, 408, 429 (honouring Retry-After), 502, 503, 504, request timeout
never retry: 400, 401, 403, 404, 409, 422, or any 5xx where the request was not idempotent

A retry budget (as in Envoy and Linkerd) is a ratio rather than a per-request count: when the dependency is genuinely down, retries stop after a percentage of traffic instead of tripling the load. A circuit breaker gets the same outcome by watching the error rate and opening after a threshold; the two combine well. Stack retries in one place, at the service mesh or the client library, not both, or three retries at each of three layers become 27 attempts.

Hedged requests send a second copy of a read to another replica when the first has not answered within the p95 latency, then take whichever returns first. They cut tail latency at the cost of a few percent extra load and are only for idempotent reads.

Timeouts need a budget that shrinks along the call path. Propagate the remaining deadline in a header (grpc-timeout, or your own X-Request-Deadline) so a downstream service can give up on work the caller has already abandoned. A request that keeps running after its caller timed out is pure waste, and under load it is the waste that tips the system over. Distinguish connect timeout (short, a few hundred milliseconds inside a datacentre) from request timeout (per operation, based on the p99 plus headroom) and idle timeout (for keep-alive connections).

Queues in practice#

A queue decouples the producer’s availability from the consumer’s, absorbs bursts, and lets you scale consumers independently. It also introduces the questions of what happens when a message cannot be processed, when a consumer crashes mid-message, and when the producer’s write to the database and its publish to the queue disagree.

The last one is the transactional outbox: write the event to an outbox table in the same transaction as the business change, then a relay process reads the outbox and publishes, marking rows sent. The database transaction guarantees the event exists if and only if the change committed; the relay guarantees at-least-once publication. Change data capture (Debezium, Postgres logical replication) does the relay without application code.

BEGIN;
INSERT INTO orders (...) VALUES (...);
INSERT INTO outbox (aggregate_id, event_type, payload) VALUES ($1, 'order.created', $2);
COMMIT;
-- relay: SELECT ... FROM outbox WHERE sent_at IS NULL ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED;

Consumers acknowledge after the work is durable, not before, so a crash redelivers rather than loses. Set a visibility timeout or ack deadline longer than the slowest legitimate processing time, and extend it during long work. Failures go to a dead-letter queue after a bounded number of attempts, with the original message, the error and the attempt count, and someone must own that queue: alert on its depth and have a replay tool. Ordering is only guaranteed within a partition or FIFO group, so choose a partition key (order ID, tenant ID) that puts the messages that must be ordered together and spreads everything else.

Competing consumers scale horizontally; one partition per consumer is the ceiling for ordered streams, so the partition count decided at creation is the maximum parallelism. Monitor consumer lag (messages behind the head, and the age of the oldest unprocessed message) rather than raw depth. A queue that is always empty is a synchronous call with extra steps; one that is never empty is a backlog whose growth rate tells you when it will become an incident. See Kafka for the log-based version of these decisions.

Schema migration#

A schema change ships with old code still running against it, because a rolling deploy runs both versions for minutes and a rollback runs the old version for hours. Every migration therefore has to be compatible with the code before and after it, which turns most changes into an expand-and-contract sequence.

ChangeSafe sequence
Add a columnAdd it nullable or with a default; deploy code that writes it; backfill in batches; then add NOT NULL if needed
Remove a columnDeploy code that stops reading it; deploy code that stops writing it; then drop it in a later release
Rename a columnAdd the new column; dual-write; backfill; switch reads; stop writing the old; drop it
Change a typeAdd a new column of the new type and treat it as a rename
Add an indexCREATE INDEX CONCURRENTLY (PostgreSQL); online DDL or gh-ost / pt-online-schema-change (MySQL)
Add a constraintNOT VALID first, then VALIDATE CONSTRAINT (PostgreSQL), so the check does not lock the table

Long-held locks are the operational danger: an ALTER TABLE that rewrites a large table or waits behind a long transaction blocks every query on it. Set lock_timeout (a few seconds) in the migration session so it fails fast instead of queueing the whole application behind it, and retry. Backfills run in batches of a few thousand rows with a pause between, keyed by primary key range, and are restartable. Keep migrations forward-only and small; a “down” migration that drops a column destroys data and is rarely tested, so plan rollback as “deploy the previous code, which tolerates the new schema” instead. Migration tooling (Flyway, Liquibase, Alembic, golang-migrate, sqitch) should run as a separate step before the deploy, with the version recorded in the database, and CI should run every migration against a copy of production-shaped data to catch the lock and the hour-long backfill before they reach production. See PostgreSQL for lock behaviour and CONCURRENTLY.

Observability#

Observability is the ability to ask a new question about the system’s behaviour without deploying new code. It rests on three signals with different strengths: metrics are cheap aggregates for alerting and dashboards, logs are discrete events with detail, and traces show one request’s path across services with timing. The join key between them is a trace ID that appears in every log line and as an exemplar on metrics.

metrics:  http_server_request_duration_seconds{route="/v1/orders", method="GET", status="2xx"}   # a histogram; alert on rate and p99
logs:     {"ts":"...","level":"error","trace_id":"4bf9...","route":"/v1/orders","msg":"db timeout","tenant":"t-42"}
traces:   GET /v1/orders (312 ms) -> SELECT orders (280 ms) -> redis GET (2 ms)

Design the signals with the system. Label metrics by what you will aggregate and alert on (route template, not raw path; status class, not code) and never by unbounded values (user ID, request ID), because each label combination is a time series. Log at the boundaries (request in, dependency call, request out) with structured fields rather than sentences, and make the request ID and tenant part of every line. Propagate trace context (traceparent, W3C) through every hop including queues, or the trace ends at the first async boundary; sample at the edge so a whole trace is kept or dropped together.

Alerts should page on symptoms users feel (error rate, latency SLO burn, backlog age), not on causes (CPU, one pod restarting). An SLO turns “is it slow” into a number: 99.9% of requests under 300 ms over 30 days, with a burn-rate alert when the error budget is being spent faster than it would last the window. Every alert needs a runbook link and an owner; one that fires without an action gets muted, and then the real incident is muted with it. Dashboards follow the request path: edge, service, dependencies, each with rate, errors and duration, so an operator can walk down until the numbers change.

Cardinality and cost are the constraints. A histogram with ten labels of ten values each is 100,000 series times its buckets. Logs at debug level in a hot path cost more than the request they describe. Decide what you keep for how long (metrics for a year at low resolution, logs for weeks, traces sampled for days) and enforce it at ingestion. The concrete tooling is in Prometheus, Grafana and OpenTelemetry.

Troubleshooting design symptoms#

SymptomLikely causeCheck
p99 latency climbs while the mean is flatQueueing near saturation, or one slow dependency in a fan-outUtilisation and queue depth per resource; per-dependency latency percentiles
Latency jumps sharply at a steady traffic increaseA resource past roughly 80% utilisationUSE metrics on CPU, connection pools, disk and the database
Database load spikes at the same moment every few minutesSynchronised cache expiry (thundering herd)TTL jitter, single-flight on miss
Error rate rises after a dependency blip and stays highRetry storm, or retries without jitterRetry counts per caller; backoff and retry budget settings
Every instance leaves the load balancer togetherReadiness check depends on a shared dependencyWhat the readiness endpoint calls
A user cannot see a change they just savedRead served by a lagging replicaReplication lag; route read-after-write to the primary
One shard or partition is hotPoor shard key, or a single large tenantPer-key request distribution
Duplicate orders, emails or chargesAt-least-once delivery or client retries without idempotencyIdempotency keys at the consumer or API
Queue depth is low but work is hours lateA stuck or poison message blocks a partitionAge of the oldest message, dead-letter queue
Deploy blocks every query for minutesA migration waiting on a lock behind a long transactionlock_timeout in the migration; pg_stat_activity for the blocker
Rollback fails with a schema errorMigration was not compatible with the previous codeExpand-and-contract; test the old release against the new schema in CI
Events missing from the queue after a crashPublish happened outside the database transactionTransactional outbox or CDC
Same event processed twice with different outcomesConsumer is not idempotentProcessed-set keyed by event ID, or an upsert
Client sees 409 on every retry of a POSTIdempotency key still marked in flight after a crashExpire in-flight keys; write key and side effect in one transaction
Tail latency doubles when one replica is slowFan-out without hedging or per-replica timeoutsHedge idempotent reads at p95; per-call deadlines
Metrics store runs out of memoryUnbounded label valuesSeries count per metric; drop the label or aggregate at ingestion
Trace stops at the queueContext not propagated in message headersInject and extract traceparent on publish and consume
Alert fires, nobody knows what to doCause-based alert without a runbookRewrite as an SLO burn alert with a runbook link

Metrics for most of these checks come from Prometheus or traces from OpenTelemetry. Test failure paths with the approaches in Testing.

Reviewing a design#

Ask these in order, and stop when an answer is missing:

  1. What is the request rate, the data volume and the growth rate?
  2. Which parts must be strongly consistent, and which may be stale?
  3. What happens when each dependency is slow, then when it is down?
  4. Where is the state, and how is it restored after loss?
  5. What is the blast radius of one bad deploy, one bad tenant, one bad key?
  6. How would an operator detect this failing, and what would they do?
  7. What is the rollback, and has anyone run it?