What Scalability Means
Scalability is a system's ability to handle growing load — more users, more data, more requests — by adding resources, ideally with near-linear cost. A scalable system does not just survive 10x traffic; it does so without a rewrite and without its per-request latency collapsing.
Load grows along several dimensions: request rate (QPS), data volume, and read/write ratio. The right scaling strategy depends on which dimension is stressed.
The core question
Are you bottlenecked on compute, memory, storage, or network? Scale the axis that is actually saturated — measure before you shard.
Vertical vs Horizontal Scaling
Vertical (scale up) Horizontal (scale out)
one bigger machine many machines behind an LB
+---------------+ +----+ +----+ +----+
| 256 GB / 64c | | S1 | | S2 | | S3 |
+---------------+ +----+ +----+ +----+
\ | /
[ load balancer ]
| Aspect | Vertical | Horizontal |
|---|---|---|
| Effort | Trivial — resize the box | Needs stateless design + LB |
| Ceiling | Hard hardware limit | Near-unbounded |
| Fault tolerance | SPOF | Survives node loss |
| Cost curve | Superlinear at the top | Roughly linear |
Start vertical (it is cheap and simple), then go horizontal once you approach the ceiling or need fault tolerance. Most large systems use both: big-ish nodes, many of them.
Stateless Services
Horizontal scaling only works cleanly if app servers are stateless — any request can hit any server. Keep session and user state out of the server's local memory; push it to a shared store (Redis, a database) or a signed token (JWT) carried by the client.
Stateful (bad for scaling) Stateless (good)
session stored in server RAM session in Redis / JWT
-> must use sticky sessions -> any server serves any user
-> losing a node drops sessions -> add/remove nodes freely
Why it matters
Stateless servers can be autoscaled, deployed, and killed at will. State becomes a separate, deliberately-scaled concern (cache + database) instead of an accident of where a request landed.
Replication
Replication keeps copies of the same data on multiple nodes for read scaling, availability, and durability. The topology defines who can accept writes.
| Topology | Writes | Trade-off |
|---|---|---|
| Leader-follower | One leader; followers read-only | Simple; leader is write SPOF, replica lag |
| Multi-leader | Multiple leaders accept writes | Great cross-region; write conflicts to resolve |
| Leaderless | Any replica; quorum reads/writes | Highly available; needs R+W>N tuning (Dynamo) |
Read replicas & replica lag
A read-heavy workload (most web apps are ~10:1 read:write) scales beautifully with read replicas: point reads at followers, writes at the leader. The catch is replication lag — a follower may be milliseconds to seconds behind, so a user might not see their own just-written data. Fix with read-your-writes routing (read from leader for a short window after a write).
Sync replication: leader waits for follower ack -> durable, slower
Async replication: leader acks immediately -> fast, may lose
last writes on failover
Semi-sync: wait for >=1 follower -> common compromise
Sharding / Partitioning
When a dataset (or write throughput) exceeds one machine, split it into shards, each holding a subset of the data on its own node. This scales writes and storage, which replication alone cannot. The hard part is choosing the partition key.
| Strategy | How | Watch out for |
|---|---|---|
| Range | Split by key ranges (A-M, N-Z) | Great for range scans; hotspots on sequential keys |
| Hash | shard = hash(key) % N | Even spread; range scans impossible; resharding pain |
| Consistent hash | Keys + nodes on a ring | Minimal movement on resize; needs virtual nodes |
The modulo trap
With hash(key) % N, changing N remaps almost every key, forcing a massive data reshuffle. Consistent hashing exists precisely to avoid this.
Consistent Hashing
Map both keys and nodes onto a circular hash space (0 to 2³²-1). A key belongs to the first node found clockwise from its position. Adding or removing a node only reassigns the keys between it and its neighbor — on average K/N keys move, not all of them.
0 / 2^32
. Node A
k4 . . k1
. .
Node D Node B
. .
k3 . . k2
. Node C
- key k1 -> walk clockwise -> Node B owns it
- add Node E between B and C: only k2 (in that arc) moves to E
- remove Node C: its keys go to the next node clockwise (D)
Virtual nodes: place each physical node at many ring points
(e.g. 100-200 vnodes) to smooth out uneven key distribution.
Consistent hashing powers Dynamo/Cassandra partitioning, distributed caches, and many load balancers. Virtual nodes are essential in practice — without them, a few nodes end up owning oversized arcs.
The AKF Scale Cube
The AKF cube frames scaling along three independent axes. Real systems combine all three.
| Axis | Meaning | Example |
|---|---|---|
| X — cloning | Duplicate identical instances behind an LB | 10 identical web servers |
| Y — functional split | Split by service/verb (microservices) | auth, feed, search as separate services |
| Z — data split | Split by data attribute (sharding) | users A-M vs N-Z, or by region |
Denormalization
Normalized schemas avoid duplication but require joins, which are expensive across shards. Denormalization duplicates data so a read hits one place — trading write complexity and storage for read speed. It is the norm in NoSQL and read-heavy systems.
The cost
Denormalized data must be kept in sync on every write (often via events/CDC). You are trading read simplicity for write fan-out and the risk of temporary inconsistency.
Caching & CDNs
Caching stores hot results close to the reader to cut latency and offload the origin. A CDN pushes static (and increasingly dynamic) content to edge locations near users, turning a 150 ms cross-continent trip into a <20 ms edge hit.
User -> CDN edge (static assets, cached pages)
-> Load balancer
-> App server -> Cache (Redis) --hit--> fast return
--miss-> DB, then populate cache
Common caching patterns: cache-aside (app reads cache, falls back to DB, populates), write-through (write cache + DB together), and write-back (write cache, flush to DB async). Always set a TTL and plan for invalidation.
Asynchronous Processing
Not everything needs to happen inside the request. Push slow or bursty work (email, image resizing, analytics, fan-out) onto a message queue and process it with worker pools. This decouples producers from consumers, absorbs spikes, and keeps user-facing latency low.
Client -> API (returns 202 fast) -> [ Kafka / SQS queue ] -> workers
|
scale workers independently
of request throughput
Hotspots & the Celebrity Problem
Even with good sharding, load is rarely uniform. A single hugely popular key — a celebrity's account, a viral post, Black Friday's top SKU — concentrates traffic on one shard: the hotspot or celebrity problem.
| Mitigation | How it helps |
|---|---|
| Key salting / splitting | Append a suffix to spread a hot key across shards |
| Dedicated caching | Serve hot items from a replicated in-memory cache |
| Hybrid fan-out | Push feeds for normal users, pull for celebrities |
| Request coalescing | Collapse identical concurrent misses into one origin call |
Capacity Planning: A Worked Example
How many app servers and cache nodes for a read-heavy service at 30,000 peak read QPS?
Given:
peak read QPS = 30,000
per-request CPU work = ~5 ms of a single core
server = 16 cores usable
Throughput per server = 16 cores / 0.005 s = 3,200 req/s
Servers needed = 30,000 / 3,200 ~= 10 servers
+50% headroom (bursts, deploys, node loss) -> ~15 servers
Cache sizing:
working set = 5M hot objects x 2 KB = 10 GB
x2 for replication + overhead = ~20 GB
-> a couple of Redis nodes with room to spare
Always add headroom
Never provision for exactly peak. Leave 30-50% for traffic spikes, rolling deploys, and losing a node or an availability zone.