skip to content

Companies like Shopify and Stack Overflow ran large parts of their production systems as a single monolithic codebase (Ruby and C# respectively) at very high traffic rather than splitting into many independently deployed services. What scaling techniques let a single-deployable-unit application handle very high load, and what limits does that approach eventually run into that a monolith alone can't solve?

level: principalimportance: nice to knowfreq 40%

answer

  1. horizontal cloning behind a load balancer
  2. vertical scaling + caching as cheap first levers
  3. shared database is the hard ceiling
  4. read replicas/sharding push the ceiling, don't remove it
  5. scaling asymmetry: cloning scales the cold 80% with the hot 20%

basics

~20 s

Run many identical copies of the app behind a load balancer, plus a stronger or split-up database. This works until the shared database, or one hot feature, can't keep up no matter how many copies you run.

solid answer

~40 s

A monolith scales mainly by horizontal cloning: run many identical, stateless instances of the same artifact behind a load balancer, so requests spread across instances while each still holds the whole application in memory. Vertical scaling (bigger machines) and caching (in-memory, CDN) add headroom on top. The classic limit is the shared database: every instance talks to it, so read replicas and query/index optimization push the ceiling out, but write throughput on a single primary eventually becomes a hard bottleneck no app-instance cloning fixes. A second limit: cloning scales the whole app uniformly even when only one feature is hot, wasting resources on the cold majority just to give the hot feature enough instances — the scaling asymmetry that eventually motivates extracting specific components into their own services.

go deeper

for a junior

Should know that you can run multiple copies of an application behind a load balancer to handle more traffic.

for a middle

Should distinguish horizontal cloning, vertical scaling, and caching as separate levers, and know that the database is typically the harder part to scale.

for a senior

Should explain concretely why the database can't be cloned as trivially as the app layer, and name at least one real mitigation (read replicas, sharding).

for a principal

Should articulate the scaling-asymmetry argument for selective extraction, reference real organizations' actual trajectories, and frame monolith-vs-distributed as an incremental spectrum rather than a single binary decision.

## The primary lever: cloning the whole artifact Scaling a monolith in production relies on a small number of well-understood techniques, all of which work by adding more of the same thing rather than changing the architecture. The primary technique is **horizontal scaling by cloning**: because the application artifact is stateless with respect to individual requests (session state, if any, lives in a shared store like Redis rather than in a particular instance's memory), you can run dozens or hundreds of identical copies of the exact same JAR or binary behind a load balancer. - Each instance is a full copy of the entire application — the whole monolith, UI and order-processing and billing code all together, all loaded into that one instance's memory. - The load balancer distributes incoming requests round-robin or by some smarter policy across all of them. This is exactly the same scaling technique used for any stateless web application, monolith or not; a monolith doesn't get a special scaling mechanism, it just applies the technique uniformly to one large artifact instead of separately to several smaller ones. ## The two supplementary levers - **Vertical scaling** supplements this: giving each instance a bigger machine (more CPU, more RAM) raises the ceiling per instance before you need another one, and is often the simplest first lever because it requires no code change, just an infrastructure change. - **Caching** is the third major lever and often the highest-leverage one: an in-process or shared cache (Redis, Memcached) absorbs repeated reads that would otherwise hit the database, and a CDN absorbs static asset traffic entirely before it reaches the application layer at all. Together, cloning plus vertical scaling plus caching can take a monolith very far — Stack Overflow famously served enormous traffic for years off a comparatively small number of powerful physical servers running a single ASP.NET application, precisely by using vertical scaling and aggressive caching rather than horizontal service decomposition, because their workload (read-heavy, cacheable Q&A content) fit that model unusually well. ## Where the story hits its hardest limit The database is where this scaling story runs into its hardest, most structural limit, because a monolith almost always centers on one shared relational database that every application instance talks to. Cloning the application instances is easy because they're stateless and identical; the database is neither — it's the one place holding true mutable state, and you can't just clone a writable database and expect both copies to agree on the current state without solving distributed consensus, which is a much harder problem than cloning a stateless web server. The standard mitigations are: - **read replicas** — route read-heavy queries to replicas, keep writes on the primary; - **connection pooling**, so hundreds of app instances don't each open their own flood of raw connections; - **query and index optimization**; - sometimes **vertical scaling the database server** itself to a very large instance. But there's a ceiling: a single primary database ultimately has a maximum write throughput, and once the application's write volume approaches that ceiling, no amount of adding more application instances helps at all, because the bottleneck has moved entirely to the one component that can't be horizontally cloned the same simple way. Shopify hit exactly this class of problem around flash-sale traffic spikes (e.g., very large simultaneous product launches), which drove specific mitigations like sharding the database by shop (splitting one large database into many smaller ones, each handling a subset of merchants) rather than abandoning the monolith itself. ## Scaling asymmetry A second, subtler limit is **scaling asymmetry**: cloning the whole monolith scales every feature uniformly, because each clone contains the entire application, whether or not that feature is under load. If checkout traffic spikes ten-fold while everything else stays flat, cloning the monolith to handle checkout load means running ten times as many copies of the reporting, admin, and every other unrelated feature too, even though those features didn't need the extra capacity — wasting compute on the 80% of the app that wasn't the actual bottleneck just to get enough headroom on the 20% that was. This asymmetry is the most common concrete, evidence-based reason organizations eventually extract a specific hot component into its own independently scaled service, rather than a general dissatisfaction with monoliths. ## The resolution pattern both examples show Both Shopify and Stack Overflow illustrate the actual resolution pattern rather than a clean 'monolith versus microservices' binary: - They scaled the monolith itself as far as cloning, vertical scaling, caching, and database sharding/replication could reasonably take it. - They extracted specific components into separately scaled or separately deployed services only where a measured, specific pressure (a genuinely divergent scaling need, or a specific operational isolation need) justified the cost of doing so — not as a wholesale rewrite. The broader lesson is that 'monolith versus distributed system' isn't a single irreversible fork taken once; it's a spectrum an organization can move along incrementally, extracting the specific piece under specific pressure while leaving the rest as a single deployable unit for as long as that continues to be the better trade.

  • Why can't the database simply be cloned the same way application instances are, to remove it as a bottleneck?
    Application instances are stateless and identical, so any clone can serve any request with the same result. A database holds mutable state that changes with every write, so multiple writable copies must agree on the current state after every write — that's a distributed consensus problem, not a simple cloning problem. Real solutions (replication, sharding, distributed databases with consensus protocols) exist but are meaningfully more complex than cloning a stateless app server, which is why the database becomes the harder scaling bottleneck first.
  • What is database sharding, concretely, and how does it help a monolith's scaling ceiling?
    Sharding splits one large database into multiple smaller databases, each holding a disjoint subset of the data — for example, splitting by merchant ID so shop A's data lives in shard 1 and shop B's data lives in shard 2. This raises the effective write-throughput ceiling because each shard is a separate, independently scalable database handling only a fraction of total load, at the cost of added complexity: queries that need to join across shards, or a merchant that somehow needs data from two shards, become much harder to express and execute efficiently.
  • If scaling asymmetry (checkout needs 10x capacity, reporting doesn't) is a real cost, why wouldn't Shopify or Stack Overflow just extract every hot component into its own service preemptively?
    Because extraction has its own fixed and ongoing cost — defining a stable API contract, standing up separate deployment and monitoring, handling the new failure modes of a network call — and that cost is only worth paying once a specific component's scaling divergence is large and persistent enough to justify it. Extracting preemptively, before a component's load pattern is actually proven to diverge, risks paying distributed-systems overhead for a scaling problem that may never materialize, which is the same domain-uncertainty argument that applies to starting a new product as a monolith in the first place.

It's like scaling a restaurant by opening many identical franchise locations that all share one central warehouse: opening more identical locations is easy and handles more diners, but eventually the one shared warehouse's loading dock becomes the bottleneck no matter how many restaurants you open, and if only the drive-through window is busy, running a full duplicate restaurant just for drive-through capacity is wasteful.

saying these in an interview costs you the question

  • Thinks a monolith cannot be horizontally scaled at all
  • Doesn't recognize the shared database as the eventual hard bottleneck distinct from the application layer
  • Proposes cloning the database the same trivial way as application instances, ignoring consistency
  • Can't explain why scaling the whole monolith to serve one hot feature is wasteful
  • Has no example of a real system that scaled a monolith successfully, or assumes monoliths inherently can't handle high traffic

context