Kafka and CQRS: Event-Driven Projections and Separated Read Models
CQRS (Command Query Responsibility Segregation) separates the write model from the read model. The write model handles commands and enforces business rules; the read model handles queries and is optimized for fast reads. Kafka supports CQRS by acting as the event backbone between the two. When the write model processes a command, it emits an event to Kafka. The read model consumes the event and updates a projection: a denormalized view of the data that is optimized for the queries the application needs. The read model can be a database, a search index, a cache, or any other store. Because the read model is updated asynchronously, it is eventually consistent with the write model, but it can be scaled independently and optimized for different query patterns. Kafka provides durability, ordering per key, and replayability, which are essential for building reliable projections. The trade-off is between consistency and scalability. CQRS gives you independent scaling of reads and writes, but it introduces eventual consistency, which the application must handle.
The mechanism of CQRS with Kafka has three parts: the command side, the event log, and the query side. The command side is the write model: it validates commands, applies business rules, and writes to its database. It also emits events to Kafka, usually using the outbox pattern to ensure consistency between the database and Kafka. The event log is Kafka: it stores the events in order per key and retains them for replay. The query side is one or more projections: each projection consumes the events and updates its read model. A projection is essentially a Kafka consumer that applies events to a store. Because Kafka retains events, you can rebuild a projection from scratch by replaying the topic. This is a powerful capability: if a projection has a bug or the read model is corrupted, you can fix the code and replay the events to rebuild the projection. The trade-off is between the number of projections and the operational cost. Each projection is a consumer group with its own offsets and its own read model; more projections mean more consumers, more storage, and more operational overhead. But each projection can be optimized for a specific query pattern, which is the main benefit of CQRS. Version note: Kafka Streams is often used to build projections because it provides state stores, windowing, and exactly-once processing. The Kafka Streams DSL makes it easy to build a projection that consumes events and updates a state store, which can be queried via interactive queries. This is a common pattern for CQRS with Kafka.
A common mistake is to use CQRS for every system, even when the read and write patterns are similar. CQRS adds complexity, so it should be used when there is a clear benefit: when reads and writes have different scaling needs, when the read model needs a different schema or store, or when the read model needs to be rebuilt independently. Another mistake is to forget that the read model is eventually consistent; the application must handle the case where a query returns stale data. A third mistake is to build a projection that is not idempotent; because Kafka delivers at least once, the projection must handle duplicates. The trade-off is between consistency and availability. CQRS with Kafka gives high availability and scalability but eventual consistency; a traditional single-model system gives strong consistency but limited scalability. For systems with high read volume or complex queries, CQRS is often the right choice. Version note: Kafka Streams interactive queries allow you to query a state store directly from the application, which is a form of CQRS where the read model is embedded in the stream processing application. This is useful for low-latency queries but requires that the state store is co-located with the application. For global queries, a separate read model (e.g., Elasticsearch, PostgreSQL) is more appropriate.
CQRS separates the write model from the read model; Kafka is the event backbone.
The write model emits events; the read model consumes them and updates a projection.
Projections are eventually consistent and must be idempotent.
Kafka retention allows projections to be rebuilt by replaying events.
Kafka Streams is a common tool for building projections with state stores and interactive queries.
Use CQRS when reads and writes have different scaling or schema needs, not for every system.
Reset consumer group offsets to rebuild a projection from scratch.
0-2 years experience
2-5 years experience
5-8 years experience