Scaling
Increase throughput by adding resources to one node (vertical) or adding identical nodes behind a load balancer (horizontal).
Vertical scaling adds CPU, memory, or IOPS to a single node. It needs no application changes and keeps consistency trivial because there is exactly one copy of state. The limits are a hard ceiling set by the largest available instance, cost that grows faster than capacity near the top, a restart to resize, and a node that remains a single point of failure.
Horizontal scaling runs N identical replicas behind a load balancer. Throughput grows roughly linearly as long as replicas are stateless and the shared dependencies (database, cache, third-party APIs) can absorb the extra connections. In production the replica count is driven by autoscaling and sized against p99 latency targets, not averages.
Production checklist
- Externalize state: sessions in Redis or signed tokens, files in object storage, nothing on local disk that must survive a restart.
- Budget connections. 50 pods with a pool of 20 each is 1,000 Postgres connections; put a pooler such as PgBouncer in front of the database.
- Start fast and shut down gracefully so scale-in and deploys don't drop in-flight requests (see connection draining).
- Run N+1 across availability zones: losing one zone must leave enough capacity to serve peak traffic.
- Load-test to find the first bottleneck before production does. It is usually the database, not the app tier.
You get
- Vertical: no code changes, single-node transactions, simplest operations
- Horizontal: near-linear throughput, zone-level fault tolerance, rolling deploys
You pay
- Vertical: hard ceiling, single point of failure, downtime to resize
- Horizontal: stateless design, a load-balancing tier, connection fan-out to shared stores, distributed debugging
Rule of thumbScale vertically while it's cheap, but ship stateless services from day one so horizontal scaling is an autoscaler setting, not a migration.
Load balancing
Distributes traffic across a pool of backends and removes unhealthy ones from rotation.
A load balancer accepts client connections and forwards each request or connection to a backend chosen by an algorithm. Its most important job is failure isolation. Active health checks probe every backend on an interval (for example every 5 s, unhealthy after 3 failures, healthy again after 2 successes), and passive checks, called outlier detection in Envoy, eject backends that start returning 5xx responses. Try it: take App B down in the live diagram on the home page.
Layer 4 vs layer 7
L4 balancers (AWS NLB, Linux IPVS, Google Maglev) forward TCP or UDP flows by their 5-tuple without parsing the payload. They handle millions of connections with microsecond overhead and can preserve the client IP. L7 balancers (AWS ALB, Envoy, NGINX, HAProxy) parse HTTP, which enables path and header routing, TLS termination, per-request retries and timeouts, and HTTP/2 or gRPC multiplexing.
One common production trap: an L4 balancer spreads connections, not requests. gRPC and HTTP/2 clients hold one long-lived connection, so every request lands on the same backend. Use an L7 balancer or client-side balancing for those protocols.
| Algorithm | How it picks | Use when |
|---|---|---|
| Round robin | Next backend in order | Homogeneous backends, uniform request cost |
| Weighted round robin | Proportional to configured weight | Mixed instance sizes, canary releases (e.g. 5% weight) |
| Least request (P2C) | Samples two random backends, picks the one with fewer active requests | Variable request cost; Envoy's default least-request mode |
| Ring hash / Maglev | Hash of a key (header, cookie, IP) selects a backend | Session affinity or per-key local caches; see virtual nodes |
Production checklist
- Separate readiness from liveness. Readiness fails when the instance can't serve; liveness fails only when the process is wedged and must be restarted.
- Enable connection draining (typically 30 to 60 s) so deploys don't reset in-flight requests.
- Set a timeout at every hop, each shorter than its caller's, and cap retries with a retry budget so a slow backend doesn't trigger a retry storm.
- Make the balancer itself redundant: managed balancers span zones; self-hosted pairs use VRRP (keepalived) or anycast.
You get
- Automatic failover and ejection of bad backends
- Zero-downtime deploys via draining
- One stable endpoint, central TLS and routing policy
You pay
- An extra network hop (sub-millisecond within a zone)
- Another tier that must be highly available
- A misconfigured health check can eject every backend at once
Rule of thumbThe health check is the product. Make readiness reflect whether an instance can actually serve, and never let one shared dependency fail every instance's liveness at the same time.
Caching
Serve repeated reads from a faster tier to cut latency and offload the source of truth.
An in-process memory lookup costs about 100 ns, a Redis round trip in the same zone roughly 0.2 to 0.5 ms, and an indexed Postgres query typically 1 to 10 ms. A cache with a 95% hit ratio reduces read load on the database by 20×. Caching happens at every layer: the browser (Cache-Control), the CDN, the application (Redis, Memcached), and the database's own buffer pool.
Write strategies
- Cache-aside (lazy loading): read the cache; on a miss, read the database and
SETthe value with a TTL. On writes, update the database, thenDELthe key. Deleting beats updating because two racing writers can't leave an older value behind. - Write-through: write the cache and the database in the same code path. Reads are fresh, writes are slower, and you may cache data nobody reads.
- Write-back (write-behind): acknowledge after writing the cache and flush to the database asynchronously. Lowest write latency, but data is lost if the cache fails before flushing.
Failure modes
- Cache stampede: a hot key expires and thousands of concurrent misses hit the database at once. Mitigate with request coalescing, TTL jitter, or early probabilistic refresh.
- Stale reads: bounded by the TTL, so a missed invalidation heals itself when the entry expires.
- Eviction: set
maxmemoryand an eviction policy such as LRU or LFU. Without a policy, Redis rejects writes once memory is full. - Cold start: an empty cache sends 100% of reads to the database. Warm critical keys before cutover and size the database to survive it.
You get
- An order of magnitude lower read latency
- Large reduction in database read load and cost
You pay
- Reads can be stale up to the TTL
- Another stateful system to operate, monitor, and fail over
- Invalidation logic spread across write paths
Rule of thumbCache data that is read far more often than it changes, always set a TTL with jitter, and size the database to survive a cold cache.
Replication & sharding
Replication copies the same data to multiple nodes; sharding partitions different data across nodes.
In leader-follower replication the leader records every change in its write-ahead log and streams it to followers. Asynchronous replication gives the lowest write latency, but a failover can lose writes that hadn't shipped yet. Synchronous or semi-synchronous replication waits for at least one follower to acknowledge, trading write latency for durability. Followers serve reads, so read throughput scales with the number of replicas, subject to replication lag. To guarantee read-your-writes, route a user's reads to the leader for a short window after they write.
Leaderless replication (Dynamo-style stores such as Cassandra) writes to N replicas and uses a quorum: with W + R > N, every read overlaps at least one replica that holds the latest write.
Sharding partitions the keyspace once write throughput or dataset size exceeds what one leader can handle. Range partitioning keeps ordered scans efficient but concentrates sequential keys, such as timestamps, on one shard. Hash partitioning spreads load evenly but gives up range queries. The shard key should appear in nearly every query so each request touches a single shard; everything else becomes a scatter-gather across all of them.
Production checklist
- Automate failover with fencing to prevent split-brain, using Patroni, Orchestrator, or a managed service such as RDS Multi-AZ or Aurora.
- Alert on replication lag in both seconds and bytes, and stop routing reads to a follower that exceeds your staleness budget.
- Pre-split into many logical shards (for example 1,024) mapped onto fewer physical nodes, so rebalancing moves whole shards instead of individual rows.
- Plan for hot keys: one very active tenant or celebrity account can saturate a single shard regardless of how many you have.
You get
- Replication: read scaling, high availability, backups and analytics without touching the leader
- Sharding: write throughput and storage beyond a single machine
You pay
- Replication: lag, stale reads, and failover edge cases
- Sharding: no cheap cross-shard joins or transactions, complex resharding, hot partitions
Rule of thumbReplicate for availability from day one. Shard only when vertical scaling and read replicas can no longer absorb the write load, and choose the shard key from your query patterns rather than your schema.
CAP & consistency
During a network partition, a replicated system must either reject requests or serve data that may be stale.
CAP is often summarized as "pick two of three," which is misleading. Network partitions are not optional in real deployments, so the trade-off only bites while one is happening: a node that can't reach a quorum must either refuse the operation (CP) or answer with data that may be stale or conflicting (AP). When the network is healthy, a well-designed system provides both consistency and availability.
PACELC: the everyday trade-off
PACELC extends CAP: if Partitioned, choose Availability or Consistency; Else, choose Latency or Consistency. The second half matters more often. Linearizable writes across regions require a cross-region round trip, commonly 50 to 200 ms, on every write.
Consistency models, strongest to weakest
- Linearizable: each operation appears to take effect atomically at one point between its start and end. Required for locks, leader election, and uniqueness guarantees.
- Causal: operations that are causally related are seen in the same order by every node; concurrent ones may differ.
- Session guarantees: read-your-writes and monotonic reads. Often sufficient for user-facing apps.
- Eventual: replicas converge once writes stop; conflicts are resolved by last-write-wins, version vectors, or CRDTs.
| System | Default | Tunable |
|---|---|---|
| PostgreSQL (single leader) | Linearizable on the primary; async replicas serve possibly stale reads | synchronous_commit, synchronous_standby_names |
| MongoDB | Reads from primary; write concern majority by default since 5.0 | readConcern, readPreference |
| DynamoDB | Eventually consistent reads | ConsistentRead=true per request (same region) |
| Cassandra | Consistency level chosen per query | ONE, LOCAL_QUORUM, QUORUM, ALL |
| etcd | Linearizable reads and writes via Raft (CP) | Serializable reads trade freshness for latency |
Rule of thumbChoose consistency per operation, not per system. Balances, inventory reservations, and idempotency keys need linearizable writes; feeds, counters, and analytics tolerate eventual consistency.
Consistent hashing
Maps keys to nodes so that adding or removing a node relocates only about 1/N of the keys.
With modulo placement, hash(key) % N, changing N from 4 to 5 remaps about 80% of all keys. For a cache cluster, that's a near-total miss storm the moment you scale. Consistent hashing places nodes and keys on the same hash space, a ring of 232 or 264 positions. A key belongs to the first node clockwise from its position, so adding a node only takes over the arc between it and its predecessor.
With a single position per node, arc sizes vary widely and one node can own several times its fair share. Production implementations give each node 100 to 256 virtual nodes, which evens out load and lets larger machines take proportionally more by weight. Replication fits naturally: store each key on the next R distinct physical nodes clockwise.
Squares are servers, circles are keys. Keys with a thick outline just changed owner.
Alternatives
Jump consistent hash needs no ring in memory and balances almost perfectly, but buckets can only be added or removed at the end. Rendezvous (highest random weight) hashing scores every node per key and picks the highest, which is simple and even but costs O(N) per lookup. Maglev hashing builds a fixed lookup table for very fast, stable load balancing. You'll find these in Cassandra and DynamoDB partitioning, Envoy's ring_hash and maglev balancers, and memcached client libraries.
You get
- Minimal data movement when nodes join or leave
- No central lookup service on the request path
- Natural placement for replicas
You pay
- Uneven load without virtual nodes
- Hot keys still overload one owner
- Membership changes must propagate consistently to every client
Rule of thumbWhenever a key-to-node mapping must survive nodes joining and leaving (caches, partitioned stores, sticky load balancing), use consistent hashing with virtual nodes.
Message queues
Buffer work between producers and consumers so each side can scale and fail independently.
Producers enqueue a message and return immediately; consumers pull at their own rate. This decouples availability (the email service can be down for ten minutes without losing a single signup), absorbs traffic bursts, and gives every task the same retry path. Queue depth and the age of the oldest message are the primary signals for scaling consumers.
Delivery semantics
Most brokers provide at-least-once delivery. SQS hides a received message for a visibility timeout and redelivers it if the consumer doesn't delete it in time. Kafka redelivers everything after the last committed offset when a consumer restarts. Exactly-once effects therefore require idempotent consumers that deduplicate on a message ID inside the same transaction as the side effect. Messages that keep failing move to a dead-letter queue after a maximum receive count.
Queue vs log
A queue (SQS, RabbitMQ) deletes each message once it's acknowledged and supports per-message visibility and redelivery. A log (Kafka, Kinesis) retains messages for a configured period, orders them within a partition, supports replay, and lets many consumer groups read independently. Consumer parallelism in a log is capped by the partition count.
Production checklist
- Publish from the same database transaction as the state change using the outbox pattern, avoiding dual-write inconsistencies.
- Apply backpressure: bounded queues, consumer concurrency limits, and autoscaling on the age of the oldest message.
- Retry with exponential backoff and jitter, and alert as soon as the dead-letter queue is non-empty.
- Choose the partition key for ordering (for example
order_id); ordering is guaranteed only within a partition.
You get
- Temporal decoupling and burst absorption
- Independent scaling of producers and consumers
- Uniform retries and failure isolation
You pay
- Duplicates and reordering to design for
- Eventual rather than immediate results
- A broker to operate, plus lag and DLQ monitoring
Rule of thumbAssume every message arrives at least twice and some arrive out of order. Give every message a unique ID and make every consumer idempotent before an incident forces you to.
Rate limiting
Bounds the request rate per client to protect capacity, enforce fairness, and contain abuse.
Limits are enforced per identity (API key, user ID, or tenant, falling back to IP address) at the edge or API gateway. Over-limit requests get 429 Too Many Requests with a Retry-After header and remaining-quota headers so well-behaved clients can slow down before they're rejected.
The token bucket is the most widely deployed algorithm. A bucket holds up to b tokens and refills at r tokens per second; each request consumes one. Clients can burst up to b requests while the long-run average stays at r. Try it below with b = 5 and r = 1/s.
5 of 5 tokens
| Algorithm | Implementation | Trade-off |
|---|---|---|
| Fixed window | INCR a per-window key with EXPIRE | Allows up to 2× the limit across a window boundary |
| Sliding window log | Sorted set of request timestamps per client | Exact, but memory grows with request count |
| Sliding window counter | Weighted blend of current and previous window counts | Approximate, constant memory; a common production choice |
| Token bucket | Token count plus last-refill timestamp | Allows controlled bursts; two parameters to tune |
| Leaky bucket | Bounded queue drained at a fixed rate | Smooth output, but adds latency under load |
Production checklist
- Across a fleet, keep counters in a shared store and run check-and-decrement atomically, for example as a Redis Lua script.
- Decide in advance whether to fail open or closed when the limiter's store is unavailable. Most public APIs fail open to preserve availability.
- Layer limits: per tenant globally, tighter per endpoint for expensive routes, and a concurrency cap for slow operations.
- Document the limits and expect clients to honor
Retry-Afterwith exponential backoff and jitter.
You get
- Protection from overload and noisy neighbours
- Predictable capacity planning and cost
- A first line of defense against abuse
You pay
- Legitimate traffic spikes get rejected
- Shared counter store on the hot path
- Limits to tune and communicate to clients
Rule of thumbLimit by the identity that matters (API key, user, or tenant, not only IP), apply it atomically, and always tell the client exactly when it can retry.