# System design

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

Canonical: https://www.wiki.jodisand.me/design/
Reviewed: 2026-09-24
Related: [HTTP and curl](https://www.wiki.jodisand.me/http/index.md), [Redis](https://www.wiki.jodisand.me/redis/index.md), [PostgreSQL](https://www.wiki.jodisand.me/postgresql/index.md), [Prometheus](https://www.wiki.jodisand.me/prometheus/index.md), [Testing](https://www.wiki.jodisand.me/testing/index.md)


## Numbers to reason with

| Operation | Order of magnitude |
| --- | --- |
| L1 cache reference | 1 ns |
| Main memory reference | 100 ns |
| SSD random read | 100 µs |
| Datacentre round trip | 500 µ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.

| Quantity | Reference point |
| --- | --- |
| 1 request/second sustained | 2.6 million per month |
| 1,000 rps | A single well-tuned service instance can often do this |
| 99.9% availability | 43 minutes of downtime per month |
| 99.99% availability | 4.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.

| Signal | Why |
| --- | --- |
| Rate, errors, duration (RED) | Per service, from the caller's perspective |
| Utilisation, saturation, errors (USE) | Per resource: CPU, memory, disk, network |
| Queue depth | Leading indicator; latency is the lagging one |

## Scaling

| Approach | Buys | Costs |
| --- | --- | --- |
| Vertical | Simplicity; no distributed reasoning | A ceiling, and a single failure domain |
| Horizontal | Headroom and redundancy | Statelessness, coordination, data partitioning |
| Read replicas | Read capacity | Replication lag and stale reads |
| Sharding | Write capacity | Cross-shard queries, rebalancing, hot keys |
| Queueing | Absorbs bursts | Latency, ordering, and at-least-once delivery |
| Caching | Latency and load | Staleness 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

| Pattern | How it works | Failure mode |
| --- | --- | --- |
| Cache-aside | App reads cache, on miss loads and populates | Thundering herd on expiry |
| Read-through | Cache loads from the store itself | Same, hidden inside the cache |
| Write-through | Write to cache and store together | Slower writes, consistent reads |
| Write-behind | Write to cache, flush asynchronously | Data loss if the cache dies |
| Refresh-ahead | Refresh before expiry | Wasted work on cold keys |

```text
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

| Layer | Sees | Can do |
| --- | --- | --- |
| L4 (TCP) | Addresses and ports | Fast, protocol-agnostic, no retries or routing by path |
| L7 (HTTP) | Method, path, headers | Routing, retries, rewriting, per-route timeouts, TLS termination |

| Algorithm | Use |
| --- | --- |
| Round robin | Homogeneous backends, uniform requests |
| Least connections | Variable request duration |
| Least time / EWMA | Heterogeneous backends; best default for HTTP |
| Consistent hashing | Cache affinity; minimises reshuffling when a node leaves |
| Random two choices | Nearly 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

| Mechanism | Prevents | Watch out for |
| --- | --- | --- |
| Timeout | Unbounded waits | Must be shorter than the caller's timeout |
| Retry with jittered backoff | Transient failure | Retry storms; only retry idempotent work |
| Circuit breaker | Hammering a dead dependency | Half-open probes need to be cheap |
| Bulkhead | One slow dependency exhausting the pool | Sizing each pool |
| Rate limit | Overload and abuse | Per-tenant, not just global |
| Load shedding | Total collapse | Shed cheaply, at the edge, with 429/503 |
| Idempotency key | Duplicate side effects on retry | Storage 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

| Question | Pick |
| --- | --- |
| Relationships, transactions, ad-hoc queries | Relational, until proven otherwise |
| Known access pattern, huge volume, simple keys | Key-value or wide-column |
| Full-text search and ranking | A search engine, fed from the store of record |
| Time series with retention and rollups | A purpose-built TSDB |
| Events consumed by many independent readers | A 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

```text
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
```

| Decision | Guidance |
| --- | --- |
| Versioning | In the path (`/v1`): visible in logs, routable, obvious |
| Pagination | Cursor, not offset: stable under concurrent writes, and cheap at depth |
| Errors | A consistent envelope with a machine-readable `code` and a human `message` |
| Partial failure | Report per-item status in bulk endpoints; do not fail the batch |
| Idempotency | Accept an `Idempotency-Key` on POST and store the result |
| Concurrency | ETag plus `If-Match`, returning 409 on conflict |
| Rate limiting | `x-ratelimit-*` headers and `retry-after` on 429 |
| Long operations | Return 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.

| Concern | Answer |
| --- | --- |
| Delivery | At-least-once in practice; make consumers idempotent |
| Ordering | Only within a partition or key; design so global order is not needed |
| Poison messages | Dead-letter queue with a retry budget and an alert on depth |
| Backlog | Alert on age of the oldest message, not just depth |
| Fan-out | Topic per event type; one queue per consumer group |
| Exactly-once | Achievable only as "at-least-once plus deduplication at the sink" |

## Multi-tenancy

| Isolation | Model |
| --- | --- |
| Shared everything | Tenant column on every row and every query; cheapest, riskiest |
| Shared database, schema per tenant | Better blast radius, harder migrations |
| Database per tenant | Clean isolation, operational overhead grows linearly |
| Cluster per tenant | Regulatory 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.

```json
{
  "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](https://www.rfc-editor.org/rfc/rfc9457) (`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:

```text
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:

```json
{ "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.

```text
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.

```text
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.

```sql
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](https://www.wiki.jodisand.me/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.

| Change | Safe sequence |
| --- | --- |
| Add a column | Add it nullable or with a default; deploy code that writes it; backfill in batches; then add `NOT NULL` if needed |
| Remove a column | Deploy code that stops reading it; deploy code that stops writing it; then drop it in a later release |
| Rename a column | Add the new column; dual-write; backfill; switch reads; stop writing the old; drop it |
| Change a type | Add a new column of the new type and treat it as a rename |
| Add an index | `CREATE INDEX CONCURRENTLY` (PostgreSQL); online DDL or gh-ost / pt-online-schema-change (MySQL) |
| Add a constraint | `NOT 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](https://www.wiki.jodisand.me/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.

```text
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](https://www.wiki.jodisand.me/prometheus/), [Grafana](https://www.wiki.jodisand.me/grafana/) and [OpenTelemetry](https://www.wiki.jodisand.me/opentelemetry/).

## Troubleshooting design symptoms

| Symptom | Likely cause | Check |
| --- | --- | --- |
| p99 latency climbs while the mean is flat | Queueing near saturation, or one slow dependency in a fan-out | Utilisation and queue depth per resource; per-dependency latency percentiles |
| Latency jumps sharply at a steady traffic increase | A resource past roughly 80% utilisation | USE metrics on CPU, connection pools, disk and the database |
| Database load spikes at the same moment every few minutes | Synchronised cache expiry (thundering herd) | TTL jitter, single-flight on miss |
| Error rate rises after a dependency blip and stays high | Retry storm, or retries without jitter | Retry counts per caller; backoff and retry budget settings |
| Every instance leaves the load balancer together | Readiness check depends on a shared dependency | What the readiness endpoint calls |
| A user cannot see a change they just saved | Read served by a lagging replica | Replication lag; route read-after-write to the primary |
| One shard or partition is hot | Poor shard key, or a single large tenant | Per-key request distribution |
| Duplicate orders, emails or charges | At-least-once delivery or client retries without idempotency | Idempotency keys at the consumer or API |
| Queue depth is low but work is hours late | A stuck or poison message blocks a partition | Age of the oldest message, dead-letter queue |
| Deploy blocks every query for minutes | A migration waiting on a lock behind a long transaction | `lock_timeout` in the migration; `pg_stat_activity` for the blocker |
| Rollback fails with a schema error | Migration was not compatible with the previous code | Expand-and-contract; test the old release against the new schema in CI |
| Events missing from the queue after a crash | Publish happened outside the database transaction | Transactional outbox or CDC |
| Same event processed twice with different outcomes | Consumer is not idempotent | Processed-set keyed by event ID, or an upsert |
| Client sees 409 on every retry of a `POST` | Idempotency key still marked in flight after a crash | Expire in-flight keys; write key and side effect in one transaction |
| Tail latency doubles when one replica is slow | Fan-out without hedging or per-replica timeouts | Hedge idempotent reads at p95; per-call deadlines |
| Metrics store runs out of memory | Unbounded label values | Series count per metric; drop the label or aggregate at ingestion |
| Trace stops at the queue | Context not propagated in message headers | Inject and extract `traceparent` on publish and consume |
| Alert fires, nobody knows what to do | Cause-based alert without a runbook | Rewrite as an SLO burn alert with a runbook link |

Metrics for most of these checks come from [Prometheus](https://www.wiki.jodisand.me/prometheus/) or traces from [OpenTelemetry](https://www.wiki.jodisand.me/opentelemetry/). Test failure paths with the approaches in [Testing](https://www.wiki.jodisand.me/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?


