Contents

Architecture & System Design › System Design Fundamentals

Horizontal Scaling

Scaling out with more machines.

Also known as: scale out, scaling out, horizontal scalability, adding more servers, scale horizontally

Horizontal scaling (scaling out) means handling more load by adding more machines, instead of making one machine bigger (vertical scaling, scaling up).

Vertical:    one 8-core server  →  one 64-core server
Horizontal:  one server         →  ten servers behind a load balancer
VerticalHorizontal
HowBigger CPU, more RAMMore machines
SimplicityVery simple: no code changesRequires the design to support it
LimitCeiling at the biggest available machine, and cost rises steeplyPractically much higher
ResilienceOne machine is a single point of failureA failed machine is replaced by others
Downtime to scaleOften a restartAdd nodes without stopping

See vertical scaling and scalability.

What makes an app scale out easily

A stateless application tier. If servers keep no per-user state locally, any server can handle any request, and you can add or remove them freely behind a load balancer. State goes into shared stores: a database, a cache, a session store (stateless services).

clients ─► load balancer ─► app server 1 ┐
                       ├─► app server 2 ├─► shared database / cache
                       └─► app server 3 ┘

Autoscaling adds and removes instances based on load.

The hard part: the data tier

Stateless servers are easy to multiply. Data isn’t:

  • Read-heavy: add read replicas (read replicas, replication).
  • Write-heavy or too big for one machine: partition the data across machines (sharding, with its trade-offs: cross-shard queries, rebalancing, hot spots).
  • Caching reduces load before you need either (caching).

Costs and cautions

  • Complexity: distributed systems have partial failures, network issues and consistency questions.
  • Don’t scale out too early. A single well-tuned server (and a replica) handles a lot. Measure and find the real bottleneck first.
  • Find the actual bottleneck: adding web servers does nothing if the database is saturated.
  • Things that silently break: in-memory sessions, local file uploads, scheduled jobs that now run on every instance, in-process caches that diverge.
  • Cost: more machines means more to run, monitor and pay for.

A common path is: optimize → scale up → add caching and replicas → scale out the stateless tier → partition data only when necessary.