Zum Inhalt springen

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.

4 Min. Lesezeit
Ein Event-Stream, der von Commands durch einen Event Store in mehrere Read Projections für unterschiedliche Datenansichten fließt

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.

tstypescript
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.

tstypescript
// ❌ 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.

tstypescript
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.

tstypescript
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.

tstypescript
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.

Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX