Hello! Let's get started with our next lesson.
Introduction
In our previous lesson, we established the conceptual foundation of Event Sourcing, defining key components like events, streams, aggregates, projections, and snapshots. We understood that instead of storing an entity's current state, we persist an immutable sequence of events, and the state is derived by replaying this history.
Today, we transition from the "what" to the "how." The learning outcome for this lesson is to implement the persistence of an aggregate's state as an event sequence in an event store. We will dissect the mechanics of the write model in an event-sourced system: loading an aggregate from its history, executing a command that produces new events, and atomically persisting those new events.
Given your background in building low-latency and high-load systems, we will focus on the critical implementation details that ensure data consistency and performance, such as optimistic concurrency control and snapshotting strategies.
1. The Core Components: Aggregate and Repository
At the heart of an event-sourced implementation are two key software components: the Aggregate and the Repository. The aggregate encapsulates the business logic and state, while the repository orchestrates the interaction with the persistence layer (the event store).
A well-structured implementation clearly separates these responsibilities. Let's examine a practical, code-based guide that walks through building these components. Although the examples are in C#, the object-oriented principles and patterns are directly applicable to Java and other languages.
Building Your First Event Store from Scratch: A Developer's ...
The article 'Building Your First Event Store from Scratch' provides an excellent step-by-step implementation of the core application-level components for Event Sourcing. We will focus on the Aggregate and Repository patterns it describes.
Please read 'Step 3: Creating the Product Aggregate' and 'Step 4: Building the Repository Layer'. As you read, focus on the following design decisions: In the Aggregate: The two distinct constructors: one for creating a new aggregate (which generates the first event) and one for rehydrating an existing aggregate from a history of events. The Apply method pattern: How state is mutated only in response to an event, ensuring the event is the single source of truth for state changes. The tracking of _uncommittedEvents: This internal list holds new events until they are successfully persisted by the repository. In the Repository: The GetByIdAsync method: How it retrieves events from the store and uses the rehydration constructor to build the aggregate. The SaveAsync method: The orchestration logic for getting uncommitted events, calculating the expectedVersion for optimistic concurrency, saving the events, and finally clearing the aggregate's list of uncommitted events.
This separation of concerns is crucial. The aggregate knows nothing about persistence; it only knows how to validate commands and produce events. The repository handles the entire persistence lifecycle, acting as the bridge to the event store.
2. Persisting Events in a Relational Database
Now that we understand the application logic, let's look at the database side. While purpose-built event stores like EventStoreDB or Axon Server exist, a standard relational database like PostgreSQL can serve as a robust and transactionally-safe event store. This approach is common in environments where leveraging existing database infrastructure is preferred.
The following resource provides a complete reference implementation using PostgreSQL.
eugene-khyst/postgresql-event-sourcing
This GitHub repository, 'postgresql-event-sourcing', demonstrates how to implement an event store using PostgreSQL. We will focus on the database schema and the mechanism for appending events.
Please read the sections 'Event sourcing and CQRS basics' and 'Solution architecture'. Pay close attention to: ER Diagram: The structure of the ES_EVENT table. Note the columns for aggregate_id, version, event_type, and the json_data payload. Optimistic Concurrency Control: The role of the ES_AGGREGATE table, which stores the latest version of each aggregate. Understand the two-step process within a single transaction for appending an event: First, UPDATE the version in the ES_AGGREGATE table, checking that the current version matches the expected version. Second, INSERT the new event into the ES_EVENT table.
This transactional approach is fundamental. If two concurrent requests try to update the same aggregate, the first one will succeed in updating the version number. The second request's UPDATE statement will fail because its expectedVersion no longer matches the version in the database, preventing a lost update and ensuring a consistent, linear history. This is a practical application of optimistic locking, a pattern you've likely encountered in trading and payment systems.
3. Performance Optimization with Snapshots
As we discussed in the last lesson, rehydrating an aggregate with a very long history can become a performance bottleneck. Snapshots are the standard solution. Let's look at how they are implemented.
The core idea is to periodically save the full state of the aggregate at a specific version. To rehydrate, you load the latest snapshot and then replay only the events that occurred after that snapshot.
Event Sourcing: Rehydrating Aggregates with Snapshots
The video 'Event Sourcing: Rehydrating Aggregates with Snapshots' by Derek Comartin (CodeOpinion) provides a clear, code-driven explanation of implementing snapshots.
Please watch from the beginning until 10:20. The video covers: The problem statement: why replaying thousands of events is inefficient (00:00 - 01:30). The concept: using a separate snapshot stream to store the aggregate's state at a specific version (01:30 - 04:21). The implementation of reading: how to load the latest snapshot first, then query for subsequent events (04:21 - 09:20). The implementation of writing: how to decide when to create a snapshot and append it to the snapshot stream (09:20 - 10:20).
The postgresql-event-sourcing repository (resource 8dfd0) also shows a concrete implementation of this. It uses an ES_AGGREGATE_SNAPSHOT table and provides the SQL query for loading the latest snapshot before a specified version, which perfectly complements the logic shown in the video.
The key takeaway is that snapshotting is an optimization. It adds complexity, so it should only be introduced when you have evidence (through measurement) that aggregate rehydration is becoming a performance issue.
4. The End-to-End Write Flow
Let's consolidate what we've learned into the complete, end-to-end sequence of operations for handling a command on the write side:
- Receive Command: An application service receives a command intended for a specific aggregate (e.g.,
UpdatePriceCommandforproduct-123). - Load Aggregate (via Repository):
- The repository first queries the snapshot store for the latest snapshot of
product-123. - It then queries the event store for all events for
product-123that occurred after the snapshot's version. If no snapshot exists, it loads all events from the beginning. - The repository instantiates the
Productaggregate, passing it the snapshot state (if any) and the subsequent event stream.
- The repository first queries the snapshot store for the latest snapshot of
- Rehydrate State: The aggregate's constructor and
Applymethods process the events, rebuilding its internal state to the present moment. - Execute Command: The service calls the appropriate business method on the now fully-hydrated aggregate (e.g.,
product.UpdatePrice(newPrice)). - Produce New Events: The aggregate's business logic validates the command. If valid, it creates one or more new
DomainEventobjects (e.g.,ProductPriceUpdatedEvent) and adds them to its internal list of uncommitted events. - Save Aggregate (via Repository):
- The service calls
repository.Save(product). - The repository retrieves the uncommitted events from the aggregate.
- It begins a database transaction.
- It performs an optimistic lock check by attempting to
UPDATEthe aggregate's version in theES_AGGREGATEtable, ensuring the version it loaded matches the current database version. - It
INSERTs the new events into theES_EVENTtable, each with the correct, incremented version number. - If the save logic determines a snapshot is due (e.g., every 100th event), it also
INSERTs a new snapshot of the aggregate's current state into theES_AGGREGATE_SNAPSHOTtable. - The transaction is committed.
- The service calls
- Finalize: Upon successful persistence, the repository calls a method on the aggregate (e.g.,
MarkEventsAsCommitted()) to clear its internal list of uncommitted events.
This flow ensures that business rules are enforced on a consistent state, and that all changes are atomically and durably recorded as an immutable history.
Conclusion
In this lesson, we moved from the theory of Event Sourcing to its practical implementation. We've seen how the Aggregate and Repository patterns work together to manage state and persistence, how a relational database can be used as a robust event store, and how snapshotting is used to optimize performance for long-lived aggregates.
Key Takeaways:
- The Aggregate is the core of the domain model, enforcing business rules and producing events. Its state is always a derivation of its event history.
- The Repository orchestrates persistence, handling the loading of historical events (and snapshots) and the atomic saving of new events.
- Optimistic Concurrency Control is essential for consistency in a concurrent environment and is typically implemented by versioning the aggregate's event stream.
- Snapshots are a crucial performance optimization, implemented by storing the aggregate's state at a point in time to reduce the number of events that need to be replayed.
Preview of the Next Lesson
We have now built a solid understanding of the write side of an event-sourced system. However, this write model is not optimized for querying. In our next lesson, we will address this by focusing on the other half of the architecture: "Design a CQRS-based architecture for a given workload, implementing separate read and write models." We will explore how to build and maintain projections (read models) by consuming the event stream we've just learned how to create.