../

System design

Building blocks, numbers and trade-offs for designing distributed systems, plus a step-by-step framework and a worked design. Protocol details live in Communication networks; code-level structure in Software architecture.

Approach framework

StepOutputWatch for
1. Functional requirements3–5 core use cases, what is out of scopebuilding features nobody asked for
2. Non-functional requirementsscale, latency (p99), availability, consistency, durability, cost"highly available" without a number
3. EstimatesQPS (avg and peak), storage/year, bandwidth, read:write ratioprecision: orders of magnitude are enough
4. APIendpoints or RPCs with request/response shapes, pagination, idempotencychatty or ambiguous contracts
5. Data modelentities, access patterns, keys, indexes, store choicechoosing a DB before knowing queries
6. High-level designboxes: clients, LB, services, caches, DBs, queues, CDNskipping the data flow for reads vs writes
7. Deep divesthe 2–3 hardest parts: hot keys, fan-out, consistency, failurespreading thin over everything
8. Trade-offs & evolutionwhat breaks at 10x, what you would monitor, alternatives rejectedpresenting one option as the only one

Availability targets and their yearly downtime budget:

SLODowntime / yearDowntime / 30 days
99%3.65 days7.2 h
99.9%8.8 h43 min
99.95%4.4 h22 min
99.99%53 min4.3 min
99.999%5.3 min26 s

Serial dependencies multiply: two 99.9% services in a chain give about 99.8%.

Back-of-envelope numbers

Order-of-magnitude latencies on current hardware (Jeff Dean's list, updated). See Latency and bandwidth for network detail.

OperationTimeRelative
L1 cache reference~1 ns
Branch mispredict~3 ns
L2 cache reference~4 ns
Mutex lock/unlock~20 ns
Main memory reference~100 ns100x L1
Compress 1 KB (Snappy/LZ4)~2 µs
Send 1 KB over 10 Gbps~1 µs
Read 1 MB sequentially from RAM~3–10 µs
Random 4 KB read, NVMe SSD~20–100 µs
Read 1 MB sequentially from NVMe~50–200 µs
Round trip within a datacenter~0.5 ms
Redis GET over the LAN~0.2–1 ms
Indexed Postgres query, warm~1 ms
HDD seek~5–10 ms
Round trip same region (cross-AZ)~1–2 ms
Round trip US east to west~60–70 ms
Round trip Europe to US west~140–150 ms
PowerExactApproxUnit
2^101,024thousandKB
2^201,048,576millionMB
2^301,073,741,824billionGB
2^40trillionTB
2^50quadrillionPB
2^324,294,967,2964.3 billionint32 range
2^641.8 × 10^19int64 range
Handy constantValue
Seconds per day86,400 ≈ 10^5
Seconds per month≈ 2.6 million
Seconds per year≈ 3.15 × 10^7 (π × 10^7)
1 million requests/day≈ 12 req/s
1 billion requests/day≈ 12,000 req/s
Peak / average2–10x (design for peak)

Estimation formulas

QPSavg=DAU×actions per user per day86 400QPSpeak≈3×QPSavg\text{QPS}_{\text{avg}} = \frac{\text{DAU} \times \text{actions per user per day}}{86\,400} \qquad \text{QPS}_{\text{peak}} \approx 3 \times \text{QPS}_{\text{avg}} storage=writes/day×bytes per write×retention days×replication factor\text{storage} = \text{writes/day} \times \text{bytes per write} \times \text{retention days} \times \text{replication factor} bandwidth=QPS×bytes per responseservers=QPSpeakQPS per server×target utilization\text{bandwidth} = \text{QPS} \times \text{bytes per response} \qquad \text{servers} = \frac{\text{QPS}_{\text{peak}}}{\text{QPS per server} \times \text{target utilization}}

Example: 10M DAU × 20 reads = 200M reads/day ≈ 2,300 QPS avg, ~7,000 peak. At 2 KB per response that is ~14 MB/s egress. With 5M writes/day × 1 KB × 365 × 3 replicas ≈ 5.5 TB/year.

Little's law for concurrency: in-flight requests = arrival rate × latency. 1,000 req/s at 200 ms means 200 concurrent requests (and connections, threads or DB sessions).

Scaling & load balancing

ApproachProsCons
Vertical (bigger box)no code change, strong consistency stays simpleceiling, single point of failure, pricey at the top
Horizontal (more boxes)near-linear scale, redundancyneeds stateless services, partitioning, coordination

Stateless services keep session, uploads and caches out of process memory (Redis, DB, object store, signed tokens) so any instance can serve any request, instances can be killed at will, and autoscaling works.

Load balancerWorks onCan doExamples
L4 (transport)IP + port, TCP/UDPvery fast, TLS passthrough, any protocolAWS NLB, IPVS, HAProxy TCP mode
L7 (application)HTTP requestspath/host routing, TLS termination, retries, header auth, gRPCnginx, Envoy, HAProxy, AWS ALB, Caddy
DNS / GSLBname resolutiongeo routing, failover between regionsRoute 53, Cloudflare
AnycastBGPone IP served from many sitesCDNs, Cloudflare, Google front ends
AlgorithmPicksGood for
Round robin / weightednext in turnuniform, short requests
Least connectionsfewest open connectionslong-lived or uneven requests
Least response time / EWMAfastest recent backendheterogeneous backends
Power of two choicesbest of 2 random backendslarge fleets, avoids herding
IP / header hashsame client, same backendsticky sessions (prefer stateless)
Consistent hashingkey's position on a ringcaches, shards: few keys move on resize

Health checks: active (LB probes /healthz) plus passive (eject on errors). Separate liveness (restart me) from readiness (send me traffic).

Caching & CDNs

LayerHoldsTTL typical
Browser (Cache-Control)static assets, API responsesimmutable assets: 1 year
CDN / edgeassets, images, cacheable HTML/APIseconds to days
Reverse proxyrendered pages, API responsesseconds
Application (in-process LRU)hot config, small lookupsseconds
Distributed (Redis, Valkey, Memcached)sessions, computed views, objectsminutes to hours
Database (buffer pool, materialized views)pages, precomputed queriesmanaged by DB
PatternRead pathWrite pathTrade-off
Cache-aside (lazy)app reads cache, on miss reads DB and fillsapp writes DB, deletes keysimple, default choice; first read is slow
Read-throughcache library loads from DB on missas cache-asidecleaner app code
Write-throughread cachewrite cache and DB synchronouslyalways warm; slower writes
Write-behind (write-back)read cachewrite cache, flush to DB asyncfast writes; data loss risk
Refresh-aheadserve cached, refresh before expiryn/alow latency for hot keys
ProblemFix
Stale data after writedelete (not update) the key after the DB commit; short TTL as a backstop
Race: old value re-cachedversioned keys, or delay-and-delete again, or CDC-driven invalidation
Stampede on expiry (dogpile)single-flight lock per key, stale-while-revalidate, probabilistic early refresh
Synchronized expiryadd jitter to TTLs (ttl × (0.9 + 0.2 × random))
Hot keylocal in-process cache in front, replicate key across shards
Cache penetration (misses for absent keys)cache negative results briefly, Bloom filter
Cold startpre-warm from logs or a snapshot
cache-aside with single-flight
interface Cache {
  get(key: string): Promise<string | null>;
  set(key: string, v: string, ttlS: number): Promise<void>;
}
 
const inFlight = new Map<string, Promise<string>>();
 
export async function cached(
  cache: Cache,
  key: string,
  ttlS: number,
  load: () => Promise<string>,
): Promise<string> {
  const hit = await cache.get(key);
  if (hit !== null) return hit;
  const pending = inFlight.get(key);
  if (pending) return pending;            // coalesce
  const p = load()
    .then(async (v) => {
      const jitter = 0.9 + Math.random() * 0.2;
      await cache.set(key, v, Math.round(ttlS * jitter));
      return v;
    })
    .finally(() => inFlight.delete(key));
  inFlight.set(key, p);
  return p;
}

CDNs

FeatureUse
Pull (origin fetch on miss)default; origin shield reduces origin hits
Push (upload to CDN)large, rarely changing files
Cache-Control: public, max-age, s-maxagebrowser vs shared-cache lifetimes
stale-while-revalidate, stale-if-errorserve stale during refresh or origin outage
Content-hashed filenamescache forever, deploy by changing names
Purge / surrogate keys (tags)invalidate groups of objects on change
Signed URLs / cookiesprivate content at the edge
Edge computeauth, redirects, A/B, personalization near users

Databases

Relational (SQL)DocumentWide-columnKey-valueGraph
ExamplesPostgres, MySQLMongoDB, FirestoreCassandra, ScyllaDB, BigtableRedis, DynamoDBNeo4j
Modeltables, joins, constraintsJSON documentspartition key + sorted columnsopaque values by keynodes + edges
Transactionsfull ACIDper document, multi-doc limitedper partition, lightweightper keyACID (varies)
Scales writes byvertical, read replicas, sharding (Citus, Vitess)built-in shardingbuilt-in, linearbuilt-in partitioningmostly vertical
Pick whendefault; relations, integrity, ad-hoc queriesself-contained aggregates, flexible schemahuge write volume, time-series, known queriessessions, caches, counters, simple lookupsdeep relationship traversal

Default to Postgres until a measured need says otherwise. See Postgres and Database design.

Replication

TopologyWritesConsistencyFailure mode
Single leader, async followersleader onlyfollowers lag (stale reads)failover can lose recent writes
Single leader, sync followerleader, ack after 1 replicano loss on failoverwrite latency, stalls if replica down
Multi-leaderany leader (per region)conflicts need resolution (LWW, CRDTs)divergent data
Leaderless (Dynamo-style)any N replicasquorum: W + R > Nsloppy quorums, read repair
Consensus (Raft/Paxos)leader via majoritylinearizableneeds majority alive (2f+1 nodes for f failures)

Read-your-writes with async replicas: route a user's reads to the leader for a few seconds after they write, or pass the commit LSN and wait for the replica to catch up.

Sharding (partitioning)

StrategyKey → shardProsCons
Rangekey ranges (a–f, dates)range scanshot spots on sequential keys
Hashhash(key) mod Neven spreadresharding moves most keys; no range scans
Consistent hashing / virtual nodesring positionadding a node moves ~1/N of keysmore moving parts
Directory (lookup table)explicit mapflexible moves, tenant placementlookup service is critical path
Geo / tenantregion or customerdata residency, isolationuneven tenant sizes

Pick a shard key that is in almost every query, has high cardinality, and spreads writes (tenant_id, user_id). Cross-shard joins and transactions become application work.

Indexes

IndexGood forCost
B-treeequality, ranges, sort, prefixdefault; slows writes a little per index
Hashequality onlyrarely worth it over B-tree
LSM-tree (storage engine)write-heavy workloadscompaction, read amplification
Inverted (GIN)full-text, JSONB, arraysslow updates
Geo (GiST, R-tree, geohash)nearest, withinspecialized
Covering / compositeindex-only scans, leftmost-prefix queriescolumn order matters

Consistency

CAP: during a network partition a replicated store must choose consistency (refuse some requests) or availability (serve possibly stale data). PACELC adds: else (no partition), choose latency or consistency.

System stylePartitionNormal operation
Postgres primary + sync replicaPCEC
Spanner, CockroachDB, etcdPCEC
Cassandra, DynamoDB (default)PAEL
ModelGuaranteeExample
Linearizable (strong)reads see the latest write, one global orderetcd, Spanner, single leader reads
Sequentialall see the same order, not necessarily real-timeZooKeeper writes
Causalcause before effect for everyone; concurrent ops may differcomment after post
Read-your-writesyou see your own writesprofile edit then reload
Monotonic readsnever go back in timepin a user to one replica
Bounded stalenessat most t seconds or k versions behindCosmos DB option
Eventualreplicas converge if writes stopDNS, Dynamo-style stores
Transaction isolationPreventsAllows
Read committed (Postgres default)dirty readsnon-repeatable reads, lost updates
Repeatable read / snapshotthe above + non-repeatable readswrite skew
Serializableall anomaliesaborts under contention (retry)

Distributed transactions: avoid 2PC across services; use a saga (local transactions plus compensating actions) orchestrated or choreographed by events.

Queues & streams

Kafka (log)RabbitMQ (broker)SQS (managed queue)
Modelappend-only partitioned logexchanges route to queueshosted queue
Consumptionpull, consumer groups track offsetspush, per-message ackpull, visibility timeout
Orderingper partitionper queue (single consumer)FIFO queues per group ID; standard: best-effort
Retentiontime/size based, replayableuntil ackedup to 14 days, until deleted
Throughputvery highhighhigh, elastic
Fitsevent streams, CDC, analytics, event sourcingtask queues, RPC, complex routingserverless jobs, decoupling on AWS
SimilarRedpanda, Kinesis, Pulsar, NATS JetStreamActiveMQ, NATSGoogle Pub/Sub, Azure Service Bus
DeliveryMeaningHow
At most oncemay lose, never duplicatesack before processing
At least oncenever lose, may duplicateack after processing (the default)
Effectively onceduplicates have no effectat-least-once + idempotent consumer

Idempotent consumers: dedupe by message ID in a processed-messages table inside the same transaction as the side effect, or make the write naturally idempotent (upsert, set not increment). Failed messages go to a dead-letter queue after N attempts.

Transactional outbox

Write the business row and an outbox row in one DB transaction; a relay (poller or CDC like Debezium) publishes outbox rows and marks them sent. No dual-write gap.

BEGIN;
INSERT INTO orders (id, user_id, total)
  VALUES ($1, $2, $3);
INSERT INTO outbox (id, topic, payload)
  VALUES (gen_random_uuid(), 'order.created',
          jsonb_build_object('orderId', $1));
COMMIT;
 
-- relay, run in a loop (multiple relays safe)
SELECT id, topic, payload FROM outbox
WHERE sent_at IS NULL
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED;

Rate limiting

AlgorithmHowProsCons
Fixed windowcounter per key per minutetrivial (INCR + EXPIRE)2x burst at window edges
Sliding window logtimestamps of each requestexactmemory per request
Sliding window counterweighted current + previous windowsmooth, cheapapproximate
Token buckettokens refill at rate r, burst up to ballows bursts, common in APIstwo parameters to tune
Leaky bucketqueue drains at fixed ratesmooth outputadds latency, drops on full
GCRAtoken bucket as one timestampone value per key, atomicless intuitive

Where: at the edge/gateway per IP and per API key, again per user in the service. Return 429 Too Many Requests with Retry-After and RateLimit headers. Distributed limits need a shared store (Redis with a Lua script) or accept per-node approximations.

token-bucket.ts
export class TokenBucket {
  private tokens: number;
  private last = performance.now();
 
  constructor(
    private readonly capacity: number,
    private readonly perSecond: number,
  ) {
    this.tokens = capacity;
  }
 
  take(cost = 1): boolean {
    const now = performance.now();
    const elapsed = (now - this.last) / 1000;
    this.tokens = Math.min(
      this.capacity,
      this.tokens + elapsed * this.perSecond,
    );
    this.last = now;
    if (this.tokens < cost) return false;
    this.tokens -= cost;
    return true;
  }
}

Reliability

TechniqueRule of thumb
Timeoutson every network call; below the caller's timeout; budget the whole request (deadline propagation)
Retriesonly idempotent operations or with idempotency keys; exponential backoff + jitter; cap attempts; retry at one layer only
Retry budgetallow retries up to ~10% of requests, stop amplifying outages
Circuit breakerclosed → open after N failures (fail fast) → half-open trial → closed
Bulkheadseparate pools/queues per dependency so one slow dependency cannot exhaust all
Load sheddingreject early (503/429) when queues or latency pass a threshold; prioritize critical traffic
Backpressurebounded queues; slow producers instead of buffering forever
Graceful degradationserve cached, partial or default content when a dependency is down
Idempotency keysclient sends Idempotency-Key; server stores result per key for replays
Health checks + auto-replacekill and replace, don't nurse
RedundancyN+1 (or N+2) instances across zones; test failover
Chaos / game daysinject failure in controlled drills
retry with full jitter
const sleep = (ms: number) =>
  new Promise<void>((r) => setTimeout(r, ms));
 
export async function retry<T>(
  fn: (signal: AbortSignal) => Promise<T>,
  { attempts = 5, baseMs = 100, capMs = 5_000,
    timeoutMs = 2_000 } = {},
): Promise<T> {
  for (let i = 0; ; i++) {
    try {
      return await fn(AbortSignal.timeout(timeoutMs));
    } catch (err) {
      if (i + 1 >= attempts) throw err;
      const max = Math.min(capMs, baseMs * 2 ** i);
      await sleep(Math.random() * max);   // full jitter
    }
  }
}

Observability

SignalAnswersTools
Metricsis it healthy, how much, how fastPrometheus, Mimir, OpenTelemetry metrics
Logswhat happened, with detailLoki, Elasticsearch; structured JSON with trace IDs
Traceswhere the time went across servicesTempo, Jaeger; OpenTelemetry SDKs
Profileswhich code burns CPU/memoryPyroscope
MethodMeasureFor
REDRate, Errors, Durationrequest-driven services
USEUtilization, Saturation, Errorsresources (CPU, disk, pools)
Four golden signalslatency, traffic, errors, saturationSRE dashboards
SLI / SLO / error budgetmeasured indicator, target, allowed failurealerting and release pace

Alert on symptoms (SLO burn rate), not causes; track p50/p95/p99 not averages. Setup details in Grafana.

IDs

SchemeSizeSortableNotes
DB auto-increment64-bityessimple; leaks counts, needs a single writer or ranges per shard
UUIDv4128-bitnorandom; poor B-tree locality at scale
UUIDv7 (RFC 9562)128-bityes (ms)48-bit Unix ms timestamp + random; Postgres 18 uuidv7()
ULID128-bityes (ms)Crockford base32 text form; same idea as UUIDv7
Snowflake64-bityes (ms)41-bit ms timestamp, 10-bit worker ID, 12-bit sequence: 4,096 IDs/ms/worker
KSUID160-bityes (s)32-bit timestamp + 128-bit random
Short codes (base62)6–8 charsno62^7 ≈ 3.5 × 10^12; from a counter or random with collision check

Prefer time-ordered IDs for primary keys (index locality) and don't expose sequential IDs where enumeration is a risk.

Search & blob storage

Search optionFits
Postgres full-text (tsvector + GIN), pg_trgmmodest scale, one less system
Meilisearch, Typesensetypo-tolerant product/site search, simple ops
Elasticsearch / OpenSearchlarge-scale text, logs, aggregations
Vector search (pgvector, dedicated vector DBs)semantic similarity, RAG

Keep the database as the source of truth and feed the index via CDC or the outbox; the index is a rebuildable projection with eventual consistency.

Blob storage practiceWhy
Object store (S3, R2, GCS) for files, DB for metadatacheap, durable (11 nines), unlimited
Presigned PUT/GET URLsclients upload and download directly, not through your servers
Multipart uploadslarge files, resumable, parallel
CDN in frontcached delivery, lower egress
Content-addressed keys (sha256/...)dedupe, immutable caching
Lifecycle rulesmove to cold tiers, expire temp files
Event on uploadqueue thumbnailing, scanning, transcoding

Trade-off cheat table

ChoicePick A whenPick B when
SQL vs NoSQLrelations, transactions, unknown queriesknown access patterns at huge scale
Strong vs eventual consistencymoney, inventory, uniquenessfeeds, counters, analytics
Sync (HTTP/gRPC) vs async (queue)caller needs the answer nowwork can happen later, spikes need smoothing
Push vs pull (fan-out)few followers: write to each feedcelebrities: merge at read time
Cache vs no cacheread-heavy, tolerates stalenesswrite-heavy, must be fresh
Monolith vs microservicessmall team, evolving domainmany teams, independent deploy and scale
Polling vs WebSockets/SSEinfrequent updates, simplereal-time, many updates
Normalize vs denormalisewrite integrityread speed
Batch vs stream processinghourly/daily results fineseconds-level freshness
Vertical vs horizontal scaleearly, simplepast one machine, need redundancy
Build vs buy (managed)core differentiatorcommodity (auth, queues, email)
Leader-based vs leaderlessordering, simple reasoningwrite availability across failures
301 vs 302 redirectcacheable forever, less loadneed per-click analytics

Worked design: URL shortener

1. Requirements

Functional: create a short link for a long URL (optional custom alias, expiry), redirect, basic click counts. Non-functional: redirects p99 under 50 ms, 99.99% available for reads, links never collide, created links durable.

2. Estimates

QuantityValue
New links100M/month ≈ 40 writes/s avg, ~200 peak
Redirects (100:1)10B/month ≈ 4,000 reads/s avg, ~20,000 peak
Row size~500 bytes (URL, code, owner, timestamps)
Storage100M × 12 × 5 years × 500 B ≈ 3 TB (before replicas)
Codes needed6B in 5 years: 7 base62 chars (62^7 ≈ 3.5 × 10^12)

3. API

POST /api/links
Idempotency-Key: 5f3c...
Content-Type: application/json
 
{"url": "https://example.com/very/long", "alias": null}
 
201 Created
{"code": "aZ3kP9q", "shortUrl": "https://sho.rt/aZ3kP9q"}
 
GET /aZ3kP9q
302 Found
Location: https://example.com/very/long

4. Data model

CREATE TABLE links (
  code        text PRIMARY KEY,       -- base62, 7 chars
  url         text NOT NULL,
  owner_id    uuid,
  created_at  timestamptz NOT NULL DEFAULT now(),
  expires_at  timestamptz
);
CREATE INDEX ON links (owner_id, created_at DESC);

Access pattern is a point lookup by code: any store works; Postgres (with read replicas) or a key-value store like DynamoDB both fit. Click events go elsewhere.

5. High-level design

client ─► CDN/edge ─► L7 LB ─► redirect service ─► Redis ─► DB replicas
                                   │ (miss)
                                   └─► click event ─► queue ─► analytics store
client ─► L7 LB ─► link API ─► ID generator ─► DB primary

6. Deep dives

ProblemDecision
Code generationcounter blocks per API instance (claim 10,000 IDs at a time from the DB), encode base62, shuffle bits so codes aren't guessable; custom aliases via unique insert
Collisionsprimary key on code; retry on conflict for random codes
Read latencycache-aside in Redis (hot 20% of links serve ~80% of traffic), plus edge caching of redirects with short TTL
Redirect status302 so every click reaches us for analytics; 301 if analytics don't matter
Analyticsfire-and-forget event to a queue; aggregate in a column store; never block the redirect
Abuserate limit creates per user/IP, scan URLs against blocklists asynchronously
Expirycheck expires_at on read; nightly job deletes expired rows and cache keys
Availabilitystateless services in 3 zones, DB with replicas and automated failover, Redis replicated

7. Trade-offs

Pre-allocated counter blocks leave gaps and need coordination on startup, but avoid a hot central counter; random codes avoid coordination but cost a uniqueness check. Eventual consistency between primary and replicas means a just-created link may 404 for a moment on a replica: read from the primary on cache miss for links younger than a few seconds.

References