skip to content

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%

answer

  1. one node orders all writes
  2. replicas replay, they don't originate
  3. fsync + index maintenance + lock contention = the cap
  4. batch and group commit collapse flushes
  5. more write capacity ⇒ more primaries ⇒ sharding

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.

solid answer

~60 s

A classic relational cluster has exactly **one node that accepts writes**. The other nodes replay its change log. So extra nodes add read capacity and failover, but every insert, update and delete still passes through one primary's CPU, buffer pool, lock manager, and log flush. Aggregate write throughput is therefore capped by that one machine — the **single-writer ceiling** — and each additional replica actually adds a little work (shipping and acknowledging the log), especially if it is synchronous. Raising the ceiling has two families of answer: 1. **Make each write cheaper** on the one primary: fewer and narrower indexes to maintain, batching and group commit so one fsync covers many transactions, shorter transactions to reduce lock hold time, faster storage for log flush latency, and removing write amplification (triggers, redundant updates, hot-row counters). 2. **Have more primaries**: shard the data so each machine owns a disjoint slice and accepts its own writes, or move a write-heavy workload (event streams, logs, metrics) out of the relational system entirely. Multi-primary/active-active relational setups exist but move the problem to conflict resolution rather than removing it.

go deeper

for a junior

Say clearly that only the primary accepts writes and the others copy it, so extra machines add reads and failover, not write capacity.

for a middle

Name what saturates on the primary — log flush, index maintenance, lock contention — and give the cheap fixes: fewer indexes, batching, shorter transactions.

for a senior

Diagnose which resource is actually the ceiling before recommending anything, and distinguish a workload that needs tuning from one that genuinely needs more primaries or a functional split.

for a principal

Discuss when the right answer is to move a write-heavy subsystem out of the relational system entirely, and the durability/latency/conflict tradeoffs of synchronous replication versus multi-primary.

## The shape of a normal relational cluster Almost every mainstream relational deployment is **single-primary**: one node accepts writes, and one or more replicas continuously apply the primary's change stream (a write-ahead log, a binary log, or a statement/row replication feed). This is not an accident. Relational semantics — serializable-ish isolation, unique constraints, foreign keys, ordered commits — are far easier to guarantee when exactly one node decides the order of changes. The consequence: **replicas are copies, not additional writers**. A write is not distributed across the cluster; it is performed once on the primary and then re-performed on each replica. Five extra machines means the same write is now done six times in total, not that six writes happen in parallel. ## Where the ceiling physically lives On the single primary, a write passes through several serializing resources, and the ceiling is whichever saturates first: - **Log flush (fsync) latency.** A durable commit must get the log record to stable storage. This is a latency-bound, largely serial step. It is why commit rate on spinning disks was historically a few hundred per second, and why NVMe and battery-backed caches changed the game so much. - **Index maintenance.** Every index on a table is another structure to update per row change. A table with eight indexes costs several times more per insert than one with two. - **Lock and latch contention.** Rows, pages, and internal structures that many transactions touch (a counter, a sequence, the tail of a hot index, a status column everyone updates) serialize work no matter how many cores exist. - **MVCC bookkeeping.** Engines that keep old row versions must also clean them up (vacuum, purge, undo). At high write rates this maintenance itself becomes a load, and if it falls behind, read performance degrades too. - **Replication back-pressure.** With synchronous replicas, each commit waits for at least one remote acknowledgement, so adding a synchronous replica *lowers* peak write throughput while raising durability. ## Why this is different from the application tier A stateless service scales out by cloning: N identical instances, any of which can serve any request, because none of them owns state. A database owns state, so an added machine must either hold **the same** data (replica — redundancy, no extra write capacity) or **different** data (shard — real extra write capacity, different semantics). There is no third option. Understanding this is the whole point of the question. ## Raising the ceiling without changing topology Before any architecture change, you can often buy a 2–10x margin on the same box: - **Drop unused and redundant indexes.** This is usually the single highest-leverage change for insert/update-heavy tables. - **Batch.** Committing 500 rows in one transaction instead of 500 transactions collapses 500 log flushes into one. Group commit does this automatically for concurrent sessions. - **Shorten transactions.** Never hold a transaction open across a network call or user think-time; the lock hold time, not the work, is what serializes others. - **Remove hot-row patterns.** A single counter row updated by every request is a serialization point; sharded counters or an append-and-aggregate pattern removes it. - **Faster storage and more memory.** Lower fsync latency raises commit rate directly; a working set that fits in memory turns random read IO into CPU. - **Reduce write amplification.** Triggers, over-broad updates that rewrite unchanged columns, and audit tables written synchronously all multiply the cost of each logical write. ## When the ceiling is real, the answers are structural If a well-tuned primary on good hardware still saturates, the options are: 1. **Shard** — multiple primaries, each owning a slice by key. Aggregate write capacity scales with shard count; cross-shard queries and transactions become the new problem. 2. **Functional split** — move a subsystem (sessions, audit, events, metrics) to its own database or to a store designed for that write pattern. Often far cheaper than sharding the whole model. 3. **Change the write pattern** — buffer through a queue and write in batches, precompute instead of updating on every event, or accept eventual materialization. 4. **Multi-primary / active-active** — possible, but it converts a capacity problem into a conflict-resolution problem: two nodes may accept conflicting writes to the same key, and someone must define the reconciliation rule. It rarely gives linear write scaling for a contended workload. ## The one-line takeaway In a single-primary system, writes do not parallelize across machines — they parallelize only across cores, disks, and locks *inside* one machine. More servers give you reads and survival, not write capacity; more primaries give you write capacity, and cost you the ability to join and transact freely across the whole dataset.

  • Does a synchronous replica help or hurt write throughput?
    It hurts throughput and latency while improving durability and failover safety. Each commit must wait for the standby to acknowledge receipt (or apply), adding at least one network round trip to the commit path. Teams usually keep one synchronous standby for safety and make the rest asynchronous.
  • You cannot shard yet and writes are saturating the primary. Name concrete changes that buy headroom.
    Remove unused and redundant indexes so each row change maintains fewer structures; batch small transactions so one log flush covers many rows; shorten transaction scope so locks are held briefly; and eliminate hot-row updates such as a global counter by sharding the counter or aggregating from an append-only table. Faster storage lowers fsync latency, which directly raises commit rate.
  • Does multi-primary (active-active) replication remove the single-writer ceiling?
    It removes the single-node bottleneck only for workloads that partition cleanly by key. Once two primaries accept writes to the same rows you must resolve conflicts — last-writer-wins, application merge, or distributed locking — and any of those either loses data or reintroduces coordination latency. It changes the problem rather than eliminating it.

saying these in an interview costs you the question

  • "Just add nodes and the cluster shares the write load" — replicas replay, they do not accept writes
  • Believing a load balancer in front of the database spreads writes
  • Ignoring that each index multiplies per-row write cost
  • Thinking a synchronous standby increases throughput
  • Proposing multi-primary as free write scaling without mentioning write conflicts

context