Hello! Welcome back.
In our last lesson, we focused on modeling a business process by defining its event sequence, producers, and consumers. This gave us a clear, business-grounded blueprint for an event stream—the "source of truth" in an event-driven system.
Today, we'll address the next logical question: how do we use this event stream to answer queries efficiently?
Lesson 2: Implementing a Projection to Create a Read Model
Learning Outcome: By the end of this lesson, you will be able to implement a projection to create a read model from an event stream.
This is a cornerstone technique in both Command Query Responsibility Segregation (CQRS) and Event Sourcing. The ability to create specialized, optimized read models from a single, canonical event stream is what gives these architectural styles their flexibility and performance, especially in high-load environments.
1. What is a Projection?
In an event-driven architecture, the write side is optimized for transactional consistency and enforcing business rules. Directly querying the write model's data structures (or the event log itself) is often inefficient for complex read requirements.
This is where projections come in. A projection is a consumer that processes an event stream to build and maintain a dedicated read model. This read model is a stateful representation of the data, specifically tailored for the query needs of a client, such as a UI, a reporting tool, or another service.
A single event stream can feed multiple, independent projections. For example, a stream of Order events could feed:
- A projection that creates a denormalized view of recent orders for a customer's dashboard.
- A projection that aggregates sales data for an analytics dashboard.
- A projection that updates a search index.
This architectural separation is central to CQRS, as illustrated below.
This diagram shows the CQRS pattern. Commands update the domain model, which generates events stored in an Event Store. In parallel, Event Handlers (the projection mechanism) consume these events to build and maintain a separate Query Database (the read model), which is then used to answer queries.
The "Event Handlers" in this diagram are the projections. They listen to the event stream and transform it into a queryable format.
2. The Core Logic: Building a View from Events
The fundamental mechanism of a projection is identical to how an event-sourced aggregate reconstructs its state: it starts with an initial state and applies each event from the stream in order.
Let's examine a practical implementation of this concept, which the author calls a "live projection."
Live projections for read models with Event Sourcing and CQRS
The article 'Live projections for read models with Event Sourcing and CQRS' by Anton Stöckl provides a clear, code-driven example of this process. First, let's see how a write-model aggregate is built from events.
Please read the section 'Some sample code for a customer aggregate: the write model'. Focus on the buildCurrentStateFrom function. This demonstrates the baseline pattern of replaying events to construct a state object.
The buildCurrentStateFrom function iterates through an event stream and uses a switch statement to apply each event to the currentState struct. This rehydrates the aggregate so it can process a new command.
Now, let's see how this exact same pattern is used to create a read model.
Live projections for read models with Event Sourcing and CQRS
The next section of the article applies the same logic to build a read model, or 'View'.
Now, read the section 'Sample code for a Customer View: the read model' and study the BuildViewFrom function.
Notice the key differences and similarities:
- Shared Source: Both the write model (
currentState) and the read model (View) are derived from the exact same event stream. - Different Purpose: The
Viewstruct is shaped for a different purpose than thecurrentStatestruct. For example, it contains a simple booleanIsEmailAddressConfirmed, which is more convenient for a UI than the type-based logic used in the write model. TheViewalso flattens thePersonNamevalue object intoGivenNameandFamilyNamefields. - Identical Logic: The core implementation—a loop iterating over events and applying them via a
switch—is the same. This is the fundamental logic of a projection.
This "live projection" approach, where the read model is built on-demand for each query, is one of several architectural choices, each with distinct trade-offs.
3. Architectural Patterns for Projections
While building a projection on-demand is simple and ensures strong consistency, it may not be performant for aggregates with long event histories. In practice, projections are often materialized and updated asynchronously.
Let's explore the spectrum of implementation patterns, which offer different trade-offs between consistency, scalability, and complexity.
Event Sourcing: Benefits, Concepts, and Practical ...
The article 'Event Sourcing: Benefits, Concepts, and Practical...' provides an excellent breakdown of three common architectural patterns for integrating projections in a CQRS/ES system.
Please read the sections 'ES & CQRS Level — 1', 'ES & CQRS Level — 2', and 'ES & CQRS Level — 3'. Pay attention to the descriptions and the accompanying Java code examples for each level.
Let's summarize the trade-offs of these patterns:
| Pattern | Description | Consistency | Performance & Scalability | Use Case |
|---|---|---|---|---|
| Level 1: Same-Transaction Update | The read model is updated within the same database transaction as the event is persisted. | Strong | Low. The write transaction is burdened with extra work, coupling the read and write models. | Simple applications where consistency is paramount and write throughput is not the primary concern. |
| Level 2: Asynchronous Projector | A separate process (the projector) polls or subscribes to the event store and updates the read model asynchronously. | Eventual | High. Decouples the write and read paths. Write operations are fast. The read model can be scaled independently. | The most common pattern for scalable distributed systems. Suitable when a small consistency lag is acceptable. |
| Level 3: Event Broker + Projector | Events are published to a message broker (e.g., Kafka). Projections are independent consumers of the broker's topics. | Eventual | Very High. Maximum decoupling. Enables easy addition of new projections and integration with other systems. | Complex, large-scale systems where multiple, diverse consumers need to react to events. Your experience with high-load payment platforms likely involved this pattern. |
The "live projection" we saw earlier is a fourth option. It provides strong consistency for reads (the article calls it "immediately consistent," as you always read from the latest events) but can have high read latency. This is often mitigated by using snapshots, where the state is periodically saved and the projection only needs to replay events since the last snapshot.
4. Practical Application
Let's connect this back to the "Termination of Employment" process from our previous lesson. The event stream might include ResignationSubmitted, ManagerApproved, EquipmentReturned, and FinalPayrollProcessed.
Imagine you need to build a dashboard for HR managers showing the status of all employees currently in the off-boarding process.
- Read Model Design: The read model for this dashboard might be a simple document or row in a database with fields like:
EmployeeID,EmployeeName,ResignationDate,LastDayOfWork,EquipmentStatus(e.g., "Pending", "Returned"),FinalPayStatus(e.g., "Calculated", "Processed"). - Projection Logic: The projection would subscribe to the event stream.
- On
ResignationSubmitted, it creates a new record for the employee. - On
EquipmentReturned, it updates theEquipmentStatusfield for that employee's record. - On
FinalPayrollProcessed, it updates theFinalPayStatus.
- On
- Architectural Choice: For an internal HR dashboard, immediate consistency is likely not a strict requirement. A few seconds of lag is acceptable. Therefore, an Asynchronous Projector (Level 2 or 3) would be an excellent choice. It decouples the core HR transaction system from the reporting dashboard, improving the scalability and resilience of both.
Conclusion
In this lesson, we demystified the process of creating read models from an event stream.
Key Takeaways:
- Projections are components that consume an event stream to create and maintain specialized read models.
- The core logic of a projection involves iterating through events and applying them to a stateful view, tailoring its structure for specific query needs.
- There are several architectural patterns for implementing projections, primarily trading off consistency vs. scalability. These range from synchronous, in-transaction updates to fully decoupled, asynchronous consumers on a message broker.
- This technique is fundamental to CQRS and Event Sourcing, enabling systems to have both a transactionally consistent write side and multiple, independently scalable, and highly optimized read sides.
Preview of the Next Lesson:
We now know how to design and implement a projection. But what happens when business requirements change and we need to alter the read model's schema? Or what if we want to introduce an entirely new projection for a system that already has years of event history? Our next lesson, "Implement a mechanism to rebuild a projection from an event store," will tackle this critical operational challenge.