CQRS: command query responsibility segregation
Chapter 9 gave order-service an event-sourced write model (the order_events log) plus one derived read projection (order_summary). This chapter generalizes that read side into full CQRS, with multiple purpose-built projections serving different query needs from the same event log.
1. Problem the Pattern Solves
order-service’s order_summary table (Chapter 9) works fine for GET /orders/{id}. But the operations team now wants a dashboard showing orders-per-hour broken down by status and region, refreshed every few seconds, and the fastest way to build it — querying order_summary with a GROUP BY — starts competing for the same table and connection pool that checkout traffic writes to. A slow analytical query during a flash sale measurably slows down POST /orders for real customers, because both are hitting the same rows under the same lock contention, even though “check out” and “show me an hourly rollup” have nothing in common as operations.
At the same time, the mobile team wants an order-history screen with cursor-based pagination and full-text search over item names — a query shape order_summary’s simple current-state schema was never designed for, and adding search indexes and pagination columns to the one table that also serves the hot checkout path risks slowing down writes to satisfy a read need.
Forces in tension:
- Read/write contention vs. a single, simple model. One table serving both writes (order placement) and reads (dashboards, search, mobile history) means every index added for a read pattern has a write-side cost, and every read competes with writes for the same locks and connections.
- Model fit vs. duplication. A read model shaped exactly for its query (denormalized, pre-aggregated, indexed for search) serves that query far more efficiently than a generic normalized table — at the cost of maintaining multiple derived copies of the same underlying facts.
- Consistency vs. read-model freedom. The write model can enforce strict invariants transactionally; a separate read model, updated asynchronously from events, is eventually consistent by construction — a real, sometimes user-visible lag that must be an accepted trade-off, not an oversight.
- Operational complexity vs. targeted performance. Multiple read models mean multiple things to build, deploy, and keep in sync — justified only where a specific read pattern’s performance or shape needs genuinely diverge from what a single model can serve well.
2. Core Idea
CQRS (Command Query Responsibility Segregation) means using separate models for writing data (commands, which change state and enforce invariants) and reading data (queries, which return a view shaped for a specific need) — rather than one model serving both. The write model is the source of truth; read models are derived, purpose-built, and often physically separate stores.
(Full detail — exact endpoints, storage engines, refresh cadence — is in the prose below; the diagram is deliberately just the shape: one write model, three independent read projections.)
Participants: the same order_events write model from Chapter 9, now feeding three independent read projections instead of one, each optimized for its specific query — a point lookup, an aggregated rollup, and a full-text search index — using whichever storage technology fits each query shape best.
Commonly confused with:
- Event sourcing (Chapter 9). CQRS doesn’t require event sourcing — you can split reads and writes over an ordinary mutable-row write model just as validly, updating read tables via database triggers, change-data-capture, or explicit application code after each write. Northwind’s implementation pairs the two because Chapter 9’s event log happens to be an ideal, already-built source for feeding multiple projections — but the pairing is a convenience, not a requirement of either pattern.
- Simply having a read replica of the same database. A Postgres read replica serves the same schema and query shape as the primary, just on different hardware — it solves read scaling, not model-shape mismatch. CQRS’s read models are often differently shaped and sometimes differently stored (Elasticsearch here, for the search projection) — the point is fitting the model to the query, not just adding read capacity.
- Microservices decomposition (Chapter 2). CQRS splits read and write models within one service’s data, not the service itself into two deployables (though a large enough read workload can justify a genuinely separate “query service” deployable — Section 4 discusses when that’s warranted versus overkill).
3. When to Use It
Strong indicators:
- Read and write workloads have measurably different performance profiles or peak times, and are visibly contending for the same resources (Northwind’s dashboard-versus-checkout contention).
- A query need’s natural shape (full-text search, complex aggregation, a different pagination model) doesn’t fit the write model’s schema without compromising the write path’s simplicity or performance.
- Multiple, materially different query patterns exist against the same underlying data (point lookup, aggregation, search) — one query pattern alone rarely justifies the added moving parts.
Concrete use cases:
- E-commerce, as here: checkout (write-heavy, needs strict consistency) versus reporting dashboards and search (read-heavy, tolerant of a few seconds’ staleness) are a textbook CQRS split.
- Logistics: a “current shipment status” write path (frequent, small updates from scanners) versus a “shipment analytics” read path (heavy aggregation across millions of historical shipments) benefit from entirely separate models and often separate storage engines.
- Banking: transaction processing (write model, strict consistency, small individual operations) versus account statement generation and spending-category analytics (read models, aggregation-heavy, can tolerate near-real-time rather than instant consistency).
- SaaS analytics products: a write path recording discrete usage events versus a read path serving pre-aggregated dashboards to customers — nearly always modeled and stored completely differently in practice, even without anyone naming it “CQRS.”
Prerequisites:
- A stable write model (Chapter 9’s event log, or any reliable source of change) to derive read models from — building CQRS on top of a write model that’s still churning means rebuilding projections repeatedly as the source shape changes.
- Monitoring for projection lag, since every read model is now eventually consistent with the write model by some measurable delay that must be visible, not assumed away.
- A clear owner and rebuild procedure per read model — each one needs to be treated as its own small piece of infrastructure, not an afterthought bolted onto the write path.
4. When Not to Use It
- A single query pattern that the write model already serves adequately. If
order_summary’s point lookups are the only read need and they perform fine, splitting further “for CQRS’s sake” adds projections nobody’s actual query pattern requires. - Strong consistency is required between a write and the very next read of it. A checkout confirmation page that must show the just-placed order’s exact state immediately, with zero tolerance for eventual-consistency lag, may need to read from the write model directly (or a synchronously-updated projection within the same transaction) rather than an asynchronously-updated read model.
- Small scale, no contention, no divergent query shapes. A service handling modest traffic with one reasonably-shaped table serving both reads and writes without any measured contention doesn’t need this pattern — it’s solving a problem that hasn’t shown up yet.
- Overengineering signal: standing up a dedicated Elasticsearch cluster for a search feature that a Postgres
LIKEquery or atsvectorfull-text index would serve perfectly well at current data volume — match the projection’s storage technology to the actual, measured query need, not to what looks impressive in an architecture diagram. - Risk: every additional read model is one more thing that can silently fall behind, corrupt, or diverge from the write model — without deliberate lag monitoring and a rebuild-from-source-of-truth capability, CQRS trades a contention problem for a hidden staleness/drift problem that’s easy to miss until a customer notices stale data.
5. Implementation Example
Three projections, three different needs, built from the same order_events log established in Chapter 9.
Projection 1 — order_summary (Postgres, point lookups) — unchanged from Chapter 9; included here as the baseline every other projection is compared against.
Projection 2 — an hourly rollup for the operations dashboard, using a Postgres materialized view refreshed on a schedule rather than per-event, since dashboard freshness of a few seconds is acceptable and per-event refresh would be wasteful for an aggregation this coarse:
CREATE MATERIALIZED VIEW order_dashboard_rollup ASSELECT date_trunc('hour', occurred_at) AS hour_bucket, (payload->>'region') AS region, event_type, count(*) AS event_countFROM order_eventsWHERE event_type IN ('OrderCreated', 'PaymentCaptured', 'OrderCancelled')GROUP BY 1, 2, 3;
CREATE UNIQUE INDEX ON order_dashboard_rollup (hour_bucket, region, event_type);package `in`.o612.eng.northwind.order.internal.cqrs
import org.springframework.jdbc.core.JdbcTemplateimport org.springframework.scheduling.annotation.Scheduledimport org.springframework.stereotype.Component
@Componentclass DashboardRollupRefreshJob(private val jdbc: JdbcTemplate) { @Scheduled(fixedRate = 30_000) fun refresh() { // CONCURRENTLY avoids locking out readers during refresh — requires // the unique index created above. jdbc.execute("REFRESH MATERIALIZED VIEW CONCURRENTLY order_dashboard_rollup") }}Projection 3 — a search index for the mobile order-history screen, using Elasticsearch specifically because full-text search over item names and cursor-based pagination at scale are exactly what Postgres’s general-purpose indexing isn’t optimized for — a deliberate, justified technology choice for this one query shape, not a default:
package `in`.o612.eng.northwind.order.internal.cqrs
import `in`.o612.eng.northwind.order.internal.eventsourcing.OrderEventimport org.springframework.kafka.annotation.KafkaListenerimport org.springframework.stereotype.Component
@Componentclass OrderSearchProjector(private val searchClient: OrderSearchClient) {
@KafkaListener(topics = ["order-service.internal-events"], groupId = "order-search-projector") fun onEvent(event: OrderEvent) { when (event) { is OrderEvent.OrderCreated -> searchClient.index(OrderSearchDocument(event.orderId, event.customerId, event.items.map { it.sku }, "PENDING")) is OrderEvent.PaymentCaptured -> searchClient.updateStatus(event.orderId, "PAID") is OrderEvent.OrderCancelled -> searchClient.updateStatus(event.orderId, "CANCELLED") else -> Unit } }}
data class OrderSearchDocument(val orderId: java.util.UUID, val customerId: java.util.UUID, val skus: List<String>, val status: String)This projector consumes from a dedicated internal Kafka topic (order-service.internal-events) carrying order-service’s own domain events for its own projections — a different, narrower-scoped topic than the cross-service integration events from Chapters 7–8, kept deliberately separate so a change to an internal projection’s needs never has to be negotiated with other services as if it were a public contract.
Query side, dispatching each read to its purpose-built projection:
package `in`.o612.eng.northwind.order.api
import org.springframework.web.bind.annotation.*import java.util.UUID
@RestControllerclass OrderQueryController( private val summaryRepository: OrderSummaryRepository, // projection 1 private val dashboardRepository: DashboardRollupRepository, // projection 2 private val searchClient: OrderSearchClient, // projection 3) { @GetMapping("/api/v1/orders/{id}") fun getOrder(@PathVariable id: UUID) = summaryRepository.findById(id)
@GetMapping("/api/v1/dashboard/orders-per-hour") fun ordersPerHour(@RequestParam region: String) = dashboardRepository.hourlyBreakdown(region)
@GetMapping("/api/v1/orders/search") fun search(@RequestParam q: String, @RequestParam(required = false) cursor: String?) = searchClient.search(q, cursor)}Each endpoint reads from exactly the model built for it — no endpoint runs an aggregation against the point-lookup table, and no endpoint runs a text search against Postgres — the whole point of the split made concrete.
6. Step-by-Step Flow
- Client action.
POST /orders, unchanged since Chapter 1. - API request. Handled entirely by the write side — the event append from Chapter 9 — with zero awareness of any read model.
- Service behavior. The append triggers (via the same
AFTER_COMMITdiscipline established in Chapters 1 and 7) updates to all three read projections, each on its own cadence. - Database interaction.
order_summaryupdates near-immediately (Chapter 9);order_dashboard_rolluprefreshes on its 30-second schedule; the search index updates near-immediately via its own Kafka consumer. - Inter-service communication. None of the read-side updates involve other services — this is entirely internal to
order-service, the same “invisible to callers” property Chapter 3’s database-per-service migration had. - Error or failure handling. If the search projector falls behind or fails,
order_summary’s point lookups are entirely unaffected — the projections are independent, so one failing doesn’t degrade the others, a direct benefit of the separation. - Observability signals. Track projection lag per read model separately —
order_summary’s lag,order_dashboard_rollup’s refresh age, and the search projector’s Kafka consumer lag are three distinct health signals, not one aggregate “is CQRS working” metric. - Final response/outcome. The dashboard shows data up to 30 seconds old (accepted), the search screen shows near-real-time results, and checkout remains fast and uncontended by either — three query needs, three fit-for-purpose answers, from one write model.
7. Production Concerns
- Timeouts, retries, idempotency. Each projector must handle redelivered events idempotently, exactly as Chapter 7 established — the search projector’s
index()call should be an upsert keyed byorderId, safe to run twice for the same event. - Data consistency and transaction boundaries. Every read model here is eventually consistent with the write model, on its own schedule — document each projection’s expected staleness (near-real-time for
order_summaryand search; up to 30 seconds for the dashboard rollup) as part of its contract with whoever consumes it, not as an implicit assumption. - Read-model rebuild. Every projection must be fully rebuildable from
order_eventsalone — if the search index is ever corrupted or needs a schema change, the recovery procedure is “replay the log into a fresh index,” not “attempt to patch the broken one.” Verify this rebuild path actually works, in a drill, before you need it in an incident. - API versioning. Each read endpoint’s response shape can evolve independently of the others and independently of the write model’s event schema — a genuine benefit, since a dashboard’s aggregation shape changing doesn’t touch the search API’s contract at all.
- Authentication and authorization. Different read models may warrant different access policies — the dashboard rollup might be internal-only (ops team), while order search is customer-facing (scoped to that customer’s own orders) — enforce this at the query API layer per endpoint, not assumed uniform across all reads.
- Logging, metrics, tracing. Correlate a write’s correlation ID through to each projection update, so “why isn’t my search result showing up yet” is answerable by tracing the specific event through the specific projector, not by guessing.
- Kubernetes deployment, autoscaling. Read-heavy projections (especially the search index, if query volume grows significantly) can eventually justify becoming their own deployable with its own scaling profile — evaluate this the same way Chapter 2 evaluated service extraction: on measured evidence, not by default.
- Testing strategy. Test each projector in isolation, feeding it a sequence of events and asserting on the resulting read-model state — these tests don’t need the write side running at all, a direct benefit of the clean separation.
- Migration strategy. Add projections incrementally, one query need at a time, starting from the pain that’s actually being felt (Northwind started with the dashboard rollup, since that was the contention that paged someone) — building all three from day one, speculatively, is the overengineering Section 4 warns against.
8. Common Mistakes
- Building a read model before a specific query need justifies it. Standing up the Elasticsearch search projection before the mobile order-history feature actually needed full-text search means maintaining infrastructure with no consumer. Fix: build each projection in response to a real, current query requirement, not anticipated future ones.
- Letting a read model silently fall behind with no lag monitoring. Without visibility into projection lag, a stuck consumer can serve stale data for hours before anyone notices — often a customer, not an engineer. Fix: monitor lag per projection explicitly, as its own metric, with alerting thresholds appropriate to that projection’s staleness contract.
- No rebuild procedure, or an untested one. Discovering during an actual incident that “replay events into a fresh index” doesn’t actually work as described is far worse than discovering it in a drill. Fix: periodically exercise the full rebuild path for every read model, not just document it.
- Reading from the write model directly when a projection would serve better. Querying
order_eventsdirectly for a dashboard need, “just this once, since the projection isn’t ready yet,” reintroduces the exact contention this pattern was adopted to remove. Fix: if a query need exists, build its projection before serving that traffic from the write side. - Treating every projection’s staleness the same way. Applying the dashboard rollup’s 30-second acceptable staleness to a payment-confirmation read (which may need much fresher, or even synchronous, data) risks showing a customer stale, wrong information about their own transaction. Fix: set staleness tolerance per read need, based on that specific consumer’s actual requirement, not a blanket policy.
- Choosing a projection’s storage technology by fashion rather than fit. Adopting Elasticsearch for the dashboard rollup (a simple aggregation Postgres materialized views handle well) just because it was already in use for search elsewhere adds an unnecessary second technology to operate. Fix: match each projection’s storage to its actual query shape, as Section 5 does — Postgres for aggregation and point lookups, Elasticsearch specifically for full-text search.
9. Decision Guide
| Problem signal | Use this pattern? | Why | Alternative |
|---|---|---|---|
| Read and write workloads measurably contend for the same resources | Yes | Separate models remove the contention at its source | — |
| A query shape (search, complex aggregation) doesn’t fit the write model well | Yes | A purpose-built projection serves that shape far more efficiently | — |
| Single query pattern, no contention, current model serves it fine | No | Nothing to solve yet; added projections would be speculative | Keep the single model |
| Consumer needs strongly consistent reads immediately after a write | No, for that specific read | Eventually-consistent projections may lag exactly when consistency matters most | Read from the write model directly, or a synchronously-updated projection |
| Multiple divergent query needs already causing pain (contention, poor fit, or both) | Yes | This is the strongest, clearest case for the pattern | — |
10. Hands-On Exercise
Extend it: add a fourth projection answering “average time from OrderCreated to PaymentCaptured, by hour” for a latency-monitoring dashboard, choosing and justifying its storage technology and refresh cadence using Section 5’s reasoning as a template.
Simulate a failure: stop the search projector’s Kafka consumer for ten minutes while orders continue to be placed, then restart it. Confirm the search index catches up correctly once it resumes, and measure how far behind it fell — this is the projection-lag monitoring from Section 7 made concrete.
Decision question, with justification required: the operations dashboard’s product owner now wants near-real-time (sub-5-second) freshness instead of the current 30-second materialized-view refresh. Should you decrease the refresh interval, switch the rollup to an event-driven update (like the search projector) instead of a scheduled one, or push back on the requirement? Weigh the trade-offs from Section 1 and Section 7.
11. Key Takeaways
- CQRS separates the model used to write data (enforcing invariants, the source of truth) from the model(s) used to read it (shaped for specific query needs) — removing contention and letting each side be optimized independently.
- It doesn’t require event sourcing, though the two pair naturally when a reliable event log already exists to derive projections from, as Chapter 9’s did here.
- Justify each read model by a real, current query need with a measured performance or shape problem — building projections speculatively is the pattern’s most common overengineering trap.
- Every read model is eventually consistent with the write model on its own schedule — document and monitor that staleness explicitly per projection, rather than treating “CQRS is working” as one binary health signal.
- Every projection must be rebuildable from the write model alone, and that rebuild path should be tested before an incident forces you to rely on it.
- Match each projection’s storage technology to its actual query shape — Postgres for aggregation and point lookups, a search engine specifically for full-text search, and so on — not to whichever technology is already fashionable on the team.
- CQRS adds real operational surface (more projections, more lag to monitor, more rebuild procedures to maintain) — it pays for itself only where divergent query needs are already causing measurable pain.