skip to content

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%

answer

  1. runway = headroom ÷ growth vs remedy lead time
  2. writes cap, reads don't
  3. restore time vs RTO forces the split
  4. cost curve steepens at top of family
  5. least drastic split: replicas → functional → shard

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.

solid answer

~60 s

Treat it as a **runway** decision, not a taste decision. The trigger is that the binding constraint can no longer be bought: - **Write ceiling** — a well-tuned single primary is saturated on commit path or contention, and reads are already offloaded. More machines cannot add writers without more primaries. - **Data volume** — backup, restore, index rebuild, and version-cleanup windows exceed the recovery objectives; a table's maintenance no longer fits the night. (Partitioning inside one server often fixes this first.) - **Cost curve** — you are near the top of the instance family, where each doubling costs disproportionately and buys less. - **Blast radius / availability** — one machine carrying the whole business exceeds the tolerable failure domain. Evidence to bring: measured headroom on the binding resource, the growth rate and hence the date the ceiling is hit, the results of cheaper interventions already applied, the cost of the next hardware step versus the engineering cost of scale-out, and a migration plan with a reversible first step. Prefer the least drastic split that removes the constraint — offload reads, then split by function, then partition, and shard last.

go deeper

for a junior

Recognise the sequence — tune, cache, scale up, replicas, then split — and that adding machines does not add writers.

for a middle

Name the specific constraints (write ceiling, maintenance windows, cost curve) and the cheaper interventions that must come first.

for a senior

Bring measurement: current headroom, growth rate, the date of exhaustion, and what each prior intervention bought; recommend the least drastic split that removes the constraint.

for a principal

Own the whole decision: runway and lead time, cost versus engineering economics, failure domain and organisational change risk, the seam the split follows, what capability the application permanently loses, a reversible first step, and readiness to operate N primaries.

## Reframe the question as runway The useful version of "should we scale out?" is: *on the current trajectory, when does the binding constraint become unmeetable, and what is the lead time of each available remedy?* Scale-out for a relational system is a multi-quarter engineering programme with permanent consequences for what the application can express, so it must be started before the wall, not at it — but starting it for a wall three years away wastes the team. So you need three numbers: current headroom on the binding resource, the growth rate, and the lead time of the remedy. That gives a date, and dates make the decision defensible. ## The four constraints that hardware cannot buy away **1. The single-writer ceiling.** Reads can be offloaded to replicas indefinitely; writes cannot. Once a well-tuned primary — sane indexes, batched writes, short transactions, no hot-row serialization, low-latency log storage — is saturated, no purchase adds writers. Only more primaries do, which means sharding or splitting the write workload out. Evidence: sustained commit-path saturation or contention at peak, with the tuning list demonstrably exhausted. **2. Data volume versus operational windows.** Very large single databases fail on *operations* before they fail on queries. Full backup and, more importantly, **restore** time; index rebuilds; schema migrations that rewrite tables; version cleanup or purge that cannot keep up; major-version upgrades. When restore time exceeds the recovery time objective the business has agreed to, size alone forces a split — though in-server table partitioning frequently rescues this case first, because it makes retention and maintenance per-partition rather than per-table. **3. The cost curve.** Instance pricing is roughly linear through the middle of a family and steepens near the top; the last doubling may cost several times the previous step while delivering a smaller relative gain because the workload does not scale linearly with cores. When the next hardware step costs more than the amortised engineering cost of the split, the economics have flipped. **4. Failure domain and change risk.** One machine holding the entire business is a single blast radius: one bad migration, one corrupted page, one resize gone wrong takes everything down. It also concentrates change risk — every team's schema change lands on the same server, and one team's runaway query is everyone's incident. At some organisational size this alone justifies splitting, independent of capacity. ## Everything you must have tried first A credible scale-out proposal is stronger for listing what it is *not*: - **Query and index work** — the top statements by total time, fixed. This is routinely a 2–10x win. - **Schema and access-pattern work** — removing write amplification, narrowing rows, eliminating hot-row counters, batching. - **Caching** — repeatedly-read, rarely-changing data served from a cache; the cheapest read scaling there is. - **Scale-up** — including a correct diagnosis of the binding resource, so the purchase actually targets it. - **Read offload** — replicas for reporting and read-heavy paths, with an explicit answer for replica staleness. - **Connection discipline** — a global pool budget and a pooler, since connection pressure masquerades as capacity exhaustion. - **Table partitioning inside one server** — turns retention and maintenance into metadata operations and enables pruning, often deferring the split by years. ## Choose the least drastic split that removes the constraint Scale-out is not one thing. In increasing order of cost and irreversibility: 1. **Read offload to replicas** — no data model change; fixes read capacity only. 2. **Functional / vertical split** — move a bounded subsystem (sessions, audit, events, search, metrics) to its own store. Often removes most of the write load with far less complexity than sharding, because the split follows an existing seam in the domain. 3. **Move the wrong-shaped workload out** — time-series, logs, blobs, and full-text often do not belong in the transactional relational store at all. 4. **Shard** — split one entity family across primaries by key. Highest capability, highest cost: cross-shard queries, global uniqueness, distributed transactions, rebalancing, and every operational task multiplied by N. The common failure is skipping to (4) when (2) or (3) would have removed the constraint for a fraction of the effort. ## What the decision document contains - The binding constraint, measured, with the wait or resource evidence. - Growth rate and the projected date of exhaustion, with uncertainty. - Interventions already applied and the headroom each bought. - Cost of the next hardware step versus engineering cost and lead time of the split. - The chosen split, the seam it follows, and explicitly what the application loses (which joins, which transactions, which constraints). - A **reversible first step** — for example, moving one subsystem's writes behind an interface before physically separating it, or introducing a routing layer while all keys still resolve to one shard. Being able to stop halfway without a rollback catastrophe is what makes the programme safe to start. - Operational readiness: N primaries means N of everything — backups, monitoring, failover drills, schema-change pipelines. ## The judgment in one line Stop buying hardware when the constraint stops being a resource and becomes a structure: one writer, one maintenance window, one failure domain, one nonlinear price. Until then, the boring interventions win, and the scale-out capital is better spent later, on a seam you understand.

  • What is the strongest single signal that scale-out is now unavoidable?
    A well-tuned single primary saturated on the write path, with reads already offloaded and the tuning list exhausted. Reads can always be moved to replicas and caches, but no purchase adds writers to a single-primary system. When commit-path or contention saturation persists at peak after indexes, batching, and transaction scope have been fixed, only more primaries change the outcome.
  • Why is a functional split often preferable to sharding?
    It follows a seam the domain already has, so the queries and transactions that cross it were usually few or already indirect. Moving sessions, audit, events, or metrics to their own store can remove the majority of write load while keeping the core transactional model intact and joinable. Sharding, by contrast, cuts through one entity family and permanently costs you cross-shard joins, global uniqueness, and single-node transactions.
  • How do operational windows force a split independently of query performance?
    Backup, restore, index rebuild, migration, and version-cleanup times scale with data volume, not with query rate. Once restore time exceeds the agreed recovery objective, or a migration can no longer complete in a maintenance window, the database is too large regardless of how fast queries are. Table partitioning inside one server often resolves this first by making retention and maintenance per-partition.
  • What makes a scale-out programme safe to start?
    A reversible first step and a clear seam. Introducing a routing layer while every key still resolves to one destination, or moving a subsystem behind an interface before physically separating it, lets the team stop or pause without a rollback catastrophe. Alongside that you need operational readiness for N primaries: backups, monitoring, failover drills, and a schema-change pipeline that works across all of them.

saying these in an interview costs you the question

  • Proposing sharding without evidence that tuning, caching, replicas, and partitioning were exhausted
  • Justifying the decision on architecture taste rather than measured headroom and growth
  • Ignoring that restore time and maintenance windows can force a split before performance does
  • Assuming scale-out reduces cost — it usually raises total cost while removing a ceiling
  • Starting an irreversible migration with no reversible first step

context