CQRS splits the code that changes state from the code that answers questions. Event sourcing stores every change as an immutable event and derives current state from them. Together they give you an audit log that cannot drift from reality, read models shaped exactly for each screen, and the ability to build a new view of history months later. They also give you new failure modes that diagrams skip: two writers racing on the same entity, a projector that silently skips events, a user who saves a form and then cannot see the change, and a rebuild that takes all weekend.

This article builds the combined architecture on a single PostgreSQL database, because that is the smallest setup that shows every moving part and the one most teams should start with. It covers the events table and the constraint that enforces optimistic concurrency, the command handler, the commit-order gap that makes naive projectors lose events, checkpointed projectors, read-your-writes, rebuilds, snapshots and schema evolution, with a wallet as the running example. The conceptual case for each pattern is in the separate CQRS and event sourcing pages linked at the end.

The architecture in one picture

Command APIWithdraw, DepositCommand handlerload, decide, appendappend v+1events tableUNIQUE (stream_id, version)read streampoll below xminProjectorcheckpoint in same txRead modelsbalances, statementsQuery APIwaits for positionposition tokenSnapshots (optional)state at version NRebuild: new read-model table, replay from zero, switchWrites are checked against the event store only; reads never touch it.
Commands are checked against the event stream and appended; projectors build read models from committed events; queries read only the read models.

The event store table

The event store is an append-only table. Each event belongs to a stream, usually one aggregate such as one wallet, and has a version that counts from 1 within that stream. A unique constraint on stream and version is the whole concurrency mechanism: two writers that both read version 7 and both try to insert version 8 cannot both succeed.

CREATE TABLE events (
  global_position bigserial   PRIMARY KEY,
  stream_id       text        NOT NULL,
  version         integer     NOT NULL,
  event_type      text        NOT NULL,
  schema_version  smallint    NOT NULL DEFAULT 1,
  payload         jsonb       NOT NULL,
  metadata        jsonb       NOT NULL DEFAULT '{}',   -- causation id, user, request id
  tx_id           xid8        NOT NULL DEFAULT pg_current_xact_id(),
  recorded_at     timestamptz NOT NULL DEFAULT now(),
  UNIQUE (stream_id, version)
);
CREATE INDEX events_tx_order ON events (tx_id, global_position);

Grant the application role INSERT and SELECT, plus USAGE on the position sequence, and nothing else; no code path should update or delete an event. The tx_id column records the writing transaction and exists for the projector problem described below; the xid8 type and pg_current_xact_id() need PostgreSQL 13 or later. Event payloads should be facts in past tense (FundsWithdrawn, not WithdrawFunds) and carry what a reader needs without looking anything up, because a lookup at replay time returns today's data, not the data as it was.

The command handler: load, decide, append

A command handler does three things: load the stream and fold its events into state, decide whether the command is allowed, and append the resulting events at the next version. There is no lock while the decision is made. The unique constraint catches the race afterwards, and the handler retries from the top.

from psycopg.errors import UniqueViolation
from psycopg.types.json import Json

def apply(state, event_type, payload):
    if event_type == "WalletOpened":
        return {"balance": 0, "open": True}
    if event_type == "FundsDeposited":
        return {**state, "balance": state["balance"] + payload["amount"]}
    if event_type == "FundsWithdrawn":
        return {**state, "balance": state["balance"] - payload["amount"]}
    return state

def withdraw(conn, wallet_id, amount, request_id, attempts=3):
    for _ in range(attempts):
        try:
            with conn.transaction():            # rolls back if the insert fails
                rows = conn.execute(
                    "SELECT version, event_type, payload FROM events "
                    "WHERE stream_id = %s ORDER BY version", (wallet_id,)).fetchall()
                state, version = None, 0
                for version, event_type, payload in rows:
                    state = apply(state, event_type, payload)
                if not state or not state["open"]:
                    raise Rejected("no such wallet")
                if state["balance"] < amount:
                    raise Rejected("insufficient funds")
                return conn.execute(
                    "INSERT INTO events (stream_id, version, event_type, payload, metadata) "
                    "VALUES (%s, %s, 'FundsWithdrawn', %s, %s) RETURNING tx_id, global_position",
                    (wallet_id, version + 1, Json({"amount": amount}),
                     Json({"request_id": request_id}))).fetchone()      # (tx_id, position) token
        except UniqueViolation:
            continue                            # someone appended first: reload and decide again
    raise Conflict("too much contention on " + wallet_id)

Make the insert the first write in the transaction. That keeps the transaction ID, which PostgreSQL assigns at the first write, later than the commit of every event the handler read, which the projector relies on. Return the transaction ID and position to the caller; the pair is the read-your-writes token, because projectors order by that pair, not by position alone. Deduplicate retried client requests by storing the request ID in metadata and checking it, or with a separate idempotency table, because the event store itself will happily accept the same withdrawal twice at two versions.

The commit-order gap

The obvious projector reads events with global_position greater than its checkpoint, in position order. It loses events. A bigserial value is taken when the row is inserted, not when the transaction commits. Transaction A inserts position 101 and is still running; transaction B inserts 102 and commits. The projector sees 102, processes it, and stores checkpoint 102. Then A commits, and 101 is behind the checkpoint forever. Nothing errors. A balance is simply wrong, and a rebuild from zero fixes it, which makes it look like a fluke.

The fix is to read only events from transactions that can no longer be in flight. pg_snapshot_xmin(pg_current_snapshot()) is the oldest transaction ID still running; every transaction below it has committed or rolled back, so no new rows with those IDs can appear. Order and checkpoint by the pair of transaction ID and position:

SELECT global_position, tx_id, stream_id, version, event_type, schema_version, payload
FROM events
WHERE (tx_id, global_position) > (%(last_tx)s, %(last_pos)s)
  AND tx_id < pg_snapshot_xmin(pg_current_snapshot())
ORDER BY tx_id, global_position
LIMIT 500;

The cost is latency: one long-running transaction anywhere in the database holds back every projector until it ends, so keep transactions short and alert on old ones. Two alternatives exist. Serialising all appends through one lock makes positions commit in order but caps write throughput. Moving events to a log such as Kafka through an outbox gives you an ordered partition per key instead. Whichever you choose, projectors should also check per-stream versions and refuse to apply version N+2 before N+1, which turns any remaining ordering bug into a loud stall instead of silent corruption.

Projectors and checkpoints

A projector turns events into a read model. When the read model lives in the same PostgreSQL database, update the read model and the checkpoint in one transaction; a crash then either keeps both or loses both, and the batch is simply re-read. That gives exactly-once effects without any deduplication logic.

def run_projector(conn, name, handlers):
    while True:
        with conn.transaction():
            last_tx, last_pos = conn.execute(
                "SELECT last_tx, last_pos FROM checkpoints WHERE name = %s FOR UPDATE",
                (name,)).fetchone()
            batch = conn.execute(READ_BELOW_XMIN, {"last_tx": last_tx, "last_pos": last_pos}).fetchall()
            for e in batch:
                handlers[e.event_type](conn, upcast(e))   # e.g. UPDATE wallet_balances ...
            if batch:
                conn.execute("UPDATE checkpoints SET last_tx = %s, last_pos = %s WHERE name = %s",
                             (batch[-1].tx_id, batch[-1].global_position, name))
        if not batch:
            wait_for_notify_or_timeout(conn, "events", seconds=1)

The FOR UPDATE on the checkpoint row means two copies of the same projector cannot interleave; the second waits, which gives you a simple active-passive pair. When a read model lives elsewhere (a search index, a cache), the checkpoint cannot share its transaction, so make every write idempotent instead: key documents by stream ID and store the stream version on them, and skip any event whose version is not newer. A trigger that calls pg_notify on insert lets idle projectors wake up quickly instead of polling hard.

Read-your-writes

With CQRS the query side lags. A user who withdraws 50 and is redirected to their balance page may see the old balance. There are three honest answers. The command can return the new state it computed, and the UI shows that. The command can return its token, and the query API waits, up to a timeout such as two seconds, until the projector's checkpoint pair is at or past the token, compared as a tuple (last_tx, last_pos) >= (token_tx, token_pos), then answers. Or the screen that must be exact can read from the event stream directly by folding the one stream it needs, which is cheap for a single aggregate. Pick per screen; most screens tolerate a second of lag and only a few need the token.

Rebuilds, snapshots and schema evolution

Rebuilds. Never fix a read model by editing it. To change a projection, create a new table, run the new projector version from position zero until it catches up, then switch the query API to it (a view or a config flag) and drop the old table later. Measure replay speed early: at 5,000 events per second, 200 million events take about 11 hours, and that number decides whether rebuilds are routine or a project.

Snapshots. If a stream grows to thousands of events, loading it for every command gets slow. Store a snapshot of the folded state with its version every few hundred events, load the newest snapshot plus the events after it, and treat snapshots as a disposable cache: tag them with the code version that produced them and throw them away when the fold logic changes. Many aggregates never need one; model short-lived streams (one statement period, one order) before reaching for snapshots.

Schema evolution. Events are immutable, so old shapes live forever. Add fields with defaults, never change the meaning of an existing field, and convert old versions to the current shape when reading, a step usually called upcasting. The schema_version column tells the upcaster which conversion to apply. When the meaning really changes, write a new event type.

Worked example: two withdrawals and a new view

Wallet w-42 has events up to version 7 and a balance of 100. Two withdrawal requests of 80 arrive at the same moment on different servers. Both load seven events, both compute 100, both decide 80 is allowed, and both try to insert version 8. The first insert commits; the second hits the unique constraint, rolls back and retries. On the retry it loads eight events, computes a balance of 20 and rejects the withdrawal with insufficient funds. No lock was held while either decision was made, and the invariant held.

The successful handler returned the token (transaction 88,412, position 9,001,337). The projector's read is held back for a moment by a slow unrelated transaction, so the balance table still says 100. The client calls GET /wallets/w-42?after=88412:9001337; the query API sees the checkpoint pair still below the token, waits 300 milliseconds, the slow transaction ends, the projector applies the event and the API answers 20. A month later finance wants a daily statement view. A new projector replays every wallet's history into a statements table, and the view exists for all past days, not just from the day it was deployed.

Failure modes

SymptomCauseFix
Read model occasionally misses an event; rebuild fixes itProjector reads by position past an uncommitted lower positionRead below snapshot xmin, checkpoint the (tx_id, position) pair
Projector lag climbs for minutes, then clearsA long transaction holds xmin backStatement and idle-in-transaction timeouts; alert on transaction age
Same withdrawal applied twiceClient retried and no idempotency checkStore and check the request ID
Many Conflict errors on one streamHot aggregate written by many usersSplit the aggregate, or queue commands per stream
Commands slow on old entitiesLong streams loaded in fullSnapshots, or model shorter-lived streams
Replay crashes on old eventsCode assumes the current schemaUpcasters with tests over a sample of every schema version
Users see stale data after savingQuery side not waiting for its own writeReturn new state, or wait on the position token

Trade-offs

Use the combination where history is the product (ledgers, orders, compliance workflows) and where several very different read shapes are needed. Do not use it for a CRUD admin screen. The costs are real: two models to maintain, eventual consistency to explain, schema evolution forever, and personal data that cannot simply be deleted from an immutable log. The usual answer to that last one is to encrypt personal fields with a per-person key and delete the key when asked; design for it before the first event is written.

Starting on one PostgreSQL database keeps transactions, backups and operations simple and is fast enough for a large share of systems. Move to a dedicated event store or a log when you need fan-out to many services or volumes one database cannot hold, and add an outbox so that publishing never disagrees with the store.

Keep learning: CQRS architecture, event sourcing architecture, the outbox pattern and Kafka partitions for when events move to a log.

What to do next

  1. Pick one aggregate whose history matters and model its events in past tense with self-contained payloads.
  2. Create the events table with UNIQUE (stream_id, version) and an insert-only role.
  3. Write the command handler with load, decide, append and retry on unique violation, plus request-ID deduplication.
  4. Build one projector that reads below snapshot xmin and commits checkpoint and read model together.
  5. Decide per screen how read-your-writes is handled, and implement the position token for the ones that need it.
  6. Time a full replay on production-sized data and write the rebuild runbook.
  7. Add an upcaster test that replays a sample of every event schema version.
Key takeaway: Build CQRS with event sourcing on an append-only events table whose unique stream-and-version constraint provides optimistic concurrency, and command handlers that load, decide, append and retry. Do not let projectors read by position alone: sequence values are assigned before commit, so read only below the snapshot xmin and checkpoint the transaction-and-position pair, in the same transaction as the read model. Handle read-your-writes per screen with a position token, rebuild projections into new tables, treat snapshots as cache, and upcast old events on read.