Command Query Responsibility Segregation splits a system's model in two: one side accepts commands that change state and enforces the business rules, the other answers queries from data shaped for reading. The idea grows out of Bertrand Meyer's command-query separation for methods; Greg Young applied it to whole models. The payoff is that each side can be designed, stored and scaled for its own job. The cost is that the read side is a copy, and copies are late, duplicated and occasionally wrong.
Most writing about CQRS pairs it with event sourcing. This article deliberately does not. It covers the far more common case: a conventional database on the write side, with changes propagated to read models through an outbox or change data capture. If you already store events, read CQRS with event sourcing instead; the propagation and rebuild problems are different there.
Commands, queries and three levels of separation
A command expresses intent and can be refused: ShipOrder(order_id, tracking). It returns success, a rejection or a conflict, and at most a small result such as the new version. A query returns data and changes nothing. Under CQRS these go to different code paths with different models. The command model is normalised and holds exactly what invariants need: status, totals, version. Read models are denormalised and exist per screen or API: an order list joined with customer names, a search index, a monthly revenue table.
CQRS does not require two databases, a message broker, microservices or events. It is a spectrum, and you should stop at the lowest level that solves your problem:
| Level | What is separated | Consistency | Use when |
|---|---|---|---|
| 1. Code | Command handlers and query handlers are separate classes over the same tables | Immediate | Reads need different shapes but the same data |
| 2. Same database | Read tables or materialised views beside the write tables, updated in the same transaction or by triggers | Immediate or near-immediate | Joins and aggregates are too slow to compute per query |
| 3. Separate stores | Read models in other engines, fed asynchronously | Eventual, with measurable lag | Read load, search or analytics need a different engine, or teams own different views |
Everything below is about level 3, because that is where the hard problems live.
The architecture end to end
Follow one command through. The client sends ShipOrder. The command API loads the order, checks that it is shippable, updates it with an optimistic version check and, in the same transaction, inserts an outbox record carrying the new version and the order's new state. A relay (or a CDC connector reading the database log) publishes that record to a topic keyed by order ID, so all changes to one order land in one partition in commit order. Each projector consumes the topic and upserts its own read store. The query API reads only read stores and contains no business rules.
Three properties of that pipeline drive every later decision: delivery is at least once, so projectors see duplicates; ordering holds only per key, and only if you keep it; and the read side always lags the write side by some amount you must measure.
The command side: invariants and versions
The command side is an ordinary transactional handler with two disciplines: check invariants against current state, and bump a version so concurrent commands conflict instead of overwriting each other. This runs as written against SQLite; PostgreSQL differs only in placeholders.
import json
class Rejected(Exception): pass
class Conflict(Exception): pass
def ship_order(conn, order_id, tracking):
"""Check the invariant, bump the version, record the change. One transaction."""
with conn: # commits on success, rolls back on any exception
row = conn.execute(
"SELECT customer_id, status, total_cents, version FROM orders WHERE id = ?",
(order_id,),
).fetchone()
if row is None:
raise Rejected("no such order")
customer_id, status, total_cents, version = row
if status not in ("PAID", "PARTIALLY_REFUNDED"):
raise Rejected(f"cannot ship an order that is {status}")
new_version = version + 1
changed = conn.execute(
"UPDATE orders SET status = 'SHIPPED', version = ? WHERE id = ? AND version = ?",
(new_version, order_id, version),
).rowcount
if changed == 0:
raise Conflict("order changed concurrently; reload and retry")
after = {"order_id": order_id, "customer_id": customer_id, "status": "SHIPPED",
"total_cents": total_cents, "version": new_version, "tracking": tracking}
conn.execute(
"INSERT INTO outbox (aggregate_id, version, type, payload) VALUES (?, ?, ?, ?)",
(order_id, new_version, "OrderShipped", json.dumps({"after": after})),
)
return new_versionTwo choices in that code matter more than they look. The outbox payload carries the full after state of the fields readers need, not a delta such as "status changed". And it carries the version. Together they let every projector be idempotent and order-tolerant with one comparison, as the next sections show.
Getting changes out: dual writes, outbox, CDC
The tempting implementation is a dual write: commit to the database, then publish to the broker. It fails in both orders. Commit then publish loses the message if the process dies between the two; publish then commit announces a change that may roll back. No retry policy fixes this, because the two systems do not share a transaction.
Two patterns close the gap. With a transactional outbox, the change and the message are one database transaction, and a relay process publishes outbox rows and marks or deletes them; see the outbox pattern for relay design. With change data capture, a connector such as Debezium tails the database's replication log and publishes row changes, so the application writes nothing extra; Debezium CDC to Kafka covers the setup. CDC exposes your table schema as the contract, which couples every read model to write-side column names. A common compromise is CDC on the outbox table itself: the application controls the payload, and the connector replaces the hand-written relay.
Either way the relay can publish a record twice (it crashes after publishing, before marking), and consumers can process a message twice (they crash after applying, before committing the offset). Design for duplicates instead of trying to prevent them.
Projectors that survive duplicates and reordering
A projector applies each change to its read store. Because messages carry full state and a version, the rule is: apply only if the incoming version is newer than the stored one. Duplicates and late arrivals become no-ops.
def apply_change(db, msg):
"""Apply one change to the read model. Safe to call twice, safe out of order."""
row = msg["after"]
cur = db.execute(
"""
INSERT INTO order_summary (order_id, customer_id, status, total_cents, version)
VALUES (:order_id, :customer_id, :status, :total_cents, :version)
ON CONFLICT(order_id) DO UPDATE SET
customer_id = excluded.customer_id,
status = excluded.status,
total_cents = excluded.total_cents,
version = excluded.version
WHERE excluded.version > order_summary.version
""",
row,
)
return cur.rowcount # 0 means stale or duplicate: skippedHere is that function run on SQLite against a realistic delivery sequence for order o-42, including a redelivered message and one that arrives late after a partition rebalance:
| Delivered | Payload | rowcount | Row afterwards |
|---|---|---|---|
| v1 | PLACED, 4200 | 1 | PLACED, 4200, v1 |
| v2 | PAID, 4200 | 1 | PAID, 4200, v2 |
| v2 again | PAID, 4200 | 0 | PAID, 4200, v2 |
| v4 | SHIPPED, 3700 | 1 | SHIPPED, 3700, v4 |
| v3 (late) | PARTIALLY_REFUNDED, 3700 | 0 | SHIPPED, 3700, v4 |
The final row is correct even though v3 was skipped, because v4 already contained the refunded total. Had the messages been deltas ("refund 500"), skipping v3 would lose money and applying it twice would lose it twice. That is why state-carrying messages are the default for state-stored CQRS; deltas need strict exactly-once, in-order processing, which you do not have.
Deletes need the same discipline. A projector that physically deletes the row on OrderDeleted v5 will happily re-create it when a late v4 arrives. Keep a tombstone (a deleted flag plus version) and purge tombstones only after a period longer than any plausible redelivery delay. For side effects that are not upserts, such as sending an email when an order ships, record processed (aggregate, version) pairs as described in idempotency keys.
Freshness: lag budgets and read-your-writes
Measure lag per projector as the time between the write transaction's commit and the moment its change is applied, and alert on a budget per read model: a search index may tolerate a minute, an order-status page perhaps two seconds. Lag is a product decision as much as a technical one, so write the budget down.
The classic complaint is a user who ships an order and sees it still marked Paid. Three answers, chosen per screen: show the state the command returned; have the command return the new version and let the query API wait briefly until the read row's version reaches it, then fall back to a "still updating" hint; or read that one screen from the write database. The version-wait technique is covered in detail on the event-sourcing page linked above; with state-stored CQRS the token is simply the row version. If a user sees several read models, they can also see them disagree with each other; causal consistency explains why and what session guarantees can restore.
Rebuilding a read model without an event log
With event sourcing you rebuild a read model by replaying the log from the start. Here there is no complete log: the outbox is pruned and topics have retention. Rebuilds instead combine a snapshot with the stream. Record the current log position, copy every row from the write database into a new read table, then start the projector from the recorded position. Changes that land in both the snapshot and the stream are harmless because the version guard discards whichever copy is older. Debezium's initial snapshot followed by streaming is exactly this procedure, automated.
Build the new read model beside the old one, compare row counts and a checksum of (id, version) per key range, switch the query API with a flag, and drop the old table later. Run the same comparison nightly as a reconciliation job even when nothing changed; it catches projector bugs, poison messages that were skipped, and manual edits that should never have happened.
Failure modes
Failure modes, with the symptom you will see first:
- Dual writes. Rare, unreproducible read models that never learned about a change. Fix with an outbox or CDC, not retries.
- Per-key ordering lost. Rows flip back to older states after a consumer scaling event. Key messages by aggregate ID and guard with versions.
- Poison message. One malformed change stops a partition and lag climbs for every order behind it. Retry a bounded number of times, then park the message in a dead-letter topic and alert.
- Schema drift. A renamed column breaks every CDC consumer at once. Version payloads, add fields before using them, and remove fields only after all projectors stop reading them.
- Business rules in projectors. Two read models compute discount eligibility differently. Projectors copy and reshape; decisions belong to the command side.
- Queries on the write side creep back. Reporting joins against the orders table slow command latency. Route them to a read model.
Trade-offs: when not to bother
Level 3 CQRS buys independent scaling, read shapes the write schema could never serve, and freedom to add views without touching the command path. It costs a broker or connector to operate, projectors to monitor, a lag budget to defend, reconciliation jobs and a support burden when users see stale data. For a CRUD service whose reads and writes share a shape, it is pure overhead; level 1 or 2 gets most of the clarity for none of the infrastructure. Reach for level 3 when read traffic dwarfs writes, when a query needs a different engine, or when separate teams own views of the same data.
What to do next
- List your top five queries and their latency and freshness needs, and decide which CQRS level each actually requires.
- Add a version column to each aggregate table and enforce it in every command handler with a conditional update.
- Replace any dual write with an outbox in the same transaction, published by a relay or by CDC on the outbox table.
- Make every message carry the full after-state and version, and write projectors as version-guarded upserts with tombstones for deletes.
- Publish per-projector lag metrics against written budgets, and add a dead-letter path for poison messages.
- Build a nightly reconciliation that compares counts and (id, version) checksums between the write store and each read model, and rehearse a full snapshot-plus-stream rebuild.