System Design (High-Level Design) Module 3 – Performance and reliability fundamentals
Scaling up, scaling out and stateless services
Compare scaling up with scaling out on cost, limits and failures, move session state out of app servers, and spot the state that needs one owner.
What you will learn
- Compare vertical and horizontal scaling on cost, limits and failure behaviour
- Move session state out of application servers so any instance can serve any request
- Recognise state that cannot simply be scaled out, and partition it instead
Before you start
On this page
When one server runs out of capacity there are two ways to grow. Scaling up, or vertical scaling, moves the service to a bigger machine with more cores and memory. Scaling out, or horizontal scaling, runs more copies of it behind a load balancer. Scaling up is simpler and often the cheapest first step, but it stops at the largest machine you can get, and that machine is still a single point of failure. Scaling out is not limited by the size of one machine and survives the loss of one, but only if the copies are stateless: they keep nothing between requests that another copy would need. This lesson compares the two, shows what breaks when state stays on the servers, and names the kinds of state that no number of copies can scale.
Scaling up and scaling out
The Kubernetes documentation draws the line in its own terms: horizontal scaling answers more load by running more Pods, and vertical scaling gives the Pods already running more memory or CPU. The same split applies to virtual machines, databases and anything else that runs on hardware.
| Criterion | Scale up (vertical) | Scale out (horizontal) |
|---|---|---|
| What changes | A bigger machine for the same program | More copies of the program, behind a load balancer |
| Code changes | None | State must leave the instances first |
| Ceiling | The largest machine available, and its price | Shared resources: the database, locks, the network |
| When one machine fails | The service is down until it is replaced | The service loses that machine's share of capacity |
| Deploys | A restart takes the whole service down | Instances are replaced one at a time |
| Running cost of the setup | One machine to watch | A load balancer, health checks, shared stores |
| When to choose | A database or a young service, until the machine or one failure domain becomes the risk | Stateless request handling that must grow and survive failures |
Most real systems use both. Databases are usually scaled up first, because splitting their data is hard, and the stateless tier in front of them is scaled out because it is easy. Scaling out also buys reliability that scaling up cannot: with eight instances, losing one costs an eighth of the capacity, and with enough headroom users notice nothing. The growth lesson followed a service through these stages; this lesson looks at what has to be true for the scale-out step to work.
Stateless instances, and what breaks without them
An instance is stateless when it keeps nothing between requests that another instance would need: no sessions in its memory, no uploaded files on its disk, no counters that only it knows. The Twelve-Factor App states it as a rule: processes share nothing, anything that must persist goes to a backing service, and an instance’s memory or disk is only a short cache for a single request. Here is a shopping cart that keeps its sessions in the instance’s memory, run on one instance, then on two behind a round-robin load balancer, then on two that share a session store:
// One shopping-cart handler, run on one instance and on two behind a round-robin load balancer: first with the
// sessions in each instance's memory, then in a shared store. A Map that holds JSON text stands in for the shared
// store (Redis or a database), so each request reads the session and writes it back, as it would over the network.
function makeInstance(sharedStore) {
const memory = new Map(); // this instance's own sessions: gone on restart, invisible to other instances
const store = sharedStore ?? memory;
const load = (id) => (store.has(id) ? JSON.parse(store.get(id)) : null);
const save = (id, session) => store.set(id, JSON.stringify(session));
return {
restart: () => memory.clear(),
handle({ path, sessionId, item }) {
if (path === "/login") {
save("s-asha", { user: "asha", cart: [] });
return "200 logged in";
}
const session = load(sessionId);
if (!session) return "401 please log in";
if (path === "/cart/add") {
session.cart.push(item);
save(sessionId, session);
return `200 ${session.cart.length} item(s)`;
}
return `200 cart: ${session.cart.join(", ") || "empty"}`;
},
};
}
const steps = [
{ path: "/login" },
{ path: "/cart/add", sessionId: "s-asha", item: "tea" },
{ path: "/cart/add", sessionId: "s-asha", item: "rice" },
{ path: "/cart", sessionId: "s-asha" },
{ deploy: true }, // every instance restarts with the new version
{ path: "/cart", sessionId: "s-asha" },
];
function run(title, instances) {
console.log(title);
let next = 0; // round robin: each request goes to the next instance in turn
for (const step of steps) {
if (step.deploy) {
instances.forEach((instance) => instance.restart());
console.log(" -- deploy: every instance restarts --");
continue;
}
const n = next++ % instances.length;
console.log(` ${step.path.padEnd(10)} on instance ${n + 1}: ${instances[n].handle(step)}`);
}
}
run("One instance, sessions in memory", [makeInstance()]);
run("Two instances, sessions in memory", [makeInstance(), makeInstance()]);
const shared = new Map();
run("Two instances, sessions in a shared store", [makeInstance(shared), makeInstance(shared)]); Output
One instance, sessions in memory /login on instance 1: 200 logged in /cart/add on instance 1: 200 1 item(s) /cart/add on instance 1: 200 2 item(s) /cart on instance 1: 200 cart: tea, rice -- deploy: every instance restarts -- /cart on instance 1: 401 please log in Two instances, sessions in memory /login on instance 1: 200 logged in /cart/add on instance 2: 401 please log in /cart/add on instance 1: 200 1 item(s) /cart on instance 2: 401 please log in -- deploy: every instance restarts -- /cart on instance 1: 401 please log in Two instances, sessions in a shared store /login on instance 1: 200 logged in /cart/add on instance 2: 200 1 item(s) /cart/add on instance 1: 200 2 item(s) /cart on instance 2: 200 cart: tea, rice -- deploy: every instance restarts -- /cart on instance 1: 200 cart: tea, rice
Recorded with Node.js 24.21.0 on macOS 26 arm64. To run it yourself: mise exec node@24.21.0 -- node sessions.mjs
Runs on this device, in your browser. The first run downloads JavaScript (about 0.6 MB), which is kept for the next runs.
Your run, in this browser
On one instance everything works until the deploy, which restarts the process and logs the user out with an empty cart. On two instances it breaks at once: the second request lands on an instance that never saw the login, so it answers 401, and the user’s “tea” is never added. Only the shared store works throughout, because every instance reads the session from the same place and writes it back.
The usual workaround, sticky sessions, makes the load balancer send each user back to the same instance. It hides the problem without solving it: a restart still loses the sessions on that instance, new instances get no traffic from existing users, and one busy user pins load to one machine. The Twelve-Factor App rules them out entirely and recommends a datastore whose entries expire on their own. The quality-attributes lesson compared the three places a session can live: the instance’s memory, a shared store, or a signed token that the client sends back.
JWT Decoder Decode a signed token to see the claims a stateless server checks on every request.The same reasoning places the rest of an application’s state:
- Sessions and carts: a shared store with expiry, or the database for carts that must survive for days.
- Uploaded files: object storage, with the file’s key in the database. A file on one instance’s disk is gone when that instance is replaced.
- Caches: a local cache on each instance is fine, because it is disposable and every instance can rebuild it; a cache that must be shared goes into a cache cluster.
- Background work: a queue, so that any worker can take the next job and a crashed worker’s job goes back to the queue.
Stateless instances keep everything that must last in shared stores
Text description of the diagram
The diagram runs from top to bottom.
- A load balancer receives every request and may send it to any instance.
- Three app instances sit side by side. They keep nothing between requests that another instance would need.
- Below them are the three places where state that outlives a request lives: a session store whose entries expire, object storage for uploads and other files, and the database primary, which is the one place that accepts writes.
Because the instances hold no state of their own, any of them can be added, replaced or redeployed without users noticing. The database primary is different: there is only one writer, so adding app instances does not add write capacity.
Autoscaling a stateless fleet
Once any instance can serve any request, the number of instances can follow the load. Kubernetes’ horizontal autoscaler does it with one rule: the desired number of replicas is the current number times the current metric divided by its target, rounded up, and it changes nothing while that ratio is within 0.1 of 1.0 by default. The loop runs every 15 seconds by default, and scaling down waits for a stabilisation window of 5 minutes, so that a brief dip does not remove instances that the next minute needs:
"""How a horizontal autoscaler picks the number of instances, following the rule the Kubernetes documentation gives
for the HorizontalPodAutoscaler: desired = ceil(current x current metric / target metric), with no change while the
ratio is within a tolerance of 0.1 of 1.0. The CPU readings are invented for the sketch; the clamp to a minimum
and a maximum stands for the autoscaler's configured bounds.
"""
import math
TARGET_CPU = 60 # the average CPU utilisation per instance the autoscaler aims for, in percent
TOLERANCE = 0.1
MIN_REPLICAS, MAX_REPLICAS = 2, 20
def desired(current, cpu):
ratio = cpu / TARGET_CPU
if abs(ratio - 1) <= TOLERANCE:
return current # close enough: no change
return max(MIN_REPLICAS, min(MAX_REPLICAS, math.ceil(current * ratio)))
replicas = 4
print(f"{'time':>6}{'avg CPU':>9}{'ratio':>7} replicas")
for time, cpu in [("09:00", 55), ("09:15", 64), ("09:30", 90), ("09:45", 75), ("10:00", 62), ("10:15", 30), ("10:30", 20)]:
new = desired(replicas, cpu)
change = "no change" if new == replicas else f"{replicas} -> {new}"
print(f"{time:>6}{cpu:>8}%{cpu / TARGET_CPU:>7.2f} {change}")
replicas = new Output
time avg CPU ratio replicas 09:00 55% 0.92 no change 09:15 64% 1.07 no change 09:30 90% 1.50 4 -> 6 09:45 75% 1.25 6 -> 8 10:00 62% 1.03 no change 10:15 30% 0.50 8 -> 4 10:30 20% 0.33 4 -> 2
Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 autoscale.py
Runs on this device, in your browser. The first run downloads Python (about 13.5 MB), which is kept for the next runs.
Your run, in this browser
At 90% CPU against a 60% target the ratio is 1.5, so four replicas become six; at 30% the ratio is 0.5 and eight become four. The 64% reading changes nothing, because 1.07 is within the tolerance. The rule assumes that load spreads evenly over the replicas, which is exactly what statelessness makes true. It also needs headroom, because new instances take time to start: if traffic can double in less time than an instance takes to boot and warm up, the fleet must already be large enough to absorb the jump.
State that does not scale out
Some state cannot be copied onto every instance, because only one place may own it at a time. Adding instances in front of it adds requests that wait for the owner, not throughput.
- A single-writer database. A PostgreSQL standby in hot standby answers read-only queries; writes go to the primary alone. Read replicas scale reads, never writes.
- Locks and hot rows. A row that every request updates, or a lock that every request takes, is held by one request at a time.
- Ordered streams. Kafka guarantees that consumers read a partition’s events in the order they were written, and events with the same key go to the same partition. Order holds per partition, so ordered processing for one key happens in one place at a time.
- Live shared state in memory. The players of a game room, or the members of a live call, must meet on the same server, so requests are routed to the room’s owner rather than to any instance.
The program below puts numbers on the second case. Every request holds a shared lock for 2 ms, so the lock passes at most 500 requests a second:
"""Throughput against instance count, with and without a contended section.
Each instance serves up to 200 requests a second. In the second workload every request also holds one shared
lock (a hot database row, a global counter) for 2 ms, and only one request can hold it at a time, so the lock
serves at most 1,000 / 2 = 500 requests a second however many instances queue for it. The numbers are assumptions.
"""
PER_INSTANCE = 200 # requests a second one instance serves
LOCK_MS = 2 # milliseconds each request holds the shared lock
LOCK_LIMIT = 1000 / LOCK_MS # requests a second the lock can pass
print(f"{'instances':>9}{'no shared lock':>16}{'with the lock':>15}{'per instance':>14}")
for n in (1, 2, 3, 4, 8, 16, 32):
free = n * PER_INSTANCE
locked = min(free, LOCK_LIMIT)
print(f"{n:>9}{free:>12,} r/s{locked:>11,.0f} r/s{locked / n:>10.0f} r/s")
# The wait at the lock, for one lock with random arrivals and service times (the M/M/1 queue of the estimation
# module): utilisation rho = rate x 2 ms, mean wait = rho / (1 - rho) x 2 ms, unbounded once rho reaches 1.
print("\nThe queue in front of the lock")
print(f"{'instances':>9}{'offered':>10}{'lock busy':>11}{'mean wait':>12}")
for n in (1, 2, 3, 4):
offered = n * PER_INSTANCE
rho = offered * LOCK_MS / 1000
wait = f"{rho / (1 - rho) * LOCK_MS:.1f} ms" if rho < 1 else "grows without limit"
print(f"{n:>9}{offered:>6,} r/s{min(rho, 1):>11.0%} {wait:>10}") Output
instances no shared lock with the lock per instance
1 200 r/s 200 r/s 200 r/s
2 400 r/s 400 r/s 200 r/s
3 600 r/s 500 r/s 167 r/s
4 800 r/s 500 r/s 125 r/s
8 1,600 r/s 500 r/s 62 r/s
16 3,200 r/s 500 r/s 31 r/s
32 6,400 r/s 500 r/s 16 r/s
The queue in front of the lock
instances offered lock busy mean wait
1 200 r/s 40% 1.3 ms
2 400 r/s 80% 8.0 ms
3 600 r/s 100% grows without limit
4 800 r/s 100% grows without limit
Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 scaling_curve.py
Runs on this device, in your browser. The first run downloads Python (about 13.5 MB), which is kept for the next runs.
Your run, in this browser
Without the lock, throughput grows with the fleet: 32 instances serve 6,400 requests a second. With it, the third instance adds only 100 of its 200 requests a second, every instance after it adds nothing, and by 32 instances each one does 16 requests a second of useful work. The second table shows why. The lock is a queue, and its mean wait follows the M/M/1 formula of the Little’s law lesson: 1.3 ms at 40% busy, 8 ms at 80%, and once the offered load passes 500 a second, the queue grows for as long as the load lasts.
The way out is the same for every item on the list: partition the state, so that each piece has one owner and different pieces live in different places. Shard the database by customer so each shard has its own writer, split a hot counter into several counters that are summed when read, partition the stream by key, and assign each game room to one server. Then adding servers adds owners, and owners add throughput. The partitioning module of this track covers how to choose the keys.
Practice
Exercise · Easy · Python
What adding instances can and cannot buy
Write three helpers in scaling.py for a fleet of identical instances.
capacity(instances, per_instance_rps, lock_ms) returns the requests a second the fleet can serve when every request also holds one shared lock for lock_ms milliseconds. The instances together serve instances × per_instance_rps, and the lock passes at most 1000 / lock_ms requests a second, so the fleet serves the smaller of the two. A lock_ms of 0 means there is no shared lock.
instances_needed(target_rps, per_instance_rps, lock_ms) returns the fewest instances that serve target_rps, or None when the lock caps the fleet below the target, because then no number of instances is enough.
survives_loss(instances, per_instance_rps, peak_rps, lost=1) returns True when the instances left after lost of them fail can still serve the peak.
For example, capacity(8, 200, 2) is 500, instances_needed(450, 200, 2) is 3, instances_needed(600, 200, 2) is None, and survives_loss(4, 200, 700) is False.
Starter code · scaling.py
import math
def capacity(instances, per_instance_rps, lock_ms):
"""Requests a second the fleet serves when every request holds one shared lock for lock_ms."""
# Replace this line with your code.
return 0
def instances_needed(target_rps, per_instance_rps, lock_ms):
"""Fewest instances that serve target_rps, or None when the lock caps the fleet below it."""
# Replace this line with your code.
return 0
def survives_loss(instances, per_instance_rps, peak_rps, lost=1):
"""Whether the instances left after `lost` failures still serve the peak."""
# Replace this line with your code.
return False The sample tests · test_scaling.py
import math
from scaling import capacity, instances_needed, survives_loss
def test_capacity():
"""is the smaller of what the instances serve and what the lock passes"""
assert math.isclose(capacity(8, 200, 2), 500)
assert math.isclose(capacity(2, 200, 2), 400)
assert math.isclose(capacity(3, 200, 0), 600)
assert math.isclose(capacity(10, 100, 4), 250)
def test_instances_needed():
"""rounds up, and is None when the lock caps the fleet below the target"""
assert instances_needed(450, 200, 2) == 3
assert instances_needed(400, 200, 2) == 2
assert instances_needed(500, 200, 2) == 3
assert instances_needed(600, 200, 2) is None
assert instances_needed(10_000, 250, 0) == 40
def test_survives_loss():
"""compares what the remaining instances serve with the peak"""
assert survives_loss(4, 200, 700) is False
assert survives_loss(5, 200, 700) is True
assert survives_loss(5, 200, 800) is True
assert survives_loss(6, 200, 800, lost=2) is True
assert survives_loss(6, 200, 801, lost=2) is False A hint
Treat a lock_ms of 0 as no limit, for example with math.inf. For the instance count, check the lock's limit first, then divide the target by one instance's rate and round up with math.ceil. For a loss, multiply the instances that are left by one instance's rate and compare with the peak.
Results of the sample tests
| Test | Result | Details |
|---|
What your code printed
The sample tests run on this device, in your browser (Pyodide): nothing is sent to mysmartcopilot.com. The first run downloads Python (about 13.5 MB), which is kept for the next runs. A check in your browser is feedback for you, not proof that the code is right for every input.
Check yourself
6 questions about this lesson. Every answer and why it is right is on the page, behind “Show the answer”. Your score stays in this browser.
Interview questions
Warm-up (fresher to mid level): what is the difference between scaling up and scaling out, and when is scaling up the better first move? Scaling up gives one machine more cores, memory or faster disks; scaling out adds more machines that share the work behind a load balancer. Scaling up needs no code changes and no new components, so it is often the better first move for a young service or for a database whose data is hard to split. Its limits are the largest machine available and the fact that one machine is one failure. Scaling out avoids both limits, but only for stateless instances, so it usually comes first for the request-handling tier and later for the data tier.
Key takeaways
- Scaling up is simple and needs no code changes, but has a ceiling and one failure domain; scaling out avoids both, at the price of a load balancer and stateless instances.
- Stateless means nothing on an instance that another instance would need: sessions, uploads and jobs live in shared stores, and local caches are only disposable copies.
- An autoscaler sets replicas to ceil(current × metric ÷ target) and only works because statelessness spreads load evenly.
- A single writer, a hot lock, an ordered stream or a live room has one owner; adding instances in front of it adds queueing, not throughput.
- Partition such state so each piece has its own owner; then more servers mean more owners.
References
- Horizontal Pod Autoscaling (Kubernetes documentation) (The Kubernetes Authors)
- The Twelve-Factor App: VI. Processes (The Twelve-Factor App (Adam Wiggins))
- PostgreSQL documentation: Hot Standby (The PostgreSQL Global Development Group)
- Introduction (Apache Kafka documentation) (The Apache Software Foundation)
- Discrete Stochastic Processes, chapter 4: Renewal processes (M/G/1 queueing delay) (MIT OpenCourseWare (Robert Gallager))
Related tools
Report a problem with this lesson
Kept only in this browser. Your Learn progress