Skip to content
BytePatterns

Vertical vs Horizontal Scaling: Scale Up or Scale Out?

8 min readBytePatterns

Vertical vs horizontal scaling explained: the ceiling of one big machine, why scaling out needs stateless servers, the availability math, and when to switch.

Almost every system design interview reaches the same fork within five minutes: traffic is growing, so do you buy a bigger machine or more machines? Vertical scaling (scaling up) gives one server more CPU, memory or disk. Horizontal scaling (scaling out) adds servers and spreads the work across them. The textbook answer is "scale out"; the good answer explains what that costs, why it needs stateless servers, and where the bottleneck moves next.

The problem it solves

One server has a fixed capacity. When requests arrive faster than it can answer, queues grow and requests time out. You have three levers:

  • Make the code cheaper first: an index or a cache can buy more than new hardware.
  • Scale up: move to a machine with more resources. No code changes, and everything stays on one box.
  • Scale out: run several identical machines behind a load balancer, which picks one for each request.

Most real systems scale up first, because it is simple, and scale out when one machine is no longer enough or no longer safe.

The intuition

Scaling up is attractive because the software does not change: same process, same memory, no network in between. Its limits are physical and operational:

  • A ceiling. There is a largest machine you can buy or rent, and you reach it eventually.
  • Diminishing returns. Four times the cores rarely means four times the throughput: locks, memory bandwidth and the parts of the work that cannot run in parallel all get in the way.
  • Downtime to resize. Moving to a bigger instance usually means a restart.
  • One failure domain. However big the box, when it dies, everything dies.

Scaling out removes the ceiling for the tier you scale and turns a machine failure into lost capacity instead of lost service. It asks for one thing in return: any machine must be able to answer any request. If a user's session lives in one server's memory, the balancer must keep sending that user to the same server, and a failure logs them out. So session data, uploads and caches move to shared stores: a database, a cache cluster, object storage.

That shared store is where the bottleneck goes next. Ten stateless web servers in front of one database are limited by the database, and scaling the data tier needs different tools: replication for reads and sharding for writes.

Watch it run

The animation starts with one machine answering everything, full at about 800 requests a second. Scaling up gives that same machine more CPU, memory and disk, with no code changes at all: 8 vCPU and 1,600 requests a second. Twice the box again reaches 3,100, and the readouts show why it is not free: resizing usually costs downtime, and four times the box did not buy four times the throughput. Then it stops: there is no larger box to rent, the hard ceiling. And one machine is still one failure that takes the whole service with it, uptime 0%. Scaling out adds machines instead and spreads the work across them: a balancer and three servers, 2,400 requests a second. It only works if a request can land on any machine, so session state moves out of the server into a shared store. Each request then goes to whichever machine the balancer picks. Lose one machine and you lose a third of the capacity, not the service. Add a fourth and a fifth: no ceiling for this tier, paid for with stateless servers and routing.

Vertical vs Horizontal Scaling

Step 1 of 10

One machine answers everything. It is full at about 800 requests a second.

The same interactive animation as the lesson — step through it with the controls.

The code

A toy model: servers are Python objects, a round-robin balancer spreads requests, and a login is a session that later requests must find. With sessions kept in each server's memory, most requests land on a server that never saw the login:

import itertools

class Server:
    def __init__(self, name, sessions=None):
        self.name = name
        self.sessions = {} if sessions is None else sessions   # local memory, or a shared store

    def handle(self, user, action):
        if action == "login":
            self.sessions[user] = "token"
            return "ok"
        return "ok" if user in self.sessions else "401"         # the session must be found

def run(servers, requests):
    pick = itertools.cycle(servers)                             # round-robin balancer
    return [next(pick).handle(user, action) for user, action in requests]

requests = [(u, "login") for u in "abc"] + [(u, "view") for u in "abc" for _ in range(3)]
local = run([Server(f"s{i}") for i in range(3)], requests)
shared_store = {}
shared = run([Server(f"s{i}", shared_store) for i in range(3)], requests)
print(local.count("401"), "of", len(local) - 3, "views failed with local sessions")
# 6 of 9 views failed with local sessions
print(shared.count("401"), "of", len(shared) - 3, "views failed with a shared store")
# 0 of 9 views failed with a shared store

Capacity and availability for the two designs, with illustrative numbers: every machine is up 99% of the time, a small server serves 800 requests a second, and the service needs 1,600. One big box is up exactly as often as that box. Three small ones survive any single failure:

from math import comb

def p_at_least(k, n, a):
    """Probability that at least k of n independent machines are up."""
    return sum(comb(n, i) * a**i * (1 - a) ** (n - i) for i in range(k, n + 1))

a, per_server, needed = 0.99, 800, 1_600
k = -(-needed // per_server)                       # servers required: ceil(1600 / 800) = 2
print(f"one big box:     {a:.4%}")                 # one big box:     99.0000%
for n in (2, 3, 4):
    print(f"{n} small servers: {p_at_least(k, n, a):.4%}")
# 2 small servers: 98.0100%
# 3 small servers: 99.9702%
# 4 small servers: 99.9996%

def capacity(n_web, per_web=800, db_limit=3_000):
    """Stateless web tier in front of one shared database: the smaller limit wins."""
    return min(n_web * per_web, db_limit)

print([capacity(n) for n in (1, 2, 3, 4, 6, 10)])  # [800, 1600, 2400, 3000, 3000, 3000]

Two servers with no spare are less available than one box: either failure breaks the service. The win comes from N + 1, and past four web servers the database caps everything. The binomial formula checked against a brute force that enumerates every up/down combination, on 500 seeded random fleets:

import random

def brute_force(k, n, a):
    total = 0.0
    for states in itertools.product([True, False], repeat=n):  # every combination of up/down
        p = 1.0
        for up in states:
            p *= a if up else 1 - a
        if sum(states) >= k:
            total += p
    return total

rng = random.Random(31)
ok = True
for _ in range(500):
    n = rng.randint(1, 12)
    k, a = rng.randint(0, n), rng.uniform(0.5, 1.0)
    ok &= abs(p_at_least(k, n, a) - brute_force(k, n, a)) < 1e-12
print(ok)                                          # True

The complexity

The costs here are operational:

  • Scale up: no new moving parts, limited by the largest machine, one failure domain, and a restart to resize.
  • Scale out: a balancer, health checks, deployment across machines and shared state stores; in exchange, capacity grows by adding machines and one failure costs 1/N of the capacity.
  • Rule of thumb: size the fleet as "servers needed for peak load, plus at least one", because N servers with no spare fail as soon as any one does.

Where it goes wrong

  • Sticky sessions as the fix. Pinning each user to one server hides the state problem until that server dies or you deploy. Move the state out instead.
  • Scaling the wrong tier. More web servers in front of a saturated database add connections, not capacity. Measure what is actually full first: CPU, memory, disk or network.
  • Local disk and local caches. Uploaded files written to one server's disk are missing on the others; per-server caches return different answers for the same key. Use object storage and a shared cache.
  • Scaling out too early. Distributed systems are harder to debug and deploy. If one machine fits the foreseeable load, scaling up is the honest answer.

When it shows up in interviews

In the first minutes of nearly every system design question, usually as "your traffic grows tenfold, what changes?". The follow-ups go straight to the consequences: how the load balancer picks a server, where sessions live, and what happens to the database. The system design interview questions article collects more of them.

How to say it in an interview

"I'd start by measuring what's saturated. If one machine can still handle the load, scaling up is the simplest move: no code changes, but there's a ceiling, it usually costs a restart, and it's a single point of failure. To scale out, I'd make the application servers stateless by moving sessions to a shared store like a cache cluster, put them behind a load balancer with health checks, and run peak capacity plus at least one spare. Then the bottleneck moves to the database, which I'd handle separately with read replicas and caching, and sharding if writes outgrow one primary."