A service keeps a per-tenant counter in a single row that is updated thousands of times per second. Read traffic scales fine, but write latency spikes and throughput plateaus. Explain why multi-version concurrency control does not help this workload, and how you would weigh the options for fixing it.
answer
- Contention is writer-vs-writer, MVCC irrelevant
- Ceiling ≈ 1 / lock-hold time
- Shorten hold: update last, no network in transaction
- Shard N rows, read SUM
- Append events, roll up later
basics
~20 sMVCC removes read locks, not write locks. Every updater of that one row serializes behind an exclusive row lock held to commit, so throughput is capped near one update per lock-hold duration. Fix by shortening the hold, sharding the counter into many rows, or appending events and aggregating on read.
solid answer
~1 minMVCC's readers-don't-block-writers property does nothing here, because the contention is **writer versus writer**. One row has one current version; each updater takes an exclusive lock on it held to commit, so writes on that row form a single queue. The ceiling is roughly `1 / lock-hold time`: hold the lock 2 ms and you get about 500 updates/s no matter how many cores or connections you add. Adding connections past that point makes latency worse, not throughput. Levers, cheapest first: 1. **Shorten the hold.** Move the counter update to the last statement before COMMIT; get all network calls, validation and other writes out of the critical section. Often a 5-10x win with no schema change. 2. **Shard the counter.** N rows per tenant, writers pick one at random or by worker id, readers SUM. Trades exact single-row reads for N-way write parallelism. 3. **Append then aggregate.** Insert one row per event (no contention at all, since inserts do not fight over an existing row) and roll up periodically or in a materialized view. 4. **Move it out of the transaction** into an in-memory counter or stream, if the value tolerates approximation and delay. The decision hinges on how exact and how fresh the counter must be.
code
sql · 7 lines-- write path
UPDATE counter_shards
SET count = count + 1
WHERE tenant_id = :tenant AND shard_id = :shard; -- shard = hash(worker) % 32
-- read path
SELECT sum(count) FROM counter_shards WHERE tenant_id = :tenant;go deeper
Recognise that MVCC only removes read locks and that all writers of one row queue behind an exclusive lock, so this row is a serial bottleneck.
Quantify the ceiling as one update per lock-hold interval and name the first fix: shorten the transaction and update the counter last, with no external calls inside.
Diagnose from lock-wait metrics, then pick between shortening the hold, sharding, and append-plus-rollup, stating the read-side and operational costs of each.
Lead with the requirement — how exact and how fresh must the number be, and who consumes it — then choose an architecture, and be explicit that this is a design limit no tuning knob removes.
## Why the read-side win is irrelevant here MVCC eliminates read locks by keeping multiple versions. But a logical row has exactly one *current* version, and two transactions cannot both produce its successor without losing one of the updates. So writers take an exclusive row lock, held to end of transaction, and first-updater-wins. On a single hot row that produces a strict queue. The property everyone quotes — readers don't block writers — is simply not the axis under stress. Worse, hot-row updating interacts badly with the rest of the MVCC machinery. Each update creates a version, so a row updated thousands of times per second generates thousands of dead versions per second. Depending on the storage design, that means either a growing version chain the engine must walk, or a heavily churned page plus reclamation pressure and index maintenance. High-frequency single-row update workloads are exactly where the version-storage cost of MVCC shows up worst. ## Do the arithmetic first The ceiling is `1 / (time the lock is held)`. If the transaction touches the counter early and then does 4 ms of other work before commit, every other writer waits those 4 ms: roughly 250 updates/s. This is a *serial* resource, so Little's law applies — adding concurrency past the ceiling only grows the queue. A common symptom is a latency graph that is flat until a threshold and then hockey-sticks, with lock-wait time dominating. Measure the hold time before choosing a redesign; the fix is often much cheaper than people expect. ## Lever 1 — shorten the critical section The lock starts at the UPDATE and ends at COMMIT. So: - Do the counter update **last**, immediately before commit. - Get every external call — payment gateway, other service, message publish — **out of the transaction**. Holding a database row lock across a network call couples your write ceiling to someone else's p99. - Keep the transaction to what must be atomic; split reporting or auditing writes out. - Ensure a consistent update ordering across all code paths, since contention makes deadlocks far more likely. This usually buys an order of magnitude and costs nothing structurally. It should always be step one. ## Lever 2 — shard the counter Replace one row with N (say 16 or 64) per tenant. Each writer picks a shard by worker id, connection id, or random choice, and updates that row. Reads become `SELECT sum(count) FROM counter_shards WHERE tenant_id = ?`. Write concurrency multiplies by roughly N. Tradeoffs to state explicitly: reads get more expensive and touch N rows (fine if reads are rarer than writes, which is the premise here); decrement-with-a-floor logic gets harder because no single shard knows the total, so quota enforcement at an exact boundary needs care; and N must be chosen — too small does not relieve contention, too large makes reads and cache footprint worse. Sharding is the standard answer when the counter must stay transactional and reasonably fresh. ## Lever 3 — append then aggregate Insert one row per event instead of updating a shared row. Inserts to different rows contend for almost nothing (page allocation and index leaf pages, which spreads out), so write throughput scales close to the hardware. The aggregate becomes a periodic rollup job, a materialized view, or an on-demand SUM over a recent window plus a stored base. This buys the most write headroom and additionally gives you an audit trail. It costs storage, a retention/cleanup policy, and a rollup pipeline — real operational surface. It also changes read freshness semantics unless you sum the un-rolled-up tail at read time. ## Lever 4 — take it out of the database If the counter is telemetry-grade — rate limiting with soft edges, dashboards, usage estimates — the right answer may be that it does not belong in an ACID row at all. An in-memory counter with periodic flush, or a stream aggregated downstream, removes the contention entirely. This is only acceptable if losing a few increments on a crash and lagging by seconds is genuinely fine. Say that condition out loud rather than assuming it. ## The judgment call The question to answer before choosing is: **how exact and how fresh must this number be, and who reads it?** A billing counter that gates a hard quota, an at-most-once idempotency guard, and a usage dashboard have three different right answers. Also ask what else the transaction does — if the counter update is one line inside a transaction that also calls an external service, the redesign may be unnecessary once that call moves out. Finally, be explicit that this is a *design* limit, not a tuning knob. No isolation level, connection-pool size, or index changes the fact that a single row's updates serialize. Raising the pool size in response to hot-row contention makes the symptom worse and is a classic wrong move.
- Why does increasing the connection pool or worker count make hot-row contention worse?The row is a serial resource, so its throughput is fixed at roughly one update per lock-hold interval regardless of how many clients are waiting. Extra concurrency only lengthens the wait queue, raising latency and tail latency while throughput stays flat. It also increases deadlock probability and consumes connections that other work needs.
- What is the downside of sharding a counter into 64 rows?Reads must sum 64 rows instead of one, so read cost and cache footprint grow, and any logic needing an exact instantaneous total — enforcing a hard quota at the boundary, for example — becomes harder because no single row knows the total. Decrements with a floor are especially awkward since an individual shard can go negative. You are trading read simplicity and exactness for write parallelism, which is only the right trade when writes dominate.
- How does a high-frequency hot-row update interact with version cleanup in an MVCC engine?Every update creates a new version, so thousands of updates per second produce thousands of dead versions per second on one row or page. The engine must reclaim them continuously, and readers may have to walk longer version chains in the meantime. If any long-running transaction pins an old snapshot, that cleanup stalls and the hot row's access cost degrades sharply.
One row is a single-lane toll booth. MVCC widened the sightseeing lane (reads) to infinity, but every car that must pay still queues at the one booth; you either speed up the booth, open more booths, or collect tolls by mail.
saying these in an interview costs you the question
- Expecting MVCC to remove write-write contention because 'readers don't block writers'
- Proposing a larger connection pool or more workers as the fix
- Suggesting a lower isolation level, which does not affect write locks
- Recommending sharding without mentioning the cost to reads and exact totals
- Moving the counter out of the database without asking whether lost increments are acceptable