Learn/system design/Scalability
Beginner~25 min read

Scalability

Vertical vs horizontal scaling, stateless services, replication, sharding, consistent hashing, the AKF scale cube, and capacity planning.

ShardingReplicationConsistent HashingAKF Cube

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 ]
AspectVerticalHorizontal
EffortTrivial — resize the boxNeeds stateless design + LB
CeilingHard hardware limitNear-unbounded
Fault toleranceSPOFSurvives node loss
Cost curveSuperlinear at the topRoughly 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.

TopologyWritesTrade-off
Leader-followerOne leader; followers read-onlySimple; leader is write SPOF, replica lag
Multi-leaderMultiple leaders accept writesGreat cross-region; write conflicts to resolve
LeaderlessAny replica; quorum reads/writesHighly 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.

StrategyHowWatch out for
RangeSplit by key ranges (A-M, N-Z)Great for range scans; hotspots on sequential keys
Hashshard = hash(key) % NEven spread; range scans impossible; resharding pain
Consistent hashKeys + nodes on a ringMinimal 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.

AxisMeaningExample
X — cloningDuplicate identical instances behind an LB10 identical web servers
Y — functional splitSplit by service/verb (microservices)auth, feed, search as separate services
Z — data splitSplit 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.

MitigationHow it helps
Key salting / splittingAppend a suffix to spread a hot key across shards
Dedicated cachingServe hot items from a replicated in-memory cache
Hybrid fan-outPush feeds for normal users, pull for celebrities
Request coalescingCollapse 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.

Section navigation