skip to content

What is the difference between scaling a relational database vertically and horizontally, and what does each approach actually buy you?

level: juniorimportance: must knowfreq 70%

answer

  1. scale-up = one bigger box; scale-out = more boxes
  2. stateful: copy = replica, split = shard
  3. replicas add reads + HA, never writes
  4. single-primary write ceiling
  5. cost curve bends at the top of the instance range

basics

~20 s

Vertical scaling means one bigger machine: more CPU, RAM, faster disks. Horizontal scaling means more machines: read replicas to spread reads, or shards to split data across independent primaries. Vertical is simple but capped; horizontal is unbounded but changes the programming model.

solid answer

~50 s

**Vertical (scale-up)** = give the single database server more CPU cores, more RAM, faster NVMe/IOPS. Nothing about the application changes: one machine, one copy of the data, full transactional semantics, joins and constraints work everywhere. It is the cheapest engineering answer and usually the right first move — modern single boxes handle a lot. **Horizontal (scale-out)** = add machines. Two flavours, and they solve different problems: - **Replicas** — copies of the same data. They add read capacity and failover headroom, not write capacity. - **Shards** — the data is split by key across independent primaries. This adds write capacity, but you lose cheap cross-shard joins, global uniqueness, and single-node transactions. The key asymmetry: a stateless app tier scales out by adding identical copies; a database holds state, so a copy is either redundant (replica) or a partial view (shard). That is why vertical scaling is done first and scale-out is a deliberate architectural commitment.

go deeper

for a junior

Define both terms cleanly and give one example each: a bigger instance versus adding replicas or shards. Mention that a database holds data, so adding machines means either copying or splitting it.

for a middle

Add the asymmetry: replicas scale reads and availability, shards scale writes, and the single primary is the write ceiling. Note the cost and complexity ordering.

for a senior

Frame it as a sequence of cheaper interventions before topology changes, and be concrete about what breaks when you shard: joins, uniqueness, and multi-row transactions across shards.

for a principal

Talk about failure domains, cost curves, the organizational cost of operating N primaries, and the point at which a workload should leave the relational system rather than be sharded inside it.

## The two words **Vertical scaling (scale-up)** means making one machine more powerful: more CPU cores, more RAM, faster and more parallel storage, more network bandwidth. The database topology does not change — there is still one server holding one copy of the data. **Horizontal scaling (scale-out)** means adding more machines and dividing work among them. For a stateless service (a web app) that is trivial: every instance is interchangeable, so a load balancer can send any request anywhere. A database is *stateful*: it owns data on disk. So "add a machine" has to answer the question "and what data does that machine hold?" There are only two answers. ## Answer one: replicas (the same data, more copies) A replica continuously applies the primary's change stream and serves a (usually slightly stale) copy. Replicas add: - **read throughput** — read-only queries can be routed away from the primary; - **availability** — a replica can be promoted if the primary dies; - **isolation** — reporting queries stop competing with transactional traffic. What replicas do *not* add is write throughput. Every write must still be executed on the primary, and then re-executed or replayed on every replica. Adding a replica adds a consumer of the write stream, not a producer. In fact each replica costs the primary a little more work (shipping and acknowledging the log). ## Answer two: shards (different data on each machine) Sharding splits rows across N independent databases by some key (customer, tenant, account, hash of an id). Each shard is a full-fledged primary that accepts writes for its slice, so aggregate write capacity really does scale roughly linearly with shard count. The price is paid in semantics, not hardware: - a query that spans shards needs a scatter-gather layer and pays the slowest shard's latency; - foreign keys, unique constraints, and joins only hold *within* a shard; - a transaction touching two shards needs distributed commit or must be redesigned to avoid it; - rebalancing when one shard grows hot is a real operational project. ## The single-writer ceiling Because of the above, a classic single-primary relational deployment has a hard ceiling: the write throughput of one machine. You can grow that ceiling (faster disks, more cores, group commit, batching) but you cannot multiply it by adding servers. Every scale-out story for writes is ultimately "more primaries", i.e. sharding — or moving that workload out of the relational system entirely. ## Where vertical scaling stops Scale-up is limited by four practical things: 1. **Physical maximums** — the largest instance a cloud sells, or the largest box you can rack. It is a big number, but finite. 2. **The cost curve** — price per core/GB is roughly linear up to the middle of the range and then rises sharply at the top. The last doubling can cost several times the previous one. 3. **Diminishing returns** — more cores do not help a workload bottlenecked on a single hot lock, a single-threaded phase, or fsync latency. Adding RAM does nothing once the working set already fits. 4. **Blast radius and downtime** — one enormous machine is one failure domain, and resizing usually means a restart or failover. ## Practical order of operations Experienced teams exhaust the cheap options in roughly this order before touching topology: fix the top queries and indexes; fix the schema and access patterns; cache what is read repeatedly; scale up; offload reads to replicas; partition large tables inside the single server for maintenance and pruning; and only then shard or split by service. Sharding first is a classic over-engineering failure — it multiplies operational surface for a load that one well-tuned machine could carry. ## The honest summary in an interview Vertical scaling buys you time with almost no engineering cost and stops at one machine's write capacity. Horizontal scaling with replicas buys read capacity and availability but never write capacity. Horizontal scaling with shards buys write capacity and pays for it in query power and operational complexity. Knowing *which* of those three you need starts with knowing whether your problem is reads, writes, or data size.

  • Why does adding read replicas not help a write-heavy workload?
    Every write still executes on the single primary and is then shipped to each replica to be replayed, so replicas are consumers of the write stream rather than additional producers. Aggregate write capacity is bounded by the primary's CPU, log flush, and lock throughput. Each extra replica actually adds a small amount of shipping and acknowledgement work to the primary.
  • What would you try before sharding?
    Fix the expensive queries and missing indexes, correct schema and access patterns, add caching for repeatedly-read data, scale the machine up, and move read-only and reporting traffic to replicas. Table partitioning inside one server also helps with huge tables by making maintenance and retention cheap and enabling pruning. Sharding is the last step because it is the only one that permanently changes what the application can express in a single query.
  • Is scaling up ever the wrong answer even when a bigger machine exists?
    Yes — when the bottleneck is not a resource the hardware provides. A workload serialized on one hot row lock, one single-threaded operation, or fsync latency will not improve with more cores. Concentrating everything on one enormous machine also enlarges the failure domain and can make resizing a downtime event.

Vertical scaling is hiring a faster chef for the one kitchen. Replicas are extra serving counters handing out copies of the same dishes — more customers served, same cooking capacity. Shards are opening separate kitchens, each cooking only part of the menu — more cooking, but now nobody can plate a dish that spans two kitchens.

saying these in an interview costs you the question

  • Claiming read replicas increase write throughput
  • Treating a database like a stateless service — "just add instances behind a load balancer"
  • Jumping straight to sharding for a workload one tuned server could carry
  • Assuming vertical scaling is unlimited in the cloud, ignoring the max instance size and the cost curve
  • Saying scale-out is always cheaper, ignoring the operational cost of cross-shard queries and rebalancing

context