Zum Inhalt springen
Deutsch

Projektionen

Ein PersistentActor schreibt Events. Eine Projektion konsumiert sie und baut eine Read-Side-View, die für Queries zugeschnitten ist:

Read-Seite

Write-Seite

PersistentActor

Commands rein

Journal

Events

ProjectionActor

Handler

Read-Model

SQL / Redis / etc.

Das Journal ist Append-only und autoritativ. Projektionen sind abgeleitet — sie können von Grund auf neu aufgebaut werden, indem das Journal abgespielt wird. Das entkoppelt Write-Durchsatz (durable, append-only) von Read-Durchsatz (denormalisiert, query-optimiert).

import { match } from 'ts-pattern';
import { ProjectionActor, ByTagProjectionOptions } from 'actor-ts/persistence';
import { SqliteJournal, SqliteJournalOptions, SqliteQuery, InMemoryOffsetStore } from 'actor-ts/persistence';
type AccountEvent =
| { kind: 'deposited'; amount: number }
| { kind: 'withdrawn'; amount: number };
// Die Read-Seite fragt dasselbe Journal ab, in das deine PersistentActors
// schreiben — SqliteQuery umschließt eine SqliteJournal-Instanz (keinen Pfad).
const eventsJournal = new SqliteJournal(SqliteJournalOptions.create().withPath('/var/lib/events.db'));
const byTagProjectionOptions = ByTagProjectionOptions.create<AccountEvent>()
.withName('account-balance-view')
.withTag('account')
.withQuery(new SqliteQuery(eventsJournal))
.withOffsetStore(new InMemoryOffsetStore())
.withHandle(async (event) => {
await match(event.event)
.with({ kind: 'deposited' }, (e) => viewDb.execute(
'UPDATE balances SET balance = balance + ? WHERE pid = ?',
[e.amount, event.persistenceId],
))
.otherwise(() => Promise.resolve());
});
const projection = ProjectionActor.byTag<AccountEvent>(system, byTagProjectionOptions);

Der Actor:

  1. Lädt seinen Offset-Cursor beim preStart aus dem Offset-Store.
  2. Pollt die Query-Schicht für Events, die zum tag passen, ab dem Cursor.
  3. Ruft handle für jedes Event auf.
  4. Persistiert den neuen Cursor.
  5. Wiederholt.
FactoryCursor-TypVerwendung
ProjectionActor.byPersistenceId(...)sequenceNr (pro pid)Lies die vollständige Historie einer Entity.
ProjectionActor.byTag(...)Offset (Timestamp + Tiebreaker)Lies alles, was mit <tag> getaggt ist, im Journal.

Tag-basiert ist der häufige Fall — der PersistentActor ruft tagsFor(event) auf, um Events zu labeln; die Projektion abonniert das Tag. Per-pid ist nützlich für schmale Views (die Aktivität eines Users).

tagsFor wird pro Event gefragt, eine Projektion, die auf ein Tag filtert, sieht also genau die Events, die es tragen — auch dann, wenn die Events gemeinsam von einem persistAll geschrieben wurden.

withTag nimmt einen TagFilter, eine Projektion ist also nicht auf ein einzelnes Tag beschränkt:

import type { TagFilter } from 'actor-ts/persistence';
const eitherOne: TagFilter = { any: ['debit', 'credit'] };
const bothOfThem: TagFilter = { all: ['account', 'production'] };
const allButCancelled: TagFilter = { all: ['orders'], not: ['cancelled'] };
const byTagProjectionOptions = ByTagProjectionOptions.create<AccountEvent>()
.withName('open-orders-view')
.withTag(allButCancelled);

Jeder Operator ist optional, verknüpft werden sie mit UND. Ein leeres {} matcht jedes Event.

Der Filter gehört zur Cursor-Identität. Ein blanker String behält den Cursor, den er schon hat, aber die beiden Schreibweisen teilen sich keinen: .withTag('orders') legt seinen Offset unter orders ab, .withTag({ all: ['orders'] }) unter all(orders). Beide wählen dieselben Events aus, das Umschreiben der einen in die andere sieht also harmlos aus — und spielt die Projektion stillschweigend von vorn ab. Entscheide dich für eine Schreibweise, bevor die Projektion das erste Mal läuft, und bleib dabei.

// Crash-Recovery-Sequenz:
// 1. Offset-Cursor auf Wert N speichern.
// 2. Handler verarbeitet Event N+1.
// 3. Crash bevor der Cursor gespeichert wird.
// 4. Beim Neustart: Cursor ist immer noch N → Handler erhält Event N+1 erneut.

Wenn der Handler läuft, aber der Cursor nicht persistiert wird, verarbeitet die Projektion dasselbe Event erneut beim Neustart. Das ist At-least-once Delivery — das Framework garantiert, dass kein Event verpasst wird, aber Duplikate sind möglich.

Handler müssen idempotent sein:

  • UPSERT in das Read-Model (kein blindes INSERT).
  • Verarbeitete Event-IDs verfolgen im Read-Model selbst für Dedup.
  • Die sequenceNr des Events als Per-pid-Dedup-Key verwenden — sie nimmt nie ab, ist monoton pro persistenceId.

Wenn du den Handler nicht idempotent machen kannst, muss die Projektion an einem 2-Phase-Commit mit dem Offset-Save teilnehmen — viel komplexer, nicht out-of-the-box bereitgestellt.

Der Cursor rückt erst vor, nachdem handle zurückgekehrt ist — ein fehlschlagender Handler blockiert also konstruktionsbedingt den Kopf der Warteschlange: jedes spätere Event wartet hinter dem, das fehlgeschlagen ist. recoveryStrategy entscheidet, wie lange das so bleibt.

StrategieRetriesDanachWähle sie, wenn
retry-and-fail (Default)jastoppt die ProjektionEine still veraltete View schlimmer ist als eine sichtbar gestoppte.
retry-and-skipjaüberspringt das EventVerfügbarkeit vor Vollständigkeit geht — ein Dashboard, ein Suchindex.
failneinstoppt sofortDer Handler bereits intern wiederholt.
skipneinüberspringt sofortManche Events schlicht erwartbar nicht zustellbar sind.
const byTagProjectionOptions = ByTagProjectionOptions.create<AccountEvent>()
.withName('account-balance-view')
.withTag('account')
.withRecoveryStrategy('retry-and-skip')
.withMaxRetries(5)
.withRetryBackoffMs(2_000)
.withMaxRetryBackoffMs(60_000)
.withOnFailure((failure) => {
console.error(
`${failure.projection}: ${failure.event.persistenceId}#${failure.event.sequenceNr} `
+ `→ ${failure.action} (attempt ${failure.attempt})`,
failure.error,
);
});

Retries erfolgen nicht sofort. Der nächste Versuch wird um retryBackoffMs verzögert und verdoppelt sich jedes Mal, bis maxRetryBackoffMs erreicht ist. maxRetries zählt die Versuche nach dem ersten — der Default von 3 bedeutet also, dass ein schlechtes Event den Handler viermal erreicht, bevor die Projektion es aufgibt.

  • onFailure — wird bei jedem Versuch aufgerufen, mit dem betroffenen PersistentEvent, dem Fehler, der Versuchsnummer und der getroffenen Aktion (retry / skip / stop). Das ist das strukturierte Signal; häng dein Alerting daran. Ein Hook, der selbst wirft, wird geloggt und verschluckt — ein kaputter Reporter darf die Projektion nicht mitreißen.

  • Ein Dead Letter — ein übersprungenes Event wird, in einen DeadLetter verpackt, auf dem System-Event-Stream veröffentlicht und nennt die Projektion als seinen Empfänger, damit nichts spurlos verschwindet:

    import { DeadLetter } from 'actor-ts';
    system.eventStream.subscribe(auditRef, DeadLetter);

    Publishen ist standardmäßig alles, was passiert — ohne Abonnent ist das Event weg. Schaltet man die Dead-Letter-Queue ein, werden übersprungene Events stattdessen aufbewahrt, filterbar nach der Projektion, die sie übersprungen hat, und erneut zustellbar, sobald der Handler repariert ist.

  • Metriken — persistence_projection_stalled, persistence_projection_failures_total und persistence_projection_events_skipped_total; siehe Stock-Metriken.

  • Logs — eine Error-Zeile, wenn eine Fehlerserie beginnt, und eine, wenn sie endet. Die Retries dazwischen gehen nach debug — eine einen Tag lang festhängende Projektion kostet also zwei Error-Zeilen statt einer pro Poll.

fail stoppt den Actor, statt den Fehler aus dem Tick entkommen zu lassen. Das ist Absicht: der Supervisor der Projektions-Gruppe startet seine Kinder neu, ein entkommender Fehler würde also denselben Cursor neu laden und am selben Event scheitern — eine Restart-Schleife, die lauter ist als das ersetzte Drehen und kein bisschen nützlicher.

Eine gestoppte Projektion behält ihren Cursor: Ursache beheben und neu starten setzt genau bei dem Event fort, das fehlgeschlagen ist.

import { InMemoryOffsetStore, DurableStateOffsetStore } from 'actor-ts/persistence';
// Default (verloren beim Neustart):
new InMemoryOffsetStore();
// Dauerhaft — speichert den Cursor in einem beliebigen DurableStateStore
// (Postgres, MariaDB, libSQL, SQL Server, MongoDB, DynamoDB, Object-Storage), sodass ein Neustart fortsetzt:
new DurableStateOffsetStore(durableStateStore);

Der Cursor ist nur eine Zahl (oder ein Compound-Offset) — pro Projection-Name + Scope gespeichert. Implementierungen:

  • InMemoryOffsetStore — okay für Tests, nutzlos in Produktion (jeder Neustart verarbeitet vom Anfang neu).
  • DurableStateOffsetStore — umschließt einen beliebigen DurableStateStore (Postgres, MariaDB, libSQL, SQL Server, MongoDB, DynamoDB, Object-Storage), sodass Offsets Neustarts überstehen.
  • Benutzerdefiniert — implementiere OffsetStore gegen deinen eigenen Store (Redis, was auch immer dein Read-Model verwendet).

Für echte Produktions-Setups kollokiere den Offset-Store mit dem Read-Model, sodass sie zusammen abstürzen — das minimiert das Re-Processing-Fenster.

Im Cluster braucht der Offset-Store dieselbe Eigenschaft wie das Journal: eine Datenbank, nicht eine pro Knoten. Ein DurableStateOffsetStore über node-lokalem Storage gibt jedem Knoten seinen eigenen Cursor — eine Projektion, die zwischen Knoten wandert, verarbeitet dann Events erneut, die ein Knoten schon abgearbeitet hat, oder überspringt die, hinter denen der Cursor des anderen bereits steht. Siehe Storage-Lokalität & Identität.

const byTagProjectionOptions = ByTagProjectionOptions.create()
.withName('...')
.withTag('...')
.withLiveOptions({
pollIntervalMs: 500, // Default 1000 ms
});
ProjectionActor.byTag(system, byTagProjectionOptions
/* .withQuery(...).withOffsetStore(...).withHandle(...) */);

Die Projektion pollt. Im Idle sind die Poll-Kosten eine Journal-Query pro pollIntervalMs. Tuning:

  • Niedriger (250-500 ms) → schnellere End-to-End-Propagation, mehr Datenbank-Last.
  • Höher (5-10 s) → langsame Sichtbarkeit, aber billig.

Für sehr niedriglatente Read-Model-Updates siehe Push-basierte Query, die Sub-Poll-Intervall-Delivery über den In-Process-Event-Bus erhält.

const byPidProjectionOptions = ByPersistenceIdProjectionOptions.create<AccountEvent>()
.withName('account-42-view')
.withPersistenceId('account-42')
.withQuery(query)
.withOffsetStore(offsetStore)
.withHandle(async (event) => {
// ... nur die Events dieses Accounts behandeln
});
const projection = ProjectionActor.byPersistenceId<AccountEvent>(system, byPidProjectionOptions);

Ein Actor pro persistenceId. Nützlich für Per-Entity-Views:

  • Eine “User-Activity-Timeline”-Projektion pro User.
  • Eine “Per-Order-Audit-Trail”-Projektion pro Order.

Für große Zahlen von pids ist das nicht der Weg zu skalieren — eine Projektion pro pid zu spawnen skaliert nicht auf Millionen von Usern. Verwende dafür eine einzelne Tag-basierte Projektion, die nach pid hasht.

Fan-out: eine starten, sobald eine Entity auftaucht

Abschnitt betitelt „Fan-out: eine starten, sobald eine Entity auftaucht“

Eine Per-pid-Projektion muss ihre pid vorab kennen — passend für eine feste Menge, nutzlos für eine wachsende. allPersistenceIds schließt diese Lücke: ein Live-Stream jeder Id, die das Journal gesehen hat, plus jeder neuen, sobald sie zum ersten Mal auftaucht:

for await (const persistenceId of query.allPersistenceIds()) {
const options = ByPersistenceIdProjectionOptions.create<AccountEvent>()
.withName(`view-${persistenceId}`)
.withPersistenceId(persistenceId)
.withQuery(query)
.withOffsetStore(offsetStore)
.withHandle(async (event) => { /* ... */ });
ProjectionActor.byPersistenceId<AccountEvent>(system, options);
}

Zwei Vorbehalte, dieselben wie bei der Per-pid-Form — hier nur sichtbar gemacht:

  • Die Schleife endet nie, führe sie also in einem eigenen Task aus und breake beim Shutdown heraus.
  • Ein Actor pro Entity bleibt ein Actor pro Entity. Fan-out ist für hunderte oder tausende langlebige Entities gedacht, nicht für Millionen kurzlebiger — dort ist eine Tag-Projektion die Antwort.

Für einen einmaligen Durchlauf statt eines Live-Streams — ein Backfill über die Entities, die es jetzt gibt — nimm currentPersistenceIdsPaginated: es endet, sobald das Journal erschöpft ist, und hält eine Seite auf einmal statt der ganzen Id-Liste. Siehe Persistence Query.

// Drei Projektionen, drei Offset-Cursors, drei unabhängige Views:
const byTagProjectionOptions = ByTagProjectionOptions.create<E>().withName('balance');
ProjectionActor.byTag<E>(system, byTagProjectionOptions /* .withTag(...).withQuery(...).withHandle(...) */);
const byTagProjection2Options = ByTagProjectionOptions.create<E>().withName('audit-log');
ProjectionActor.byTag<E>(system, byTagProjection2Options /* .withTag(...).withQuery(...).withHandle(...) */);
const byTagProjection3Options = ByTagProjectionOptions.create<E>().withName('monthly-stats');
ProjectionActor.byTag<E>(system, byTagProjection3Options /* .withTag(...).withQuery(...).withHandle(...) */);

Jede hat ihren eigenen Offset-Cursor; sie lesen unabhängig vom Journal. Das ist die Stärke von Event Sourcing — ein Event-Stream, viele abgeleitete Views, keine an die andere gekoppelt.

Die ProjectionActor-API-Referenz deckt alle Settings ab.