Command Palette

Search for a command to run...

Hectal
PHASE 15Advanced ~8 min· topic 4 of 5

Topic 15.4

CQRS and Read Models

In one line

CQRS separates the write model (validates commands, stores truth) from read models optimised for queries (search indexes, denormalised views, caches), kept up to date by consuming change events from Kafka. Read models are eventually consistent and can be rebuilt from events at any time.

0/5 · 0%

Think of it like this

A library's acquisition ledger (write model) and its public catalogue (read model). New books are recorded in the ledger first; the catalogue is updated from it and can be reprinted from the ledger if the catalogue is lost.

Key ideas

  1. 01

    Write side: a normalised database with invariants and transactions. Every change produces an event (outbox or CDC) into Kafka.

  2. 02

    Read side: consumers project events into stores shaped for queries: OpenSearch for search, a denormalised table for order history, Redis for hot lookups, a warehouse for analytics.

  3. 03

    Eventual consistency: a user may not see their change in the read model for milliseconds to seconds. Handle with read-your-writes tricks (return the new state from the command, or read from the write side briefly after writing).

  4. 04

    Rebuilds: to add a field or fix a bug, create a new index, replay events into it from the beginning (or from a snapshot plus events), and switch reads via an alias. This requires events retained long enough or a snapshot source.

  5. 05

    Idempotent, versioned projections: upsert by entity ID, ignore events older than the stored version, so replays and out-of-order retries are safe.

Code & diagrams

cqrs.mermaiddiagram
Rendering diagram…
SearchProjector.javajava
@KafkaListener(topics = "product.events", groupId = "search-projector-v7")
public void project(ProductUpdated e) {
    // Versioned upsert: stale or duplicate events can't overwrite newer data.
    IndexRequest<ProductDoc> req = IndexRequest.of(r -> r
        .index("products-v7")
        .id(e.productId())
        .version(e.version())
        .versionType(VersionType.External)       // rejected if version <= stored version
        .document(ProductDoc.from(e)));
    try {
        search.index(req);
    } catch (ElasticsearchException ex) {
        if (!isVersionConflict(ex)) throw ex;    // version conflict = already newer: ignore
    }
}

Interview problem

The problem

Search read model for products

Products live in PostgreSQL; ProductUpdated events flow through Kafka; Elasticsearch/OpenSearch serves search. Design synchronisation and a rebuild strategy.

Explain it without notes

01

Why use a compacted topic as the source for rebuilding read models?

Practice

01

Rebuild a read model into a new index while the old one keeps serving, and switch with an alias.

Trade-offs

  • ↔

    CQRS optimises reads independently and allows rebuilds, at the cost of eventual consistency and more moving parts.

Done when you can

  • I can build and rebuild CQRS read models from Kafka events safely.