skip to content

You're designing an event-sourced system where a small number of aggregates — say, the top few 'hot' inventory SKUs during a flash sale — receive a disproportionate share of all writes, causing constant optimistic-concurrency retries on those specific streams. What are the real options for fixing this, and what does each one cost?

level: principalimportance: nice to knowfreq 30%

answer

  1. conflict rate scales with concurrent writers per stream
  2. sub-partition/shard the hot aggregate
  3. make the operation commutative to skip version checks
  4. decouple via intent event + async settlement
  5. single-writer handler removes races but caps throughput

basics

~20 s

When too many people try to update the same one record at once, retries pile up. Fixes usually mean either spreading that one record's updates across several smaller pieces, changing the kind of update so order stops mattering, or accepting some updates outside the strict record and reconciling later — each trades away some strict consistency or simplicity for more speed.

solid answer

~60 s

Since optimistic concurrency conflicts are proportional to concurrent writers on the same stream, the fix has to reduce how many writers race for the same expected-version slot, or remove the need for that check at all. Common options: (1) sub-partition the hot aggregate into multiple sub-streams (e.g., per-warehouse or per-shard stock counts for one SKU) and combine them for the total, trading a slightly more complex read for less contention; (2) redesign the operation as a commutative update (e.g., a decrement-only reservation counter) that doesn't need strict ordering to be correct, so concurrent writers no longer conflict at all; (3) move to an async reservation/queue-then-settle pattern, where the hot path appends a lightweight 'intent' event to a less contended stream and a background process later applies it to the aggregate, trading immediate consistency for throughput; (4) simply serialize writes to that one stream through a single in-memory handler instance holding it in memory, so there's no concurrent access to conflict in the first place, at the cost of that stream becoming a scalability ceiling by design.

go deeper

for a junior

Not expected to design a fix; may only recognize that 'too many people editing the same thing at once' causes conflicts.

for a middle

Should suggest at least one concrete lever, like splitting into smaller pieces or doing the update asynchronously, even if the trade-offs aren't fully articulated.

for a senior

Should describe two or more of the levers with concrete trade-offs and connect the fix back to why optimistic-concurrency conflict rate is proportional to concurrent writers per stream.

for a principal

Should compare all the major levers with their specific costs, pick an appropriate one for a stated business constraint (e.g., must show instant confirmation vs. can tolerate async), and flag the invariant-precision or reconciliation burden each introduces.

## Why a hot aggregate hurts The problem starts from a mechanical fact already established by optimistic concurrency control: every append to a stream must match the stream's current expected version, checked atomically, and a conflict forces the losing writer to reload and retry. **Conflict probability rises directly with the number of concurrent writers targeting the same stream in the same short window.** A hot aggregate — one SKU's inventory-count stream during a flash sale, one popular auction's bid stream in its last seconds, one celebrity's follower-count stream — can see hundreds or thousands of concurrent append attempts against a single stream, and once concurrent writers exceed what a single stream can serialize and retry through in real time, users start seeing failed checkouts, timeouts, or unacceptable latency, even though the aggregate is perfectly healthy in isolation. ## The first lever — re-partitioning the aggregate The first lever is re-partitioning the aggregate itself. Instead of one `sku-4471` stream holding 'current stock: 40', you split it into N sub-streams — say, one per warehouse or one per logical shard — each independently tracking a portion of the total (`sku-4471-shard-0`, `sku-4471-shard-1`, ...), and the 'current stock' the business cares about becomes a sum computed by a projection over all shards rather than a single stream's folded state. This directly reduces contention because writers are now spread across N streams instead of one, at the cost of real complexity: - a single 'reserve one unit' business operation may now need to pick a shard (round-robin, or 'first shard with stock'), - decisions about stock allocation across shards get fuzzier (a shard could show 0 while another has 5, requiring either rebalancing logic or accepting some false negatives), - and any invariant that spans the whole aggregate (e.g., 'never go negative in total') has to be re-derived from the sum rather than enforced by a single stream's own history. ## The second lever — a commutative operation The second lever is changing what you're modeling into something commutative — an operation where the order two concurrent writes are applied in doesn't change the result, so you no longer need strict per-write ordering to be correct, and therefore no longer need optimistic-concurrency serialization to protect it. A raw 'set stock to 39' write is not commutative (order matters enormously), but 'decrement stock by 1, reject if it would go below 0' can, with the right underlying structure (e.g., a database-level atomic decrement), be made safe under high concurrency without a version check per write, because each decrement is independently valid against the current count rather than against a specific expected prior version. The cost here is that this pattern doesn't transfer to every kind of update — many business operations are genuinely order-sensitive (you can't commute 'ship the order' and 'cancel the order') — so it only relieves contention on the subset of operations on the hot aggregate that are actually commutative, and modeling that correctly is real design work, not a drop-in fix. ## The third lever — decoupling the hot write path The third lever is decoupling the hot write path from the aggregate's own stream entirely: instead of every 'buy this SKU' request racing to append directly to `sku-4471`, the request appends a lightweight, always-independent 'reservation intent' event to its own low-contention stream (e.g., keyed by the request/session, not the SKU), and a single background consumer processes intents against the actual SKU stream serially, at a pace it controls, possibly batching several intents into one append. This removes contention entirely from the user-facing write path — nothing there does an expected-version check on the hot stream anymore — at the cost of moving from immediate consistency ('you know instantly if you got the item') to eventual/async consistency ('your reservation is pending, confirmed a moment later'), which is a genuine user-experience and correctness trade-off that has to be surfaced to the business, not hidden as an implementation detail. ## The fourth lever — removing concurrency at the source The fourth lever, mostly relevant at extreme scale or extreme correctness requirements, is to remove concurrency at the source: pin the hot aggregate to a single in-memory handler instance (a **single-writer-per-key architecture**) that processes all commands for that one stream strictly sequentially in memory, appending in order with no expected-version races because there's structurally only ever one writer. This trades away horizontal write scalability for that specific aggregate — it becomes, by design, bottlenecked on one process's throughput — in exchange for zero wasted retries and simpler reasoning, and only works if that one process can actually keep up with the load (which is exactly the assumption that breaks during a true flash-sale spike unless it's been load-tested at the target concurrency). ## What each lever costs None of these are free: | Lever | What it trades | |---|---| | (1) | trades read simplicity and cross-shard invariant strength | | (2) | trades applicability (only commutative operations benefit) and modeling effort | | (3) | trades immediate consistency and adds an async settlement path with its own failure modes (stuck intents, ordering of the settlement) | | (4) | trades horizontal scalability for that one aggregate | A concrete real-world example of levers (1) and (2) combined: high-throughput ticketing and flash-sale systems commonly model 'available inventory' as a set of small, independently-decrementable shard counters summed by a read model, specifically to survive the exact hot-aggregate write storm this question describes.

  • How would sharding a hot aggregate's stock count affect an invariant like 'stock can never go negative'?
    The invariant now has to be enforced per-shard (each shard independently refuses to go below zero) rather than against one true global total, which can produce a false 'out of stock' when one shard is empty while another still has units, unless you add rebalancing or an allocation strategy across shards. The trade is fewer conflicts for a slightly less precise, eventually-reconciled notion of total stock.
  • Why doesn't making an operation commutative eliminate the need for optimistic concurrency control entirely for that aggregate?
    Only the commutative subset of operations on that aggregate benefits — most aggregates still have some order-sensitive transitions (like cancel-after-ship) that genuinely need strict ordering, so optimistic concurrency control typically stays in place for those while the hot, high-volume commutative operations get carved out to bypass it.
  • What new failure mode does the async 'intent event, settle later' pattern introduce that direct synchronous appends didn't have?
    Stuck or lost intents: if the background settlement consumer crashes, falls behind, or an intent is malformed, a user's request can be accepted (their intent was recorded) but never actually applied to the aggregate, requiring monitoring, timeouts, and a reconciliation/dead-letter path that a synchronous append-and-fail-fast design didn't need.

Like one bank teller window with a huge line versus opening several windows that each handle part of the queue and reconcile totals at day's end — more windows means less waiting, but now you need a process to make sure the branch's total cash still adds up correctly across all of them.

saying these in an interview costs you the question

  • proposes only 'add more database servers' without addressing that a single stream's writes are inherently serialized regardless of hardware
  • doesn't recognize sharding trades away a precise global invariant for reduced contention
  • assumes commutative operations are a free universal fix applicable to any write
  • ignores that async/intent-based decoupling introduces its own reconciliation and failure-handling burden

context