skip to content

Vertical vs horizontal scaling of an RDBMS

Scaling up a single server versus scaling out to many, seen from the engine's side: what a bigger box buys you, where it stops (connections, write throughput, failover blast radius), and what scaling out costs in complexity. A standard warm-up before deeper sharding questions.

part ofRelational database conceptsoverview, primer and where to startread it →
on this pageshow

questions

5

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

open as a page

A team adds five more database servers to a relational cluster and is surprised that overall write throughput barely changes. Why does adding machines fail to multiply write capacity in a typical single-primary relational deployment, and what actually raises that ceiling?

level: middleimportance: must knowfreq 60%

basics

~20 s

In a single-primary cluster every write is executed once on the primary and then replayed on the others, so the extra machines are copies, not additional writers. The ceiling is one machine's commit path. You raise it by making that path cheaper, or by having more primaries — i.e. sharding.

open as a page

A relational database server shows plenty of idle CPU and free memory, yet query latency degrades badly after the application tier grows from 20 instances to 200, each holding its own connection pool. Explain what limits the number of usable database connections and how connection pooling changes the picture.

level: seniorimportance: should knowfreq 50%

basics

~20 s

Each connection costs memory and a scheduling slot, and concurrency beyond the number of cores and disks only adds context switching, lock contention, and cache thrash. Throughput plateaus then falls. Fix it with a small pool per instance, or a shared external pooler that multiplexes many clients onto few server connections.

open as a page

Before ordering a larger database machine, how do you determine which resource — CPU, memory, or storage IOPS and latency — is actually limiting the server, and which of those a hardware upgrade will genuinely fix?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Measure where time goes, not what looks busy. Sustained high CPU with a hot working set in memory means CPU-bound; heavy read IO with a low cache hit ratio means memory-bound; high commit latency or IO wait means storage-bound. Lock waits and single-threaded hotspots are not fixed by bigger hardware.

open as a page

How do you decide that a relational system has reached the end of useful vertical scaling and must move to a scale-out architecture, and what evidence would you bring to that decision?

level: principalimportance: should knowfreq 38%

basics

~20 s

Scale out when the constraint is structural, not purchasable: the single primary's write ceiling, data volume that breaks maintenance and recovery windows, or a cost curve that has gone nonlinear — and only after tuning, caching, replicas, and partitioning have been exhausted. Bring measured headroom, growth rate, and the date the ceiling arrives.

open as a page