Skip to content
English

Replicated event sourcing overview

Standard event sourcing has one writer per persistenceId — a single actor instance appends events; another instance (after failover) replays them. This is fine for sharded entities where the framework guarantees one-actor-per-key.

Replicated event sourcing removes that constraint. Multiple replicas of the same entity can be active at once — on different nodes, in different regions — and each persists independently. The framework uses vector clocks to detect concurrent edits + conflict resolvers to merge them.

persist event_A1

persist event_B1

gossip /

async replication

Replica A

eu-west

Replica B

us-east

Shared journal

Shared journal

converge via vector clock + resolver

This is the niche persistence pattern. Most apps shouldn’t need it. Use cases:

  • Multi-region active-active — same entity writable in EU + US. Network partition between regions doesn’t stop either side.
  • Edge-style replication — entities replicate close to users, reconcile centrally.
  • Cluster-spanning concurrent writers — same entity edited on multiple cluster nodes without singleton coordination.

For typical sharded-entity setups, ClusterSharding + PersistentActor gives exactly-one-writer per key automatically — simpler than this.

Replicated event sourcing trades simplicity for availability:

Single-writer ESReplicated ES
Total event order per pidPartial order — concurrent events can be unordered
State is a deterministic foldState is a fold + conflict resolution
Commands are validated against latest stateCommands are validated against local replica’s view
Restart replays the logRestart replays the log + reconciles concurrent branches

The mental model is CRDT-like for events — multi-writer convergence 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> = {
// When two replicas concurrently set different values, the higher wins:
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.
}

The actor extends ReplicatedEventSourcedActor instead of PersistentActor. Nothing is passed into the constructor: the base class reads the Cluster off the ActorSystem (see this.cluster). Two additional things to specify:

  • resolver() — how to merge concurrent events (optional; defaults to last-writer-wins).
  • The journal must be shared across replicas (Cassandra, shared object-storage, etc.).

replicaId — this replica’s stable identifier, distinct from persistenceId — defaults to this node’s cluster address. Override it only when the id has to survive a re-address, e.g. a fixed region name:

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

Two replicas that ever share an id will share a vector-clock component, so the default is the safe one.

Cross-replica delivery rides on DistributedPubSub, so every field of an arriving event envelope was written by another node. Anything that can complete the cluster handshake can reach a replicated actor — through the pub-sub mediator, or straight at its actor path — and membership is not a precondition. Three properties follow, and all three matter operationally.

An envelope is validated whole, before anything is applied. replica and eventId must be non-empty and bounded, seqAtReplica a positive integer, timestamp finite, and the vector clock a plain object of finite non-negative components with a bounded entry count. A failure drops that one envelope and logs a WARN naming the field; it never fails the actor and never leaves a half-applied event in the history.

Event identity comes from entropy, not from the payload’s counters. Each event carries an eventId minted at persist time — the replica id, a #, and 96 bits of randomness — and that is what deduplication compares. It has to be unguessable because a deduplication hit means silently discard: an identity derived from replica and seqAtReplica is arithmetic to predict, so a peer could claim identities a replica had not issued yet and its real events would then be dropped by everyone, permanently, since the identity set is snapshotted. seqAtReplica is still on the envelope, but only as the last tie-break in the canonical order.

An envelope’s author is the node that sent it. replica is a value the sender chose, so it is held against the node the connection authenticated rather than believed. A replicated actor subscribes to its topic asking for the publishing node, and isAuthorizedAuthor decides whether that node may write under the arriving replica id; the default is the equality that the default replicaId makes true. An envelope that reaches the actor without an authenticated origin — which is what a frame aimed straight at the actor’s path produces, bypassing pub-sub — is refused for that reason alone.

The binding is sound because a replicated event crosses at most one hop: persist is the only publisher and no mediator re-forwards what it received, so the node an arriving envelope came from is the node that wrote it.

ComponentPurpose
VectorClockTracks causality across replicas — detects concurrent writes.
ConflictResolverDecides how to merge concurrent events into a single state.
Single-writer leaseOptional — gates writes via a lease for stronger consistency.
Replicated snapshotsSnapshots that include the vector clock for full recovery.

Each gets its own deep-dive page.

Sharding + PersistentActor Replicated ES
Multiple writers per entity? No (exactly one) Yes
Conflict resolution needed? No Yes
Cross-region active-active? Sharding favors one region Yes
Operational complexity? Low High
Use when Default You really need it