Event Sourcing and CQRS
Storing state as an append-only log of events, rebuilding it by replay, and serving reads from separate projections, plus the costs: versioning and stale reads.
Event sourcing stores what happened instead of what is: every change is appended to a log as an immutable event (ItemAdded, OrderPaid), and the current state of anything is a fold over its events. CQRS (Command Query Responsibility Segregation) splits the model that handles writes from the models that answer reads; with event sourcing, the read models are projections built by consuming the log. You gain full history, replay and cheap new views, and you pay with event versioning and reads that lag the writes.
Context
The idea is older than software. An accountant never erases a ledger line; a mistake is fixed by a correcting entry, and the balance is the sum of all entries. A bank statement, a Git history and a database's write-ahead log work the same way. Martin Fowler described the pattern for application code as "Event Sourcing" in 2005. The read/write split goes back to Bertrand Meyer's CQS (command-query separation) principle from 1988: a method either changes state or returns data, never both. Greg Young lifted that from methods to whole models and named it CQRS around 2010, usually together with event sourcing and DDD (domain-driven design). Tooling followed: EventStoreDB (now KurrentDB), Axon on the JVM (Java virtual machine), Marten on PostgreSQL for .NET.
You have met the shape if you have used Redux (state is a reducer over actions), read an audit log that can explain every balance, or rebuilt a search index by replaying a Kafka topic. The core fits in three lines:
const events = [
{type: 'Deposited', amount: 100},
{type: 'Withdrawn', amount: 30},
{type: 'Deposited', amount: 5},
]
const balance = events.reduce(
(sum, e) => e.type === 'Deposited' ? sum + e.amount : sum - e.amount, 0)
// balance === 75, and the log explains how it got there- Event
- An immutable fact in the past tense: OrderPlaced, ItemAdded. It records a decision that was already made, so it is never rejected or edited.
- Command
- A request in the imperative: PlaceOrder, AddItem. It can be rejected. If accepted, it produces one or more events.
- Stream / aggregate
- The events of one entity (order 8f3c), with a version number per event. The aggregate is the in-memory object rebuilt from its stream that decides whether a command is allowed.
- Projection / read model
- A table, document or index built by consuming events and shaped for one kind of query. Disposable: it can be rebuilt from the log.
- Snapshot
- A saved copy of an aggregate’s state at some version, so loading it does not replay thousands of events.
Why it matters
Overwriting a row loses the reason it changed. Domains where the history is the product (payments, ledgers, inventory, compliance, logistics) end up bolting audit tables onto CRUD (create, read, update, delete) and still cannot answer "what did this order look like last Tuesday?". Event sourcing makes that history the source of truth, and any new report becomes a new projection over the existing log. The price is real: every reader must cope with data that is a moment behind, events stored years ago must stay readable forever, and simple features take more code. It is a strong fit for a few core domains and a costly default for everything else.
From command to events to read models
On the write side, a command loads one aggregate by replaying its stream, checks the business rules against that state, and appends new events. Nothing else is written: no current-state row, no read table. On the read side, projections subscribe to the log and update whatever shapes the queries need. Queries never touch the aggregate or the log directly.
The aggregate is a pure fold. Each event type has one function that moves state forward, and that function contains no validation: by the time an event exists, the decision was already made and must always apply.
type OrderEvent =
| {type: 'OrderPlaced'; customerId: string}
| {type: 'ItemAdded'; sku: string; qty: number; price: number}
| {type: 'OrderPaid'; amount: number}
type Order = {status: 'new' | 'placed' | 'paid'; total: number}
const initial: Order = {status: 'new', total: 0}
export function apply(s: Order, e: OrderEvent): Order {
switch (e.type) {
case 'OrderPlaced':
return {...s, status: 'placed'}
case 'ItemAdded':
return {...s, total: s.total + e.qty * e.price}
case 'OrderPaid':
return {...s, status: 'paid'}
}
}
// current state = left fold over the stream, from version 1
export const load = (events: OrderEvent[]) => events.reduce(apply, initial)The event store and concurrent writes
An event store needs surprisingly little: append events to a stream, read a stream in order, and read all events in global order for projections. A relational table does it. The one property that matters is the version per stream: a command remembers the version it loaded, and the append only succeeds if nobody appended after that. That is optimistic concurrency, enforced here by the primary key.
CREATE TABLE events (
stream_id text NOT NULL, -- 'order-8f3c'
version integer NOT NULL, -- 1, 2, 3 ... within the stream
type text NOT NULL, -- 'ItemAdded'
data jsonb NOT NULL,
metadata jsonb NOT NULL DEFAULT '{}', -- user, causation id
global_pos bigserial, -- order across all streams
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (stream_id, version) -- rejects a concurrent append
);export async function handleAddItem(cmd: AddItem) {
const stream = `order-${cmd.orderId}`
const {events, version} = await store.read(stream)
const order = load(events)
// business rules run against rebuilt state
if (order.status !== 'placed') throw new Error('order is not open')
const newEvents: OrderEvent[] = [
{type: 'ItemAdded', sku: cmd.sku, qty: cmd.qty, price: cmd.price},
]
// inserts versions version+1..n in one transaction; a duplicate
// (stream_id, version) means someone else appended first
await store.append(stream, newEvents, {expectedVersion: version})
}
// on ConcurrencyError: reload the stream and run the command again- 1Alice and Bob both load
order-8f3cat version 3. - 2Alice's command appends version 4. Bob's command tries to append version 4 as well, and the primary key rejects it.
- 3Bob's handler reloads the stream (now version 4), re-checks the rules against the new state, and appends version 5, or rejects the command if Alice's change made it invalid.
Snapshots
Replaying a stream of 20 events costs nothing; replaying 200,000 on every command does. A snapshot stores the folded state together with the version it reflects, for example every 500 events. Loading becomes "latest snapshot, then events after its version". Snapshots are a cache: they can be deleted and rebuilt, and they must be discarded whenever the fold logic changes. Long streams are often a modelling smell first: an account that lives forever can be closed per period, with the closing balance carried into a new stream, the way accountants close books.
Projections, stale reads and CQRS on its own
A projection is a loop: read events from a checkpoint, update a read model, save the checkpoint. Because each projection keeps its own checkpoint, adding a new report is a matter of writing one and letting it replay from position zero, and a broken one can be dropped and rebuilt. Delivery is at-least-once, so handlers must be idempotent, for example by storing the last applied position per row.
for await (const e of store.readAll({from: await checkpoint()})) {
if (e.type === 'OrderPaid') {
await db.query(
`UPDATE order_list SET status = 'paid', last_pos = $2
WHERE id = $1 AND last_pos < $2`, // idempotent re-run
[e.streamId, e.globalPos])
}
await saveCheckpoint(e.globalPos)
}The gap between the append and the projection catching up is usually milliseconds, but it is never zero, and that produces the classic bug: the user saves, the page reloads the list, and the new item is missing. Fixes, from simplest: update the UI from the command's own result; return the new global position from the command and have the query wait until the projection's checkpoint passes it; or, for the one screen that needs it, read the aggregate itself.
CQRS without event sourcing
CQRS only says that writes and reads use different models. That can be one PostgreSQL database where commands go through a rich domain model and queries read denormalized tables or views updated in the same transaction, with no events and no lag. Event sourcing, in turn, is almost always paired with CQRS, because the log is terrible to query directly: "all unpaid orders over $100" would mean replaying every stream.
| CRUD | CQRS | ES + CQRS | |
|---|---|---|---|
| Truth | Current-state rows | Current-state rows | The event log |
| Writes | UPDATE in place | Commands via a write model | Append events to a stream |
| Reads | Same tables | Separate read tables or views | Projections built from events |
| History | Lost unless audited | Lost unless audited | Complete, replayable |
| Read lag | None | None if same transaction | Always some, eventually consistent |
| Cost | Lowest | Moderate | Highest: versioning, tooling |
Changing events and deleting data
Events written in 2026 must still load in 2030, so the stored history is a public API you can never break. Additive changes (a new optional field) are free. For anything else, the usual tool is an upcaster: a function that converts an old event shape into the current one as it is read, so the fold only knows the latest version. Renames and meaning changes get a new event type rather than an edited old one, and a large restructuring is done by copying a stream into a new one through a transformation, never by updating rows in place.
// v1 stored {price} in dollars; v2 stores {priceCents, currency}
function upcast(raw: StoredEvent): OrderEvent {
if (raw.type === 'ItemAdded' && raw.schema === 1) {
return {...raw.data, priceCents: Math.round(raw.data.price * 100),
currency: 'USD'}
}
return raw.data
}Immutability collides with the GDPR (General Data Protection Regulation) right to erasure. Two approaches work. Keep personal data out of events and reference it by ID from a store you can delete from. Or use crypto-shredding: encrypt personal fields with a key per person, keep the keys elsewhere, and delete the key to make that data unreadable everywhere it was copied, including backups. Either way, projections that copied the data must be cleaned or rebuilt too.
Pitfalls
- Event sourcing the whole system
Most of an application is CRUD: settings, profiles, catalog edits. Event sourcing them adds projections, versioning and stale reads with no benefit, because nobody needs their history. Apply it to the few domains where history and auditability are the point, and keep the rest as plain tables.
- Events that describe data instead of decisions
OrderUpdatedwith the whole row is CRUD written as a log: it hides why the change happened and couples every consumer to the table layout. Name events after business facts (ItemAdded,AddressCorrected) so projections and other services can react to meaning. - Using Kafka as the event store of record
Kafka distributes events well, but it has no per-entity append with an expected version, so it cannot stop two commands from both acting on the same stale state, and loading one aggregate means scanning a partition. Keep the source of truth in a store with optimistic concurrency, and publish to Kafka from it, for example through an outbox.
- Ignoring the read lag in the UI
Teams test locally, where projections catch up instantly, and ship screens that save and immediately reload a projection. Under load the reload wins the race and users see their change vanish. Design every write path with an explicit answer to "what does the user see right after this?".
- Side effects inside the fold
The apply function runs on every load and every replay. If it sends an email or calls an API, rebuilding a projection or loading an aggregate repeats that action. Side effects belong in event handlers with their own checkpoints and idempotency, never in the code that rebuilds state.
Interview questions
Q1What is event sourcing, and how is it different from storing current state?
The source of truth is an append-only log of events, and current state is derived by folding over them. A CRUD system overwrites the row and loses how it got there; an event-sourced one can reconstruct any past state, explain every change and build new views by replaying. In exchange, reads come from projections that lag slightly, and stored events can never change shape freely.
Q2Are CQRS and event sourcing the same thing?
No. CQRS separates the write model from the read models, which can be done in one database with no events at all. Event sourcing is a storage choice for the write side. They are usually combined because an event log is hard to query, so event-sourced systems almost always need separate read models.
Q3Walk me through implementing a command handler in an event-sourced system.
Load the aggregate's stream and its current version, fold the events into state, and run the business rules against that state. If the command is valid, create the new events and append them with the loaded version as the expected version, in one atomic write. If the append fails on a version conflict, reload and retry or reject. The handler writes nothing else; read models are updated by projections.
Q4What happens when two commands modify the same aggregate at the same time?
Both load version N, both decide, and both try to append version N + 1. The store accepts the first and rejects the second because the expected version no longer matches, typically via a unique key on stream and version. The loser reloads, re-checks its rules against the new state and either appends N + 2 or fails, so no decision is made on stale state.
Q5A user creates an item and the list page does not show it. Why, and how do you fix it?
The list is a projection that had not processed the new event yet when the page queried it. Options: update the UI from the command result without re-reading; return the event's global position and make the query wait until the projection has passed it; or read that screen from the aggregate. What you should not do is pretend the lag does not exist.
Q6How do you change the schema of an event that has already been stored?
Never edit stored events. Add optional fields for additive changes, and for breaking changes add a schema version and an upcaster that converts old shapes to the current one at read time, so the domain code only sees the latest version. If the meaning changes, introduce a new event type, and for large restructurings copy streams into new ones through a transformation.
Q7How do you delete a user’s personal data from an immutable event log?
Either keep personal data out of events and reference it from a deletable store, or encrypt it with a per-user key and delete the key, which is called crypto-shredding. The events stay, but the personal fields become unreadable everywhere, including backups and copies. Projections that hold the plain data must be cleaned or rebuilt as part of the same erasure.
Q8When would you not use event sourcing?
For CRUD-shaped domains where nobody needs the history, for small teams without the capacity to run projections and event versioning, and where users need every read to reflect their write immediately across many screens. I would reach for it in ledgers, orders, inventory or compliance-heavy domains, and even then only for those bounded contexts, not the whole system.
- Event sourcing stores immutable events; current state is a fold over a stream, so history, audit and replay come for free.
- Concurrency is a per-stream version: append with the version you loaded, and a conflict means reload and decide again.
- CQRS separates write and read models and does not require events; event sourcing practically requires CQRS because the log is hard to query.
- Projections are disposable, rebuildable and always slightly behind. Design every write path for the read-your-writes gap.
- Stored events are a permanent contract: evolve them additively or with upcasters, and handle erasure with external PII or crypto-shredding.
- Use it for the few domains where history is the product, not as a default for the whole system.