Patrones de Event Sourcing para Aplicaciones con Auditoría Exigente
Usa event sourcing para un registro de auditoría completo e inmutable de cada cambio, con proyecciones, snapshots y reproducción de eventos.

Por qué el estado por sí solo no basta
La mayoría de aplicaciones guardan el estado actual de las entidades. El saldo de un usuario es de $500. El estado de un pedido es "shipped". Pero cuando un auditor pregunta "¿cómo llegó este saldo a $500?" o "¿quién cambió el estado del pedido y cuándo?", terminas deduciéndolo a partir de logs incompletos. Event sourcing resuelve esto almacenando cada cambio de estado como un evento inmutable, haciendo del historial completo la fuente de la verdad.
La base del event store
Un event store es un log de solo apéndice. Los eventos nunca se actualizan ni eliminan. Para conocer el estado actual de una entidad, reproduces sus eventos desde el inicio.
interface DomainEvent {
eventId: string;
aggregateId: string;
aggregateType: string;
eventType: string;
payload: Record<string, unknown>;
metadata: {
timestamp: string;
userId: string;
correlationId: string;
causationId: string;
};
version: number;
}
class EventStore {
constructor(private readonly db: Database) {}
async append(
aggregateId: string,
events: DomainEvent[],
expectedVersion: number
): Promise<void> {
await this.db.transaction(async (tx) => {
// Optimistic concurrency check
const current = await tx.query<{ max_version: number }>(
"SELECT COALESCE(MAX(version), 0) as max_version FROM events WHERE aggregate_id = $1",
[aggregateId]
);
if (current.rows[0].max_version !== expectedVersion) {
throw new ConcurrencyError(
`Expected version ${expectedVersion}, ` +
`but found ${current.rows[0].max_version}`
);
}
for (const event of events) {
await tx.query(
`INSERT INTO events (event_id, aggregate_id, aggregate_type,
event_type, payload, metadata, version)
VALUES ($1, $2, $3, $4, $5, $6, $7)`,
[
event.eventId,
event.aggregateId,
event.aggregateType,
event.eventType,
JSON.stringify(event.payload),
JSON.stringify(event.metadata),
event.version,
]
);
}
});
}
async getEvents(
aggregateId: string,
fromVersion?: number
): Promise<DomainEvent[]> {
const result = await this.db.query<DomainEvent>(
`SELECT * FROM events
WHERE aggregate_id = $1 AND version > $2
ORDER BY version ASC`,
[aggregateId, fromVersion ?? 0]
);
return result.rows;
}
}Agregados que producen eventos
Los agregados son el modelo de escritura en event sourcing. Llegan comandos, las reglas de negocio los validan y el agregado emite eventos que describen qué ocurrió, no lo que se solicitó.
// ❌ Mutable state with no history
class Account {
balance: number = 0;
withdraw(amount: number): void {
this.balance -= amount; // History is lost
}
}
// ✅ Event-sourced aggregate with full audit trail
class AccountAggregate {
private balance: number = 0;
private status: "active" | "frozen" = "active";
private uncommittedEvents: DomainEvent[] = [];
private version: number = 0;
static fromEvents(events: DomainEvent[]): AccountAggregate {
const account = new AccountAggregate();
for (const event of events) {
account.apply(event);
account.version = event.version;
}
return account;
}
withdraw(amount: number, userId: string, correlationId: string): void {
if (this.status === "frozen") {
throw new BusinessRuleError("Cannot withdraw from a frozen account");
}
if (amount <= 0) {
throw new BusinessRuleError("Withdrawal amount must be positive");
}
if (this.balance < amount) {
throw new BusinessRuleError("Insufficient funds");
}
this.emit({
eventType: "MoneyWithdrawn",
payload: { amount, previousBalance: this.balance },
userId,
correlationId,
});
}
private apply(event: DomainEvent): void {
switch (event.eventType) {
case "AccountOpened":
this.balance = event.payload.initialDeposit as number;
this.status = "active";
break;
case "MoneyDeposited":
this.balance += event.payload.amount as number;
break;
case "MoneyWithdrawn":
this.balance -= event.payload.amount as number;
break;
case "AccountFrozen":
this.status = "frozen";
break;
}
}
getUncommittedEvents(): DomainEvent[] {
return [...this.uncommittedEvents];
}
}Proyecciones para modelos de lectura
Event sourcing separa escrituras de lecturas. El event store gestiona las escrituras. Las proyecciones se suscriben a los eventos y construyen modelos de lectura optimizados para las consultas. Es el patrón CQRS aplicado de forma natural.
interface Projection {
name: string;
handle(event: DomainEvent): Promise<void>;
}
class AccountBalanceProjection implements Projection {
name = "account-balance";
constructor(private readonly readDb: Database) {}
async handle(event: DomainEvent): Promise<void> {
switch (event.eventType) {
case "AccountOpened":
await this.readDb.query(
`INSERT INTO account_balances (account_id, balance, owner_name, updated_at)
VALUES ($1, $2, $3, $4)`,
[
event.aggregateId,
event.payload.initialDeposit,
event.payload.ownerName,
event.metadata.timestamp,
]
);
break;
case "MoneyDeposited":
case "MoneyWithdrawn": {
const delta =
event.eventType === "MoneyDeposited"
? (event.payload.amount as number)
: -(event.payload.amount as number);
await this.readDb.query(
`UPDATE account_balances
SET balance = balance + $1, updated_at = $2
WHERE account_id = $3`,
[delta, event.metadata.timestamp, event.aggregateId]
);
break;
}
}
}
}
class AuditLogProjection implements Projection {
name = "audit-log";
constructor(private readonly readDb: Database) {}
async handle(event: DomainEvent): Promise<void> {
await this.readDb.query(
`INSERT INTO audit_log (event_id, aggregate_id, event_type,
user_id, timestamp, correlation_id, details)
VALUES ($1, $2, $3, $4, $5, $6, $7)`,
[
event.eventId,
event.aggregateId,
event.eventType,
event.metadata.userId,
event.metadata.timestamp,
event.metadata.correlationId,
JSON.stringify(event.payload),
]
);
}
}Snapshots para el rendimiento
Reproducir miles de eventos para cada comando es costoso. Los snapshots capturan el estado del agregado periódicamente, de modo que la reproducción solo necesite los eventos desde el último snapshot.
class SnapshotStore {
constructor(private readonly db: Database) {}
async save(
aggregateId: string,
state: Record<string, unknown>,
version: number
): Promise<void> {
await this.db.query(
`INSERT INTO snapshots (aggregate_id, state, version, created_at)
VALUES ($1, $2, $3, NOW())
ON CONFLICT (aggregate_id) DO UPDATE
SET state = $2, version = $3, created_at = NOW()`,
[aggregateId, JSON.stringify(state), version]
);
}
async load(
aggregateId: string
): Promise<{ state: Record<string, unknown>; version: number } | null> {
const result = await this.db.query(
"SELECT state, version FROM snapshots WHERE aggregate_id = $1",
[aggregateId]
);
return result.rows[0] ?? null;
}
}
class AccountRepository {
private static SNAPSHOT_INTERVAL = 100;
constructor(
private readonly eventStore: EventStore,
private readonly snapshotStore: SnapshotStore
) {}
async load(aggregateId: string): Promise<AccountAggregate> {
const snapshot = await this.snapshotStore.load(aggregateId);
const fromVersion = snapshot?.version ?? 0;
const events = await this.eventStore.getEvents(aggregateId, fromVersion);
let aggregate: AccountAggregate;
if (snapshot) {
aggregate = AccountAggregate.fromSnapshot(snapshot.state, snapshot.version);
aggregate.replayEvents(events);
} else {
aggregate = AccountAggregate.fromEvents(events);
}
return aggregate;
}
async save(aggregate: AccountAggregate): Promise<void> {
const events = aggregate.getUncommittedEvents();
await this.eventStore.append(
aggregate.id,
events,
aggregate.version - events.length
);
if (aggregate.version % AccountRepository.SNAPSHOT_INTERVAL === 0) {
await this.snapshotStore.save(
aggregate.id,
aggregate.toSnapshot(),
aggregate.version
);
}
}
}Versionado y migración de eventos
Los eventos son inmutables, pero tu comprensión del dominio evoluciona. Cuando cambian los esquemas de eventos, usa upcasters para transformar los eventos antiguos al formato actual durante la reproducción.
type Upcaster = (event: DomainEvent) => DomainEvent;
const upcasters: Map<string, Upcaster[]> = new Map([
[
"MoneyWithdrawn",
[
// v1 → v2: Added 'channel' field
(event) => {
if (!event.payload.channel) {
return {
...event,
payload: { ...event.payload, channel: "unknown" },
};
}
return event;
},
// v2 → v3: Renamed 'previousBalance' to 'balanceBefore'
(event) => {
if ("previousBalance" in event.payload) {
const { previousBalance, ...rest } = event.payload;
return {
...event,
payload: {
...rest,
balanceBefore: previousBalance,
},
};
}
return event;
},
],
],
]);
function upcastEvent(event: DomainEvent): DomainEvent {
const eventUpcasters = upcasters.get(event.eventType) ?? [];
return eventUpcasters.reduce((e, upcaster) => upcaster(e), event);
}Puntos clave
Event sourcing convierte cada cambio de estado en un registro permanente y consultable. Para dominios con auditoría exigente —finanzas, salud, cumplimiento— no es opcional, es la arquitectura correcta. El event store es de solo apéndice e inmutable: los eventos registran hechos que ya ocurrieron.
Los agregados aplican las reglas de negocio y emiten eventos. Las proyecciones consumen eventos para construir modelos de lectura optimizados para cualquier patrón de consulta. Los snapshots evitan la degradación del rendimiento a medida que crecen los streams de eventos. Los upcasters gestionan la evolución del esquema sin corromper el registro histórico.
Empieza con un event store simple, un agregado y una proyección. Añade complejidad —snapshots, upcasters, múltiples proyecciones— solo cuando el dominio lo exija. El registro de auditoría que construyas hoy es la respuesta de cumplimiento que darás mañana.


