Event-Sourcing-Patterns für auditintensive Anwendungen
Event Sourcing für einen vollständigen, unveränderlichen Audit-Trail jeder Zustandsänderung: Projektionen, Snapshots und Event-Replay.

Warum der Zustand allein nicht ausreicht
Die meisten Anwendungen speichern den aktuellen Zustand von Entitäten. Der Kontostand eines Nutzers beträgt $500. Der Status einer Bestellung ist "shipped". Doch wenn ein Auditor fragt "wie ist dieser Kontostand auf $500 gekommen?" oder "wer hat den Bestellstatus wann geändert?", rätst du anhand lückenhafter Logs. Event Sourcing löst das, indem jede Zustandsänderung als unveränderliches Event gespeichert wird — die gesamte Historie wird zur Source of Truth.
Das Fundament des Event Stores
Ein Event Store ist ein Append-only-Log. Events werden nie aktualisiert oder gelöscht. Um den aktuellen Zustand einer Entität zu ermitteln, spielst du ihre Events von Anfang an wieder ab.
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;
}
}Aggregate, die Events erzeugen
Aggregate sind das Write Model in Event Sourcing. Commands kommen rein, Geschäftsregeln validieren sie, und das Aggregate emitiert Events, die beschreiben, was passiert ist — nicht, was angefragt wurde.
// ❌ 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];
}
}Projections für Read Models
Event Sourcing trennt Schreiben von Lesen. Der Event Store kümmert sich um Schreibvorgänge. Projections subscriben auf Events und bauen optimierte Read Models für Queries. Das ist CQRS, natürlich angewendet.
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 für die Performance
Tausende Events für jeden Command wieder abzuspielen ist teuer. Snapshots erfassen periodisch den Zustand des Aggregates, sodass das Replay nur noch die Events seit dem letzten Snapshot benötigt.
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
);
}
}
}Event-Versioning und Migration
Events sind unveränderlich, aber dein Domänenverständnis entwickelt sich weiter. Wenn sich Event-Schemas ändern, verwendest du Upcaster, um alte Events während des Replays in das aktuelle Format zu transformieren.
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);
}Kernaussagen
Event Sourcing macht aus jeder Zustandsänderung einen dauerhaften, abfragbaren Datensatz. Für auditintensive Domänen — Finanzen, Gesundheitswesen, Compliance — ist das keine Option, sondern die richtige Architektur. Der Event Store ist append-only und unveränderlich: Events dokumentieren Tatsachen, die bereits eingetreten sind.
Aggregate setzen Geschäftsregeln durch und emitieren Events. Projections konsumieren Events und bauen optimierte Read Models für beliebige Query-Patterns. Snapshots verhindern Performance-Einbußen, wenn Event-Streams wachsen. Upcaster behandeln Schema-Evolution, ohne den historischen Datensatz zu beschädigen.
Beginne mit einem einfachen Event Store, einem Aggregate und einer Projection. Füge Komplexität — Snapshots, Upcaster, mehrere Projections — nur hinzu, wenn die Domäne es erfordert. Der Audit-Trail, den du heute aufbaust, ist die Compliance-Antwort, die du morgen gibst.


