Zum Inhalt springen
Deutsch

Replicated Event Sourcing im Überblick

Standard-Event Sourcing hat einen Writer pro persistenceId — eine einzelne Actor-Instanz hängt Events an; eine andere Instanz (nach Failover) spielt sie ab. Das ist in Ordnung für Sharded Entities, wo das Framework Ein-Actor-pro-Key garantiert.

Replicated Event Sourcing entfernt diese Einschränkung. Mehrere Replicas derselben Entity können gleichzeitig aktiv sein — auf verschiedenen Nodes, in verschiedenen Regionen — und jede persistiert unabhängig. Das Framework verwendet Vector Clocks, um nebenläufige Edits zu erkennen + Conflict Resolver, um sie zu mergen.

persist event_A1

persist event_B1

Gossip /

asynchrone Replikation

Replica A

eu-west

Replica B

us-east

geteiltes Journal

geteiltes Journal

konvergieren über Vector Clock + Resolver

Das ist das Nischen-Persistenz-Muster. Die meisten Apps sollten es nicht brauchen. Use Cases:

  • Multi-Region Active-Active — dieselbe Entity schreibbar in EU + US. Netzwerk-Partition zwischen Regionen stoppt keine Seite.
  • Edge-artige Replikation — Entities replizieren nah am User, gleichen zentral ab.
  • Cluster-spanning gleichzeitige Writer — dieselbe Entity wird auf mehreren Cluster-Nodes ohne Singleton-Koordination bearbeitet.

Für typische Sharded-Entity-Setups gibt ClusterSharding + PersistentActor automatisch Exactly-One-Writer pro Key — einfacher als das hier.

Replicated Event Sourcing tauscht Einfachheit gegen Verfügbarkeit:

Single-Writer-ESReplicated ES
Totale Event-Order pro pidPartielle Order — nebenläufige Events können ungeordnet sein
State ist ein deterministischer FoldState ist ein Fold + Conflict Resolution
Commands werden gegen den neuesten State validiertCommands werden gegen die View der lokalen Replica validiert
Neustart spielt das Log abNeustart spielt das Log ab + gleicht nebenläufige Branches ab

Das Mental-Model ist CRDT-artig für Events — Multi-Writer- Konvergenz by Design.

import {
ReplicatedEventSourcedActor,
VectorClock,
type ConflictResolver,
} from 'actor-ts/persistence';
type State = { value: number };
type Event = { kind: 'set'; value: number };
const maxWinsResolver: ConflictResolver<Event> = {
// Wenn zwei Replicas gleichzeitig unterschiedliche Werte setzen, gewinnt der höhere:
resolve(a, b) {
return a.event.value >= b.event.value ? a.event : b.event;
},
};
class Counter extends ReplicatedEventSourcedActor<Command, Event, State> {
readonly persistenceId = 'counter-42';
protected override resolver(): ConflictResolver<Event> { return maxWinsResolver; }
initialState() { return { value: 0 }; }
onEvent(state: State, event: Event) {
return { value: event.value };
}
// ... onCommand etc.
}

Der Actor erweitert ReplicatedEventSourcedActor statt PersistentActor. In den Konstruktor wird nichts hineingereicht: die Basisklasse liest den Cluster vom ActorSystem (siehe this.cluster). Zwei zusätzliche Dinge zu spezifizieren:

  • resolver() — wie nebenläufige Events gemergt werden (optional; Default ist Last-Writer-Wins).
  • Das Journal muss über Replicas geteilt sein (Cassandra, geteilter Object Storage, etc.).

replicaId — der stabile Identifier dieser Replica, anders als persistenceId — hat als Default die Cluster-Adresse dieses Nodes. Überschreibe ihn nur, wenn die Id eine Adressänderung überleben muss, z.B. ein fester Regionsname:

override get replicaId(): string { return process.env.REPLICA_ID!; }

Zwei Replicas, die sich jemals eine Id teilen, teilen sich eine Vector-Clock-Komponente — der Default ist also der sichere.

Die Cross-Replica-Auslieferung läuft über DistributedPubSub, jedes Feld eines eintreffenden Event-Envelopes wurde also von einem anderen Node geschrieben. Alles, was den Cluster-Handshake abschließen kann, erreicht einen replizierten Actor — über den PubSub-Mediator oder direkt über seinen Actor-Pfad — und Mitgliedschaft ist keine Voraussetzung. Daraus folgen drei Eigenschaften, und alle drei sind betrieblich relevant.

Ein Envelope wird vollständig validiert, bevor irgendetwas angewendet wird. replica und eventId müssen nicht-leer und begrenzt sein, seqAtReplica eine positive Ganzzahl, timestamp endlich und die Vector Clock ein einfaches Objekt aus endlichen, nicht-negativen Komponenten mit begrenzter Anzahl an Einträgen. Ein Fehlschlag verwirft genau diesen Envelope und loggt ein WARN, das das betroffene Feld benennt; er lässt den Actor nie fehlschlagen und hinterlässt nie ein halb angewendetes Event in der Historie.

Die Event-Identität kommt aus Entropie, nicht aus den Zählern des Payloads. Jedes Event trägt eine zur persist-Zeit erzeugte eventId — die Replica-Id, ein # und 96 Bit Zufall — und genau die vergleicht die Deduplizierung. Sie muss unerratbar sein, denn ein Deduplizierungs-Treffer bedeutet stilles Verwerfen: eine aus replica und seqAtReplica abgeleitete Identität ist per Arithmetik vorhersagbar, ein Peer könnte also Identitäten beanspruchen, die eine Replica noch nicht ausgegeben hat — ihre echten Events würden dann von allen verworfen, dauerhaft, weil das Identitäts-Set gesnapshottet wird. seqAtReplica steht weiterhin im Envelope, aber nur noch als letzter Tie-Break in der kanonischen Ordnung.

Der Autor eines Envelopes ist der Node, der ihn gesendet hat. replica ist ein vom Sender gewählter Wert und wird deshalb gegen den Node geprüft, den die Verbindung authentifiziert hat, statt geglaubt zu werden. Ein replizierter Actor abonniert sein Topic so, dass der publizierende Node mitgeliefert wird, und isAuthorizedAuthor entscheidet, ob dieser Node unter der eintreffenden replica-Id schreiben darf; der Default ist genau die Gleichheit, die der Default-replicaId wahr macht. Ein Envelope, der den Actor ohne authentifizierte Herkunft erreicht — was ein direkt an den Actor-Pfad adressierter Frame erzeugt, PubSub also umgeht —, wird allein deswegen abgelehnt.

Die Bindung trägt, weil ein repliziertes Event höchstens einen Hop zurücklegt: persist ist der einzige Publisher und kein Mediator leitet Empfangenes weiter — der Node, von dem ein Envelope kommt, ist also der Node, der ihn geschrieben hat.

KomponenteZweck
VectorClockVerfolgt Kausalität über Replicas — erkennt nebenläufige Writes.
ConflictResolverEntscheidet, wie nebenläufige Events in einen einzigen State gemergt werden.
Single-Writer-LeaseOptional — gattert Writes über ein Lease für stärkere Konsistenz.
Replicated SnapshotsSnapshots, die die Vector Clock für volle Recovery enthalten.

Jedes bekommt seine eigene Deep-Dive-Seite.

Sharding + PersistentActor Replicated ES
Mehrere Writer pro Entity? Nein (genau einer) Ja
Conflict Resolution nötig? Nein Ja
Cross-Region Active-Active? Sharding bevorzugt eine Region Ja
Operative Komplexität? Niedrig Hoch
Verwenden wenn Default Du brauchst es wirklich