skip to content

In a high-throughput event-sourced system where multiple processes might concurrently write events and take snapshots for the same aggregate, what specific race conditions and failure modes can corrupt or invalidate a snapshot, and how do you design around them?

level: principalimportance: should knowfreq 35%

answer

  1. read-state-vs-append race -> atomic derive from event store
  2. concurrent snapshot writers -> optimistic concurrency / version-keyed storage
  3. torn writes -> validate + fallback
  4. event store = source of truth, snapshot = disposable cache
  5. lost-update pattern

basics

~20 s

If two things try to save or read a snapshot for the same object at the same time, you can end up with a snapshot that doesn't match reality. Fix this by always deriving snapshots from a confirmed, already-saved event position, and by checking that the snapshot's numbers actually line up before trusting it.

solid answer

~50 s

The core race is between the process appending new events and the process (often async) reading current state to build a snapshot: if the snapshot job reads state and version separately, or from an in-memory copy rather than the durable event store, a concurrent append can land in between, producing a snapshot tagged with a version that doesn't match what it actually contains. A second race is between two concurrent snapshot writers themselves (e.g., two instances both decide to snapshot the same aggregate), which can produce out-of-order or overwritten snapshot writes if there's no serialization/last-writer-wins discipline. The fix is to derive every snapshot from a single, durable read of the event store at a specific version (never from in-process state), use optimistic concurrency (only accept a snapshot write if its version is newer than what's stored, or use the aggregate's version as part of the storage key), and treat any detected inconsistency as 'discard and fall back to replay' rather than attempting to reconcile.

go deeper

for a junior

Not expected to reason about this; awareness that concurrency can cause bugs generally is sufficient.

for a middle

Should recognize that concurrent writes could cause snapshot problems, even without naming the specific mechanisms.

for a senior

Should describe at least one specific race (e.g., read-state-vs-append) and propose deriving snapshots atomically from the event store.

for a principal

Should distinguish all three failure modes, design optimistic-concurrency or version-keyed storage, and connect the design to the broader cache-staleness principle.

## Snapshotting becomes a concurrent system At scale, snapshotting stops being a single-threaded implementation detail and becomes a concurrent system with its own race conditions, because the process taking a snapshot and the process(es) appending new events to the same aggregate's stream are usually not the same code path, and often not even the same machine. Three distinct races matter here, and each has a different concrete failure signature. ## Race one — reading state while an append lands The first and most dangerous is the **read-state-vs-concurrent-append race**. A snapshot job — whether it's inline in the command handler, a background worker polling the event store, or a stream consumer subscribed to appends — needs to capture 'the aggregate's state, and the version it corresponds to' as a single, atomic fact. If instead it reads the version counter and the state separately (for example: read current version = 5,000, then separately fold/read events to build state, then write the snapshot tagged 5,000), a new event can be appended between those two reads. The resulting snapshot is now tagged 5,000 but actually reflects only 4,999 events' worth of state — or the reverse, tagged 5,000 but with a stray in-progress write already partially reflected, depending on exactly where the race lands. Either way, a reader that later loads this snapshot and replays events after version 5,000 will either skip event 5,000's effects entirely or double-apply it, both of which corrupt state silently, with no error raised anywhere. The fix is to make the state-and-version capture atomic with respect to the event store: derive the snapshot exclusively from a single durable read (e.g., 'the state after folding events 1 through the store's confirmed current version, both determined in that one read'), never from separately-tracked in-memory counters that could be stale relative to a concurrent writer. ## Race two — two snapshot writers for the same aggregate The second race is between multiple concurrent snapshot writers for the same aggregate — plausible in a horizontally scaled system where more than one process instance might independently decide 'this aggregate needs a snapshot now' (for example, two replicas of a snapshot-triggering consumer group both processing the same event, due to an at-least-once delivery guarantee and a bug in dedup logic). If both write to the same snapshot storage key without coordination, you can get an out-of-order write: writer A, which read at version 4,800, finishes writing after writer B, which read at version 4,900, and A's older, stale write overwrites B's newer, correct one — the **'lost update' problem**, familiar from any concurrent-write-to-shared-key scenario. Two defenses: - **Optimistic concurrency on the snapshot write itself** — the standard one: include the version being written in a conditional write (only succeed if the stored version is currently lower than the version this writer is about to write), so a stale writer's attempt is rejected rather than silently overwriting a newer, valid snapshot. - **Keying snapshots by version** — alternatively, some designs sidestep the problem structurally by keying snapshots by version rather than overwriting a single 'latest' slot (store 'snapshot-4800' and 'snapshot-4900' as distinct records, and have the loader explicitly query for the highest version at or below what it needs); this trades extra storage and a slightly more complex read query for eliminating the overwrite race entirely. ## Torn and partial writes The third failure mode is more operational than race-condition-specific: **partial or torn writes**, where a snapshot write is interrupted mid-flight (a process crash, a network partition to the snapshot store, an out-of-memory kill) and leaves a corrupted or incomplete payload behind — readable as bytes but not valid under the deserializer, or valid-looking but actually missing trailing fields. Defensive loaders should validate a snapshot's integrity before trusting it (a checksum, a required-fields check, or simply catching deserialization exceptions) and treat any validation failure identically to a missing snapshot: discard and fall back to full replay from the event store, which remains authoritative and unaffected by any of these snapshot-side failures. ## The unifying principle The unifying design principle across all three is that the event store, not the snapshot store, is the single source of truth, and every defense here is really the same move applied in different places: - never let a snapshot be trusted unless its provenance can be verified against the event store's own version ordering; - always have 'discard this snapshot, replay from events instead' as a cheap, always-correct fallback rather than attempting to repair or reconcile a suspect snapshot in place. This is analogous to how distributed caches generally handle staleness — the cache (snapshot) is allowed to be wrong or missing because there's always an authoritative, slower path (the event store) to fall back to; the engineering effort goes into detecting staleness/corruption reliably and cheaply, not into making the cache itself perfectly consistent, which would reintroduce much of the coordination cost snapshotting was meant to avoid in the first place.

  • Why is keying snapshots by their version number (rather than always overwriting a single 'latest' slot) sometimes preferred despite using more storage?
    It structurally eliminates the lost-update race between concurrent writers, since each write lands at a distinct key rather than contending for the same slot; the loader just queries for the highest version at or below what it needs, and old versions can be pruned separately on a retention schedule rather than needing write-time coordination.
  • What's the danger of 'repairing' a detected corrupted snapshot instead of discarding it?
    Any repair logic is an additional piece of code that has to be correct under exactly the failure conditions (crashes, partial writes) that are hardest to test, whereas discard-and-replay reuses the event store's already-proven correctness; repair attempts risk quietly manufacturing plausible-but-wrong state instead of failing safe.
  • How does at-least-once event delivery interact with the concurrent-snapshot-writer race?
    At-least-once delivery means the same triggering event can be processed more than once by different consumer instances, so a naive 'snapshot when I see this event' consumer can end up with two instances racing to write a snapshot for the same version; deduplication and optimistic-concurrency writes are both needed, not just one.

Like two people updating the same shared spreadsheet cache at once without checking each other's timestamp — whoever saves last silently overwrites the other's more current number unless you add a 'only save if my version is newer' check.

saying these in an interview costs you the question

  • assumes snapshot writes are inherently atomic without design effort
  • doesn't distinguish the read-race from the writer-vs-writer race
  • proposes fixing a corrupted snapshot in place rather than falling back
  • unaware that at-least-once delivery can trigger duplicate snapshot writers

context