skip to content

Write throughput is saturating your primary database while the read replicas sit idle. Walk through how you would decide among batching, queueing, bigger hardware, and splitting the data across multiple primaries.

level: principalimportance: should knowfreq 42%

answer

  1. diagnose first: fsync / log / indexes / locks / connections
  2. every extra index multiplies write cost
  3. batch = fewer round trips + fewer flushes
  4. queue only what can be applied later, idempotently
  5. split across primaries last — permanent cross-node cost

basics

~20 s

Measure what is actually saturated first — commit/fsync, log volume, index maintenance, lock contention, or connections. Then take the cheap levers in order: batch writes, drop redundant indexes, shorten transactions, move non-critical writes to a durable queue, add IOPS. Splitting across primaries is last because it costs cross-shard queries and transactions forever.

solid answer

~1 min

**Step 1 — diagnose.** "Writes are slow" has several distinct causes: commit/fsync latency, log (WAL/redo) volume and checkpoint storms, index maintenance amplification, lock or latch contention on hot rows and index pages, connection oversubscription, or replication overhead. Each has a different fix; guessing wastes a quarter. **Step 2 — reduce the work per write.** - Batch: multi-row inserts and `COPY`-style bulk loads cut round trips and log flushes dramatically. - Drop redundant indexes — every extra index is another structure to maintain on every write. - Keep transactions short and never hold locks across network calls. - Relax durability *selectively* (group commit, deferred flush) for data you can afford to lose on crash — and only for that data. **Step 3 — move work off the synchronous path.** Non-critical writes (audit, analytics, denormalised counters, search indexing) go through a durable queue and are applied in batches. This changes the consistency contract, so it needs idempotency, ordering rules, and backpressure. **Step 4 — buy capacity.** Faster storage and more IOPS is often the cheapest year of runway; partitioning can shrink index depth and turn deletes into partition drops. **Step 5 — split across primaries** only when a single node genuinely cannot hold the write rate, because cross-node queries, transactions, and rebalancing become permanent costs.

go deeper

for a junior

Know that replicas do not help writes and that batching many small writes into fewer statements is the first easy win.

for a middle

Name the main write-path costs — commit flushes, log volume, index maintenance — and the corresponding cheap fixes.

for a senior

Diagnose from real signals, quantify index amplification, and design an asynchronous write path with idempotency and backpressure.

for a principal

Present an ordered decision framework with explicit costs at each step, classify data by durability and immediacy requirements, and state the gate that must be met before accepting the permanent complexity of multiple primaries.

## Why the ordering matters Each step down this list is roughly an order of magnitude more expensive in engineering time and permanent operational complexity than the one above it. Teams that jump to splitting the data first pay that cost forever, often to solve a problem that three redundant indexes were causing. ## Step 1: find the actual bottleneck Candidate bottlenecks, and how they present: - **Commit / fsync latency.** Throughput is a multiple of 1/commit-latency per session; storage flush time dominates. Symptom: high wait time on log flush, low CPU. Fix: faster storage, group commit, batching several logical writes into one transaction. - **Log volume.** Every write generates log records; full-page writes after a checkpoint multiply it. Symptom: log throughput near device limits, checkpoint spikes. Fix: fewer/narrower updates, fewer indexes, checkpoint tuning, better storage. - **Index maintenance amplification.** A table with 8 indexes turns one row insert into 9 structure modifications, most of them random writes. Symptom: insert cost far above row size, high buffer churn. Fix: delete unused indexes — the single highest-yield write optimisation in practice. - **Lock/latch contention.** Many sessions serialising on the same row or the same index leaf. Symptom: throughput flat or falling as concurrency rises, lock waits climbing. Fix: contention-specific (see the hot-row topic) — batching, sharded counters, append-only ledgers. - **Connection oversubscription.** Thousands of connections thrash the scheduler and memory. Symptom: high CPU with low useful throughput. Fix: a pooler with a sane pool size, not more hardware. - **Replication overhead / sync commits.** Every commit waits on a network round trip. Fix: reconsider synchronous replica placement and count. Without this step, every later decision is a guess. ## Step 2: reduce work per write - **Batching.** 1,000 single-row inserts, each its own transaction, pay 1,000 round trips and 1,000 commit flushes. One statement inserting 1,000 rows pays one of each. Order-of-magnitude gains are routine. Cost: latency for the batching window, larger failure blast radius, and bigger chunks for replicas to replay (which spikes lag), so size batches deliberately — typically hundreds to low thousands of rows. - **Index diet.** Audit index usage statistics and delete what is unused or redundant (a prefix of another index). Also consider whether some indexes only exist for a report that belongs on a replica. - **Narrower, fewer updates.** Wide updates rewrite more; in MVCC engines an update that touches an indexed column forces index maintenance that an update to non-indexed columns can avoid. Splitting a hot mutable counter out of a wide row is a classic win. - **Transaction hygiene.** Short transactions release locks sooner and shrink replay chunks. Never call an HTTP API inside a transaction. - **Durability tuning, selectively.** Group commit amortises fsync across concurrent commits at no correctness cost. Relaxing synchronous commit trades a window of possible loss for throughput — appropriate for telemetry, unacceptable for payments. Make it a per-workload decision, never a global switch flipped in a hurry. ## Step 3: change the shape of the write - **Queue non-critical writes.** Put audit events, activity feeds, denormalised aggregates and search-index updates onto a durable log/queue; a consumer applies them in batches. The database sees far fewer, larger transactions. Requirements: the write must be safe to apply later (the user need not see it immediately), consumers must be idempotent (retries are certain), ordering rules must be explicit, and the queue must have backpressure and a monitored lag — otherwise you have moved the outage rather than removed it. - **Append-only instead of update-in-place.** Inserting events and rolling them up periodically avoids contention and turns random updates into sequential appends. Costs read-side complexity. - **Staging + merge.** Bulk-load into an unlogged or temporary table, then merge in one pass. ## Step 4: buy capacity Storage IOPS and latency, more memory to keep the working set and index upper levels cached, more CPU for commit work. Modern single nodes handle very large write rates; buying a year of runway is often correct and lets you fix the real problem deliberately rather than in an incident. Partitioning also belongs here: smaller per-partition indexes are cheaper to maintain, and retention becomes a partition drop instead of a mass delete that generates enormous log volume. ## Step 5: split the data across primaries When one node's write ceiling truly is the limit, splitting data across independent primaries multiplies write capacity — and permanently imposes cross-node queries, distributed transactions or their avoidance, rebalancing, and per-node operations. It is the right answer at sufficient scale and the wrong answer as a first move. The gate I would apply: the write path has been measured and optimised, a bigger box has been costed and rejected, growth projections show the ceiling arriving within the planning horizon, and there is a natural key along which the data separates cleanly. ## The judgement being tested This question is not looking for a list of techniques; it is looking for **ordering discipline** — measure, then reduce work, then change the consistency contract deliberately, then spend money, then accept permanent architectural complexity — and for candour about what each step costs. The strongest answers also state explicitly which data may lose durability or immediacy and which may not, because that classification is what makes queueing and durability relaxation safe.

  • Batching improves throughput. What does it cost you?
    Latency, because writes wait for the batch window; a larger failure blast radius, since one bad row can fail a batch and retries must be idempotent; and bigger replay chunks, which spike replica lag and lengthen lock hold times. Batch sizes therefore have a sweet spot — usually hundreds to low thousands of rows — chosen against both the throughput curve and the lag and latency budgets.
  • When is relaxing durability an acceptable way to get write throughput?
    Only for data classes whose loss over the flush window is genuinely tolerable — telemetry, page-view counters, caches of derived data — and only as a per-workload setting, never a global one. Group commit is different and should be on by default: it amortises fsync across concurrent commits without weakening any guarantee. For financial or audit data, relaxing durability trades a real correctness property for throughput and is not a legitimate scaling technique.
  • How do you decide the moment has come to split writes across multiple primaries?
    When the write path has been measured and optimised, a larger instance and faster storage have been costed and shown insufficient within the planning horizon, and there is a key along which the data partitions cleanly with most transactions staying inside one partition. I would also want the operational maturity to run several primaries — backups, failover, schema change, and rebalancing — because those costs are permanent and arrive on day one.

saying these in an interview costs you the question

  • Proposing to split data across primaries before measuring what is saturated.
  • Suggesting more read replicas as a fix for a write bottleneck.
  • Turning off synchronous commit globally to get throughput.
  • Ignoring index count as a driver of write amplification.
  • Treating a queue as free — forgetting idempotency, ordering, backpressure and lag monitoring.

context