skip to content

When two commands try to append to the same event stream at nearly the same time, how does optimistic concurrency control in an event store prevent one write from silently clobbering the other?

level: middleimportance: must knowfreq 75%

answer

  1. expected version on every append
  2. atomic compare-and-swap at write time
  3. conflict = reject, not merge
  4. client reloads + retries, not blind resubmit
  5. hot aggregate = high conflict rate

basics

~20 s

Each write says 'I'm adding this assuming the stream currently has N events.' If someone else already added an event in between, the store rejects the write instead of accepting it blindly, so the app can retry with fresh data.

solid answer

~50 s

Optimistic concurrency control works by having every append call include the expected version (or expected stream position) the caller believes the stream is currently at, based on the events it read before deciding what to append. The event store atomically checks that expected version against the stream's actual current version at write time; if they match, the new events are appended and the version advances. If they don't match — because another writer appended events in between the read and the write — the store rejects the append with a version-conflict error rather than silently interleaving or overwriting anything. The client is then expected to reload the stream, re-run its business logic against the fresh state, and retry the append with the new expected version. This gives per-stream serializability without needing a lock held across the read-decide-write cycle, which is why it's called 'optimistic' — no lock is taken, conflicts are just detected and retried.

go deeper

for a junior

Should be able to explain, in plain terms, that the store checks 'is this still version 7' before writing and rejects if not.

for a middle

Should know the retry-with-reload pattern and why blindly resubmitting the same event is wrong.

for a senior

Should compare OCC to pessimistic locking with real throughput/contention reasoning and identify hot-aggregate contention as a design smell.

for a principal

Should propose concrete redesigns for high-contention streams (sharding, commutative operations) and reason about bounded-retry/backoff strategy at the platform level.

## The race it prevents Optimistic concurrency control (OCC) in an event store solves a specific race: two commands targeting the same aggregate stream — say, `account-482` — are processed concurrently by two different request threads, or two different service instances, and each one reads the stream, decides on new events to append based on what it read, and then writes. Without any coordination, if both reads happen before either write, both writers believe the stream is at version 7 and both decide it's safe to append their own event as position 8. If the store just accepted both appends blindly, you'd end up with either: - a **lost update** (one writer's business decision silently overwritten), or - depending on the store's design, two conflicting events both claiming to be position 8, corrupting the ordering guarantee the whole model depends on. ## How the check works Mechanically, OCC works by making 'expected version' an explicit, required parameter of the append call. When a command handler loads a stream to decide what to do, it reads not just the events but also the stream's current version/sequence number (say, 7). After running its business logic and producing new events to append, it calls something like `append(streamId: account-482, expectedVersion: 7, events: [...])`. 1. The store, at the moment of the write, atomically compares the `expectedVersion` the caller supplied against the stream's actual current version in storage. 2. If they match, the append succeeds and the version moves to 8. 3. If a different writer already advanced the stream to 8 in the meantime, the comparison fails, and the store rejects the append entirely with a concurrency-conflict error — it **does not partially apply, does not silently renumber, does not merge**. This check-and-append has to be a single atomic operation at the storage layer, otherwise you've just moved the race into the check itself. ## Why optimistic rather than pessimistic The reason OCC — rather than pessimistic locking — is the standard approach is that most event-sourced workloads have relatively low contention per aggregate: two commands hitting the exact same account in the same instant are rare, so paying the cost of a lock held across an entire read-decide-write round trip (including any external calls the business logic makes) would throttle throughput for a conflict that usually doesn't happen. OCC instead assumes success, and only pays a cost — a rejected write and a retry — in the rare case a conflict actually occurred. This retry is cheap and mechanical: on a version-conflict error, the caller - reloads the stream (getting the fresh version and the events it missed), - re-executes its business/validation logic against that fresh state, - and re-attempts the append with the new expected version. Well-designed command handlers make this retry loop transparent to the end user, often completing in single-digit milliseconds. ## The trade-off it pushes onto the client The trade-off is that OCC pushes complexity onto the client: every writer must be prepared to retry, and that retry logic has to be correct — **re-running validation against fresh state, not just re-submitting the stale event payload**, since the fresh state might mean the original decision is no longer valid at all (e.g., a withdrawal that was fine against a balance of 400 might now overdraw a balance of 40 after a concurrent charge landed). Systems that get this wrong either loop forever under contention, or — worse — catch the conflict error and blindly retry the exact same append with a bumped version number, which reintroduces the lost-update bug OCC was supposed to prevent, because the retried event was decided against stale data. ## Failure modes Failure modes in production tend to cluster around 'hot' aggregates: - a single account, inventory SKU, or feature-flag stream that many callers write to concurrently will see a high conflict rate, and if retry logic isn't bounded (no backoff, no max-attempts), those retries themselves become a **load-amplifying feedback loop** — more concurrent retries generate more conflicts, generate more retries. - Another common failure is a client caching a stream's version across requests (e.g., in a long-lived in-memory object) and reusing a stale `expectedVersion` long after the underlying stream moved, which produces a wall of spurious conflict errors on writes that would have been fine had the client reloaded first. ## What a real client API looks like A concrete real-world example: EventStoreDB's client API takes an `expectedRevision` on every `appendToStream` call — you can pass a specific revision number, or sentinel values like `NO_STREAM` (I expect this stream doesn't exist yet, for a first-time create) or `STREAM_EXISTS` (I expect it exists, but don't care at what revision) — and returns a `WrongExpectedVersionException` on mismatch, which application code is expected to catch and handle with a reload-and-retry loop, exactly the pattern described above.

  • What's the difference between optimistic concurrency control here and a database row-level lock (pessimistic locking)?
    Pessimistic locking holds a lock for the entire read-decide-write cycle, blocking other writers until it's released, which is safe but throttles throughput under any concurrency. OCC takes no lock; it lets both writers proceed and only rejects at the final write if the version moved, trading a rare retry cost for no blocking in the common uncontended case.
  • Why is naive retry-with-the-same-event dangerous after a version conflict?
    The event was computed by business logic that ran against stale state; blindly resubmitting it with a bumped version ignores whatever changed in the meantime, which can violate an invariant the fresh state would have caught (e.g., an overdraft). The correct retry reloads the stream and re-runs the decision logic from scratch against current state.
  • How would you handle a stream that sees very high write contention, like a global counter?
    You typically redesign the model to avoid needing strict per-event ordering on a single hot stream — for example, partitioning the counter into per-shard sub-streams and summing them, or moving to a commutative operation that doesn't need optimistic-concurrency serialization at all.

Like editing a shared Google Doc offline: you note 'I'm editing based on version 7,' and if someone else already saved version 8 before you, your save is bounced back instead of silently overwriting their change — you have to pull the new version and reapply your edit.

saying these in an interview costs you the question

  • thinks the event store silently merges concurrent writes
  • doesn't know a version-conflict error must trigger reload+re-decide, not blind resubmit
  • confuses optimistic concurrency with database transaction isolation levels
  • assumes locking has no throughput cost so OCC is 'just extra complexity'

context