Event Sourcing
Event Sourcing is an architectural pattern where changes to application state are stored as a sequence of events. Instead of storing the current state of an entity in a table, you store the entire history of what happened to it.
Overview
In a standard database, if a user updates their email, you UPDATE users SET email = 'new' WHERE id = 1. The old email is lost forever.
In Event Sourcing, you never update records. Instead, you append a new event: { type: 'UserEmailUpdated', newEmail: 'new', timestamp: '...' }. The current state of the user is derived by playing all these events in order from the beginning of time.
This pattern is natively supported in NestJS via the @nestjs/cqrs package.
Key Concepts
- Event Store: An append-only database (like EventStoreDB, Apache Kafka, or a specialized SQL table) that stores the sequence of events.
- Aggregate Root: The domain entity that reconstructs its state by applying past events to itself.
- Projections: Because querying a list of users by reading a million events is incredibly slow, you create “Read Models” (Projections). As events are saved, an Event Handler updates a standard SQL table specifically optimized for reading.
Code Examples
1. The Event
An event is a fact that happened in the past.
// user-created.event.ts
export class UserCreatedEvent {
constructor(
public readonly aggregateId: string,
public readonly email: string,
) {}
}
// user-email-updated.event.ts
export class UserEmailUpdatedEvent {
constructor(
public readonly aggregateId: string,
public readonly newEmail: string,
) {}
}
2. The Aggregate Root
The Aggregate Root inherits from AggregateRoot provided by @nestjs/cqrs. It handles its own state reconstruction.
// user.aggregate.ts
import { AggregateRoot } from '@nestjs/cqrs';
import { UserCreatedEvent } from './user-created.event';
import { UserEmailUpdatedEvent } from './user-email-updated.event';
export class UserAggregate extends AggregateRoot {
public id: string;
public email: string;
constructor(id?: string) {
super();
this.autoCommit = false; // We commit manually after saving to the Event Store
}
// Business Logic: Creating a new user
createUser(id: string, email: string) {
// 1. We don't set this.email directly. We apply an event!
this.apply(new UserCreatedEvent(id, email));
}
// Business Logic: Updating email
updateEmail(newEmail: string) {
if (this.email === newEmail) throw new Error('Same email');
this.apply(new UserEmailUpdatedEvent(this.id, newEmail));
}
// ----------------------------------------------------
// STATE RECONSTRUCTION (Called automatically by apply())
// ----------------------------------------------------
onUserCreatedEvent(event: UserCreatedEvent) {
this.id = event.aggregateId;
this.email = event.email;
}
onUserEmailUpdatedEvent(event: UserEmailUpdatedEvent) {
this.email = event.newEmail;
}
}
3. The Command Handler (Writing)
The Command Handler creates the aggregate, calls the business logic, saves the uncommitted events to the Event Store, and then publishes them.
// update-email.handler.ts
import { CommandHandler, ICommandHandler, EventPublisher } from '@nestjs/cqrs';
@CommandHandler(UpdateEmailCommand)
export class UpdateEmailHandler implements ICommandHandler<UpdateEmailCommand> {
constructor(
private readonly eventStore: EventStoreService,
private readonly publisher: EventPublisher,
) {}
async execute(command: UpdateEmailCommand) {
// 1. Fetch all past events for this user from the DB
const eventHistory = await this.eventStore.getEventsForAggregate(command.userId);
// 2. Reconstruct the Aggregate state
const user = this.publisher.mergeObjectContext(new UserAggregate());
user.loadFromHistory(eventHistory);
// 3. Execute business logic (which generates a new event)
user.updateEmail(command.newEmail);
// 4. Save the NEW events to the database
const uncommittedEvents = user.getUncommittedEvents();
await this.eventStore.saveEvents(command.userId, uncommittedEvents);
// 5. Publish the events to the event bus (for projections to read)
user.commit();
}
}
Best Practices
- Never Change Past Events: Events are immutable facts. If a bug caused the wrong email to be saved, you don’t edit the old
UserEmailUpdatedEventin the database. You issue a newUserEmailCorrectedEvent. This creates a perfect, auditable ledger (essential in banking/finance). - Snapshots: If a user has 10,000 events, pulling them all from the database to reconstruct the Aggregate will be slow. Take a “Snapshot” of the aggregate state every 100 events. Then, you only load the snapshot + the events that occurred after the snapshot.