Skip to main content
Ahmed Salama

Architecture Lab

Horizontal scaling

Add more machines instead of a bigger one, once the application allows it.

The shape of the decision

System diagram: Horizontal scaling: Add more machines instead of a bigger one, once the application allows it.The stable address in front of a changing fleet.EDGELoad balancerInterchangeable instances, added and removed against a load signal.SERVICEApp ×NShared state moved out of the process so instances stay disposable.CACHESession storeUploads live here rather than on an instance disk that is about to be terminated.INFRAObject storageThe layer that does not scale by duplication, and therefore the one that usually binds first.DATADatabasedistributeshared statefilesthe real ceiling
Horizontal scaling: Add more machines instead of a bigger one, once the application allows it.
Read this diagram as text
Load balancerEdge
The stable address in front of a changing fleet.
App ×NService
Interchangeable instances, added and removed against a load signal.
Session storeCache
Shared state moved out of the process so instances stay disposable.
Object storageInfra
Uploads live here rather than on an instance disk that is about to be terminated.
DatabaseData
The layer that does not scale by duplication, and therefore the one that usually binds first.

Connections

  • Load balancer → App ×N (distribute)
  • App ×N → Session store (shared state)
  • App ×N → Object storage (files)
  • App ×N → Database (the real ceiling)
The problem it solves
Vertical scaling runs out. There is a largest instance, it is expensive out of proportion to its size, and it is still one machine that can fail.
How it works
Make instances interchangeable (no local state, no local sessions, no local uploads), then run as many as the load needs, adding and removing them automatically against a signal like CPU or queue depth.
What it costs
Statelessness is a design constraint that reaches into the whole application. Sessions move to a shared store, uploads move to object storage, background jobs need locks so two instances do not run the same one, and every scheduled task needs to be leader-elected or idempotent.
When not to reach for it
The bottleneck is the database. Adding application instances to a saturated database makes the problem arrive faster and with more connections. The fix belongs at the constrained layer, however convenient the other one is to scale.

Stateless is a prerequisite, not a preference

Horizontal scaling adds interchangeable copies instead of more capacity inside one bigger machine. A balancer sends the next request to whichever instance is free, with no memory of which one handled the request before it, so every instance has to be able to answer any request entirely on its own. I treat that as a precondition to verify before a second instance ever joins the rotation, not a gap to discover once traffic has already found it. Nothing about the balancer's routing logic can make up for an instance that only knows how to answer requests it has already seen.

The first thing to break is anything held in a process's own memory as session state. A user's next request has no guarantee of returning to the same instance, so a login flag or a shopping cart set in one process is simply absent when the balancer happens to send the following request somewhere else. It works perfectly with one instance, which is exactly why the bug does not surface until the day a second one is added and half of every user's requests start landing on a process that has never heard of them.

The same failure hits anything written to an instance's own disk. A file uploaded through one instance exists only on that instance's own storage, and cloud instances in particular treat that storage as disposable: it can vanish entirely the next time the instance is replaced, not merely become briefly unreachable. A user who uploads an avatar through one instance and requests it back through another gets a missing file, and the failure looks intermittent right up until someone works out that it correlates exactly with which instance answered which request.

An in-process cache produces a subtler version of the same problem, because it fails quietly rather than with a missing file or a lost session. Each instance warms its own copy independently, so adding instances for capacity multiplies how many separate copies of the same cache exist, and multiplies how often each one is cold. The aggregate hit rate against the origin can fall as instances are added specifically to relieve load on that origin, the opposite of what the extra capacity was meant to achieve, until the cache is moved somewhere every instance shares.

Externalising the state

Moving sessions out of process memory means choosing a store every instance can reach, typically a small, fast database or an in-memory store built for exactly this. The part that is easy to skip is deciding what happens when that store is briefly unavailable: whether an authenticated user is treated as logged out until it recovers, or whether the application accepts a request as authenticated without being able to confirm it. Either answer is defensible. Not deciding it in advance means the application discovers its own answer during the incident.

Moving uploads out of instance disk means changing the upload path itself, not just where the file eventually lands: writing directly to object storage from the request, or streaming through the instance without ever committing the file to local disk on the way. Every downstream piece of code that assumed a local file path, a thumbnail generator, a virus scanner, an export job, needs the same change, because each one was built against the assumption the upload migration just removed.

Moving a cache from in-process memory to a shared store trades one cost for another. A hit now costs a network round trip instead of a memory read, slower in absolute terms even though it is shared across the whole fleet instead of duplicated per instance. It also concentrates the stampede risk the caching concept describes elsewhere in this lab: instead of each instance separately warming a cold key, the whole fleet can now miss on the same key at the same moment and land on the origin together.

Sessions tend to come first because a broken session is a correctness bug a user notices immediately. Uploads follow because a lost file is also a correctness problem, just a rarer one. A shared cache can come last, since leaving it in process only costs speed and how evenly load spreads across the fleet, right up until traffic grows enough that the uneven spread becomes the more pressing problem. I default to that order, sessions first, then uploads, then cache, because it clears correctness bugs before spending effort on what is merely a speed and balance problem.

The bottleneck moves to the database

Making the application tier interchangeable relocates the constraint the application used to absorb instead of removing it from the system. Every instance, however many there are, still talks to the same database unless that layer has been scaled separately, so capacity that used to be spread thin across a handful of processes now converges on one place that was never redesigned when the fleet in front of it started growing.

The arithmetic that catches teams by surprise is on connections, not on query volume. Each instance opens its own pool, sized for what that one instance expects to handle concurrently, and the database sees the sum of every pool across the fleet whether or not those connections are doing anything at a given moment. An autoscaler that doubles the fleet during a burst doubles the connection count against the database in the same moment, and a database has a hard ceiling on how many connections it will accept well before it has a ceiling on how much query work it can do. The fleet can hit that ceiling and start refusing connections while the queries themselves would have been comfortably within capacity. I check that arithmetic, pool size per instance multiplied by how many instances the autoscaler could plausibly add, against the connection ceiling before I let application-tier autoscaling run unbounded.

Read replicas are the first lever, and only once the queries reaching the database are already efficient, because a replica multiplies a slow query across more machines rather than making it faster, a point the replication concept in this lab makes in more detail. Sharding is the more drastic step, splitting the data itself rather than just the read path, and it is reserved for write load a single primary genuinely cannot absorb, because resharding later is a much harder migration than adding another replica ever is.

None of this argues against scaling the application tier. It argues for treating the database's own headroom, on connections as much as on query time, as a number I check before increasing fleet size, not one I want discovered only once the fleet has already grown past what the database can hold.

Autoscaling on the right signal

CPU utilisation is the default signal for most autoscalers because it is the easiest number to get, but it measures how busy a process is, not whether a user is waiting on it. An instance blocked on a slow downstream call can sit at low CPU while every request behind it queues, because the process is idle waiting on a response rather than doing work the processor would register. A dashboard reading comfortably low CPU can be sitting directly above a fleet that is failing every user currently waiting on it.

For anything fronted by a queue, queue depth is a far more direct signal than CPU will ever be, because it states plainly how much work is waiting instead of how busy the workers currently processing it happen to be. Scaling against queue depth adds capacity at the moment work starts backing up, instead of waiting for that backlog to eventually show up as elevated CPU on whatever workers are already running.

Scaling against a high percentile of response time ties the signal directly to what a user experiences rather than to a resource reading that only correlates with it. A high percentile matters more than an average here, because an average can stay flat while a growing share of requests quietly slips past an acceptable wait, hidden by everything else in the average that is still fast.

Concurrency, the count of requests an instance is handling at any given moment, moves earlier than either of those two signals. It rises before latency visibly degrades and before CPU saturates, because a request can be counted as concurrent while it waits on a downstream call rather than consuming the processor at all. Scaling against concurrency gives the autoscaler room to add capacity ahead of the point where a user notices anything, rather than reacting once the degradation is already visible further downstream. No single signal covers every failure mode alone, and I combine the ones that track a user's actual wait, queue depth, latency, concurrency, keeping CPU utilisation as a secondary check instead of the primary trigger, since it only correlates with what the user experiences rather than measuring it directly.

Let’s talk

Have a product, platform or delivery challenge? Let’s talk about turning it into a structured, scalable solution.

Open to technical leadership, product delivery and senior engineering roles, and available for architecture consulting, technical reviews and mentorship. Engagements run as project-based work, contracts, consulting, freelance engagements, remote collaboration and long-term partnerships.

Based in Cairo, Egypt, working remotely with clients across the MENA region and internationally.

Also on LinkedIn (opens in a new tab)