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.
the 2–3 hardest parts: hot keys, fan-out, consistency, failure
spreading thin over everything
8. Trade-offs & evolution
what breaks at 10x, what you would monitor, alternatives rejected
presenting one option as the only one
Availability targets and their yearly downtime budget:
SLO
Downtime / year
Downtime / 30 days
99%
3.65 days
7.2 h
99.9%
8.8 h
43 min
99.95%
4.4 h
22 min
99.99%
53 min
4.3 min
99.999%
5.3 min
26 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.
Operation
Time
Relative
L1 cache reference
~1 ns
Branch mispredict
~3 ns
L2 cache reference
~4 ns
Mutex lock/unlock
~20 ns
Main memory reference
~100 ns
100x 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
Power
Exact
Approx
Unit
2^10
1,024
thousand
KB
2^20
1,048,576
million
MB
2^30
1,073,741,824
billion
GB
2^40
trillion
TB
2^50
quadrillion
PB
2^32
4,294,967,296
4.3 billion
int32 range
2^64
1.8 × 10^19
int64 range
Handy constant
Value
Seconds per day
86,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 / average
2–10x (design for peak)
Estimation formulas
QPSavg=86400DAU×actions per user per dayQPSpeak≈3×QPSavgstorage=writes/day×bytes per write×retention days×replication factorbandwidth=QPS×bytes per responseservers=QPS per server×target utilizationQPSpeak
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
Approach
Pros
Cons
Vertical (bigger box)
no code change, strong consistency stays simple
ceiling, single point of failure, pricey at the top
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.
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)
Strategy
Key → shard
Pros
Cons
Range
key ranges (a–f, dates)
range scans
hot spots on sequential keys
Hash
hash(key) mod N
even spread
resharding moves most keys; no range scans
Consistent hashing / virtual nodes
ring position
adding a node moves ~1/N of keys
more moving parts
Directory (lookup table)
explicit map
flexible moves, tenant placement
lookup service is critical path
Geo / tenant
region or customer
data residency, isolation
uneven 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
Index
Good for
Cost
B-tree
equality, ranges, sort, prefix
default; slows writes a little per index
Hash
equality only
rarely worth it over B-tree
LSM-tree (storage engine)
write-heavy workloads
compaction, read amplification
Inverted (GIN)
full-text, JSONB, arrays
slow updates
Geo (GiST, R-tree, geohash)
nearest, within
specialized
Covering / composite
index-only scans, leftmost-prefix queries
column 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 style
Partition
Normal operation
Postgres primary + sync replica
PC
EC
Spanner, CockroachDB, etcd
PC
EC
Cassandra, DynamoDB (default)
PA
EL
Model
Guarantee
Example
Linearizable (strong)
reads see the latest write, one global order
etcd, Spanner, single leader reads
Sequential
all see the same order, not necessarily real-time
ZooKeeper writes
Causal
cause before effect for everyone; concurrent ops may differ
comment after post
Read-your-writes
you see your own writes
profile edit then reload
Monotonic reads
never go back in time
pin a user to one replica
Bounded staleness
at most t seconds or k versions behind
Cosmos DB option
Eventual
replicas converge if writes stop
DNS, Dynamo-style stores
Transaction isolation
Prevents
Allows
Read committed (Postgres default)
dirty reads
non-repeatable reads, lost updates
Repeatable read / snapshot
the above + non-repeatable reads
write skew
Serializable
all anomalies
aborts 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)
Model
append-only partitioned log
exchanges route to queues
hosted queue
Consumption
pull, consumer groups track offsets
push, per-message ack
pull, visibility timeout
Ordering
per partition
per queue (single consumer)
FIFO queues per group ID; standard: best-effort
Retention
time/size based, replayable
until acked
up to 14 days, until deleted
Throughput
very high
high
high, elastic
Fits
event streams, CDC, analytics, event sourcing
task queues, RPC, complex routing
serverless jobs, decoupling on AWS
Similar
Redpanda, Kinesis, Pulsar, NATS JetStream
ActiveMQ, NATS
Google Pub/Sub, Azure Service Bus
Delivery
Meaning
How
At most once
may lose, never duplicates
ack before processing
At least once
never lose, may duplicate
ack after processing (the default)
Effectively once
duplicates have no effect
at-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 outboxWHERE sent_at IS NULLORDER BY idLIMIT 100FOR UPDATE SKIP LOCKED;
Rate limiting
Algorithm
How
Pros
Cons
Fixed window
counter per key per minute
trivial (INCR + EXPIRE)
2x burst at window edges
Sliding window log
timestamps of each request
exact
memory per request
Sliding window counter
weighted current + previous window
smooth, cheap
approximate
Token bucket
tokens refill at rate r, burst up to b
allows bursts, common in APIs
two parameters to tune
Leaky bucket
queue drains at fixed rate
smooth output
adds latency, drops on full
GCRA
token bucket as one timestamp
one value per key, atomic
less 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.
on every network call; below the caller's timeout; budget the whole request (deadline propagation)
Retries
only idempotent operations or with idempotency keys; exponential backoff + jitter; cap attempts; retry at one layer only
Retry budget
allow retries up to ~10% of requests, stop amplifying outages
Circuit breaker
closed → open after N failures (fail fast) → half-open trial → closed
Bulkhead
separate pools/queues per dependency so one slow dependency cannot exhaust all
Load shedding
reject early (503/429) when queues or latency pass a threshold; prioritize critical traffic
Backpressure
bounded queues; slow producers instead of buffering forever
Graceful degradation
serve cached, partial or default content when a dependency is down
Idempotency keys
client sends Idempotency-Key; server stores result per key for replays
Health checks + auto-replace
kill and replace, don't nurse
Redundancy
N+1 (or N+2) instances across zones; test failover
Chaos / game days
inject 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
Signal
Answers
Tools
Metrics
is it healthy, how much, how fast
Prometheus, Mimir, OpenTelemetry metrics
Logs
what happened, with detail
Loki, Elasticsearch; structured JSON with trace IDs
Traces
where the time went across services
Tempo, Jaeger; OpenTelemetry SDKs
Profiles
which code burns CPU/memory
Pyroscope
Method
Measure
For
RED
Rate, Errors, Duration
request-driven services
USE
Utilization, Saturation, Errors
resources (CPU, disk, pools)
Four golden signals
latency, traffic, errors, saturation
SRE dashboards
SLI / SLO / error budget
measured indicator, target, allowed failure
alerting and release pace
Alert on symptoms (SLO burn rate), not causes; track p50/p95/p99 not averages. Setup
details in Grafana.
IDs
Scheme
Size
Sortable
Notes
DB auto-increment
64-bit
yes
simple; leaks counts, needs a single writer or ranges per shard
UUIDv4
128-bit
no
random; poor B-tree locality at scale
UUIDv7 (RFC 9562)
128-bit
yes (ms)
48-bit Unix ms timestamp + random; Postgres 18 uuidv7()
ULID
128-bit
yes (ms)
Crockford base32 text form; same idea as UUIDv7
Snowflake
64-bit
yes (ms)
41-bit ms timestamp, 10-bit worker ID, 12-bit sequence: 4,096 IDs/ms/worker
KSUID
160-bit
yes (s)
32-bit timestamp + 128-bit random
Short codes (base62)
6–8 chars
no
62^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 option
Fits
Postgres full-text (tsvector + GIN), pg_trgm
modest scale, one less system
Meilisearch, Typesense
typo-tolerant product/site search, simple ops
Elasticsearch / OpenSearch
large-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 practice
Why
Object store (S3, R2, GCS) for files, DB for metadata
cheap, durable (11 nines), unlimited
Presigned PUT/GET URLs
clients upload and download directly, not through your servers
Multipart uploads
large files, resumable, parallel
CDN in front
cached delivery, lower egress
Content-addressed keys (sha256/...)
dedupe, immutable caching
Lifecycle rules
move to cold tiers, expire temp files
Event on upload
queue thumbnailing, scanning, transcoding
Trade-off cheat table
Choice
Pick A when
Pick B when
SQL vs NoSQL
relations, transactions, unknown queries
known access patterns at huge scale
Strong vs eventual consistency
money, inventory, uniqueness
feeds, counters, analytics
Sync (HTTP/gRPC) vs async (queue)
caller needs the answer now
work can happen later, spikes need smoothing
Push vs pull (fan-out)
few followers: write to each feed
celebrities: merge at read time
Cache vs no cache
read-heavy, tolerates staleness
write-heavy, must be fresh
Monolith vs microservices
small team, evolving domain
many teams, independent deploy and scale
Polling vs WebSockets/SSE
infrequent updates, simple
real-time, many updates
Normalize vs denormalise
write integrity
read speed
Batch vs stream processing
hourly/daily results fine
seconds-level freshness
Vertical vs horizontal scale
early, simple
past one machine, need redundancy
Build vs buy (managed)
core differentiator
commodity (auth, queues, email)
Leader-based vs leaderless
ordering, simple reasoning
write availability across failures
301 vs 302 redirect
cacheable forever, less load
need 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
Quantity
Value
New links
100M/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)
Storage
100M × 12 × 5 years × 500 B ≈ 3 TB (before replicas)
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 storeclient ─► L7 LB ─► link API ─► ID generator ─► DB primary
6. Deep dives
Problem
Decision
Code generation
counter 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
Collisions
primary key on code; retry on conflict for random codes
Read latency
cache-aside in Redis (hot 20% of links serve ~80% of traffic), plus edge caching of redirects with short TTL
Redirect status
302 so every click reaches us for analytics; 301 if analytics don't matter
Analytics
fire-and-forget event to a queue; aggregate in a column store; never block the redirect
Abuse
rate limit creates per user/IP, scan URLs against blocklists asynchronously
Expiry
check expires_at on read; nightly job deletes expired rows and cache keys
Availability
stateless 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.