skip to content

A single counter row — say the like count on one viral post — is updated thousands of times per second and overall throughput collapses. What is happening inside the engine, and how would you redesign it?

level: seniorimportance: should knowfreq 44%

answer

  1. one row = serial section; rate ≈ 1/hold time
  2. hold spans update → commit (fsync + round trips)
  3. waiters eat connections, latency blows up
  4. shard the counter / append-only ledger / Redis INCR
  5. touch the hot row last; no network calls under lock

basics

~20 s

Updates to one row serialise on its exclusive row lock, so throughput is capped at roughly one update per lock-hold time — including commit and network round trips — while everyone else queues, and MVCC version chains and index churn make it worse. Fix by removing the single point: sharded counters, an append-only ledger with rollups, or an in-memory counter flushed periodically.

solid answer

~1 min

Writers to the same row must take an exclusive lock on it, so they execute strictly one at a time. Throughput is bounded by 1 / (lock hold time), and the lock is held from the update until commit — which includes the log flush and any application round trips still inside the transaction. At a few milliseconds per hold you are capped in the low hundreds per second regardless of how many cores the database has. Worse, the queue of waiters consumes connections and memory, latency fans out across the whole system, and in MVCC engines every update creates a new row version, driving update chains, index bloat, and vacuum/purge pressure. Redesigns, in order of preference: 1. **Shard the counter**: N rows per post, each writer picks one at random; readers `SUM`. Contention drops ~N-fold. 2. **Append-only ledger**: insert one row per like — no contention at all — and maintain the total with a periodic rollup or a materialised aggregate. 3. **Aggregate outside the database**: `INCR` in Redis or in-process, flush a delta every second. 4. **Coalesce in the application**: one writer batches N increments into a single `+N` update. Also reduce lock hold time: update the hot row last, keep transactions short, never make external calls while holding it.

code

sql · 13 lines
sql
CREATE TABLE post_like_counts (
  post_id bigint  NOT NULL,
  bucket  smallint NOT NULL,
  n       bigint  NOT NULL DEFAULT 0,
  PRIMARY KEY (post_id, bucket)
);

-- writer picks a random bucket 0..63
UPDATE post_like_counts SET n = n + 1
 WHERE post_id = 1 AND bucket = 17;

-- reader
SELECT SUM(n) FROM post_like_counts WHERE post_id = 1;

go deeper

for a junior

Explain that concurrent updates to one row must wait for each other's row lock, so they run one at a time.

for a middle

Quantify it — throughput is roughly one update per lock hold time, and the hold lasts until commit — and propose sharded counters.

for a senior

Add MVCC version-chain and index-churn effects, the connection-pool blast radius, and compare sharded counters, append-only ledgers, and external counters with their tradeoffs.

for a principal

Separate counters that may drift from counters that must be exact, and design the write path so the exact ones stay exact by construction (immutable ledger, reservation rows) rather than by holding one lock.

## What the engine is doing When two transactions update the same row, the second must wait for the first to commit or roll back — that is the row lock, and it is required for correctness. So updates to a single row are **serialised**. The maximum rate is: ``` max updates/sec ≈ 1 / (time the lock is held) ``` The lock is acquired at the `UPDATE` and released at commit. That interval includes the log flush (an fsync, often 0.1–2 ms) plus, crucially, any time the application spends between the update and the commit — extra statements, an HTTP call, a slow client. Even at 1 ms of hold time, the ceiling is ~1,000/s for that row; at 5 ms it is 200/s. Nothing about the hardware changes that: it is a serial section. Beyond the raw ceiling, three secondary effects turn a bottleneck into an outage: 1. **Queue amplification.** Waiters hold connections. When each connection blocks for tens of milliseconds, the pool drains, unrelated queries can't get a connection, and the incident spreads far beyond the counter. Lock waits also raise deadlock probability when transactions touch several hot rows in different orders. 2. **MVCC costs.** In PostgreSQL each update writes a new row version; if the table's page has no room, or an indexed column changed, the update is non-HOT and every index gets a new entry. Thousands of versions per second on one row means bloat, index churn, and vacuum working hard on a row whose logical content is one integer. In InnoDB and Oracle the equivalent cost is a long undo/rollback-segment chain that readers must walk to reconstruct their snapshot, so *reads* of the hot row get slower too. 3. **Latency variance.** With serialised access, waiting time grows non-linearly as arrival rate approaches the service rate — the classic queueing curve. Small traffic increases produce large latency jumps, which is why this fails suddenly rather than gradually. ## The redesigns ### Sharded (striped) counters Store N rows per logical counter — `(post_id, bucket)` for bucket 0..N-1. Each writer updates a randomly chosen bucket; readers `SELECT SUM(count) ... GROUP BY post_id`. Contention falls by roughly N. With N = 64 a 200/s ceiling becomes ~12,000/s. - Pros: keeps the counter in the database and transactional with the like itself; simple. - Cons: reads cost an aggregate over N rows (cache it, or maintain a rolled-up total); N is fixed at design time; the exact total requires reading all buckets. ### Append-only ledger plus rollup Insert a row per event (`likes(post_id, user_id, created_at)`). Inserts of distinct rows do not contend on a row lock at all. The total comes from a periodic rollup job, a materialised view, or an incrementally maintained summary table updated in batches. - Pros: zero contention on the write path, and you keep the underlying facts (who liked, when), which usually has product value anyway — and it makes de-duplication (one like per user) a unique constraint rather than a race. - Cons: storage grows with events; exact real-time totals need a cached/rolled-up value; the rollup job is another moving part. - Watch out for a different hot spot: many concurrent inserts with a monotonically increasing key all target the rightmost index leaf page, which becomes a latch contention point. Hash or reverse-key indexing, or an unordered key like a random UUIDv4 with the tradeoffs that brings, addresses it. ### Aggregate outside the database `INCR post:1:likes` in Redis is a single-node in-memory operation at hundreds of thousands per second; a background job flushes deltas to the database every second or so. - Pros: enormous headroom, and the read is instant. - Cons: the counter is no longer transactional with the underlying write; a crash between increment and flush loses or double-counts unless flushes are idempotent (flush a delta with a token, or read-and-reset atomically). Fine for engagement counters, wrong for money. ### Coalesce in the application Accumulate increments in memory per process and issue one `UPDATE ... SET n = n + k` per second. Same serialisation, but 1/k as many acquisitions. Cheap to build; loses in-flight counts on crash. ### Reduce the hold time (do this regardless) - Touch the hot row **last** in the transaction, immediately before commit. - Never make a network call, wait on user input, or run a slow query while holding the lock. - Keep the transaction to the minimum statements. - Avoid `SELECT ... FOR UPDATE` followed by application arithmetic; use an atomic `SET n = n + 1` so the lock is held for one statement rather than a round trip. ## When strict correctness matters If the counter is a balance or an inventory count, you cannot silently trade accuracy for throughput. The legitimate patterns there are the append-only ledger (the balance is the sum of immutable entries, and a reservation row provides the constraint) or explicit partitioning of the resource into buckets that are individually claimed. Both preserve exactness while removing the single serial point. ## What a strong answer includes The engine explanation (row lock serialisation, hold time as the throughput bound, MVCC version chains), the systemic amplification (connection pool exhaustion, latency blowup), at least two redesigns with their tradeoffs, and the distinction between counters that may drift and counters that may not.

  • How many buckets should a sharded counter have, and what does a larger number cost?
    Enough that the per-bucket update rate falls well under the serial ceiling: divide the peak rate for the hottest key by a safe per-row rate (a few hundred per second) and round up, so a 10,000/s key wants something like 64 buckets. Larger N costs read work — every total is an aggregate over N rows — plus storage per key, and it is awkward to change later, so people usually cache or roll up the total and pick N generously once.
  • Why does an incident like this take out queries that have nothing to do with the counter?
    Because each waiting transaction occupies a connection and a worker for the whole wait. As the queue on the hot row grows, the connection pool drains and unrelated requests cannot get a connection, so the failure presents as a site-wide outage. Bounded pools, short lock-wait timeouts, and per-endpoint concurrency limits keep the blast radius local.
  • When can you not simply move the counter into Redis?
    When the count must be exact and transactional with the underlying data — balances, inventory, quota enforcement — because an external counter can drift from the database on crash, retry, or partial flush. In those cases the correct patterns are an append-only ledger whose sum is the authoritative value, or explicit reservation rows that make each claim a distinct row, both of which stay exact while removing the single serial point.

One turnstile for a stadium. It does not matter how many staff you hire — everyone must pass through single file, and the queue grows until the gate is what people talk about. You fix it by opening more turnstiles, not by making the queue more orderly.

saying these in an interview costs you the question

  • Proposing a bigger instance or more replicas — neither removes a serial section.
  • Using SELECT ... FOR UPDATE plus application arithmetic, which holds the lock across a round trip.
  • Blaming deadlocks when the actual mechanism is lock-wait serialisation.
  • Forgetting the MVCC side effects: version chains, index churn, vacuum/undo pressure.
  • Moving an exact counter (balance, inventory) into an external cache and calling it solved.

context