Aller au contenu
Français

Replicated snapshots

Ce contenu n’est pas encore disponible dans votre langue.

In single-writer event sourcing, snapshots save the state at a particular seqNr — recovery loads the snapshot, replays events after that seqNr.

In replicated event sourcing, the picture is more complex — there’s no single linear seqNr; instead, events are partially ordered across replicas via vector clocks. Snapshots must carry the vector clock alongside the state.

Snapshot contents:
- state (after applying all events causally seen at the time)
- vector clock ({ A: 100, B: 95, C: 60 })

On recovery:

1. Load snapshot.
2. Read events from journal AFTER the snapshot's vc.
3. Apply each (re-running conflict resolution if any are concurrent).
4. Ready.

The vc lets the recovering replica skip events that causally precede the snapshot (already incorporated).

Same picking heuristics as single-writer ES, but bias toward more frequent:

  • Replicated workloads accumulate events from multiple replicas at once.
  • The journal grows faster (N replicas × per-replica rate).
  • Recovery has to re-run conflict resolution for each not-snapshotted concurrent event.
class Account extends ReplicatedEventSourcedActor<...> {
override snapshotPolicy() {
return everyNEvents(100); // every 100 events
}
}

For replicated entities accumulating 1000 events/day across all replicas, snapshot every 100 means at most 100 events to re-process on recovery — sub-second.

A ReplicatedSnapshot carries more than state + clock:

{
state: State,
vc: { A: 100, B: 95, C: 60 }, // current vector-clock view
seenIds: [ ... ], // dedupe set: (replica, seqAtReplica)
events: [ ... ], // canonical sorted history — the heaviest field
localSeq: 205, // _localSeq at snapshot time
journalSeqAtSnapshot: 4096, // recovery resumes from here + 1
takenBy: 'A', // replica that took it (informational)
takenAt: 1716297600000 // wall-clock ms (informational)
}

The snapshot is a normal snapshot blob persisted to the standard snapshot store; existing stores (in-memory, SQLite, object-storage) handle it without changes. Note it carries the full deduped events history and the seenIds dedupe set — not just state + clock — so an out-of-order remote arrival can refold without re-reading the whole journal. events is the heaviest field: a 100k-event actor snapshots 100k envelopes on disk.

preStart():
↓ load latest snapshot
↓ state = snapshot.state
↓ vc = snapshot.vc
↓ read events from journal
↓ for each event in journal:
↓ if event.vc <= snapshot.vc: skip (already incorporated)
↓ else if event.vc concurrent with vc: invoke resolver
↓ else: apply via onEvent
↓ ready

The “skip” case is what makes snapshots bound recovery time — events written before the snapshot are skipped.

All standard snapshot stores work:

  • InMemorySnapshotStore — tests.
  • SqliteSnapshotStore — single-node (rare for replicated ES).
  • ObjectStorageSnapshotStore — shared across replicas.

For replicated ES with multiple replicas across regions, a shared snapshot store is critical — each replica recovers faster if it can load the latest snapshot from any replica, not just its own.

{
const objectStorageSnapshotStoreOptions = ObjectStorageSnapshotStoreOptions.create().withBackend(new S3ObjectStorageBackend(S3ObjectStorageOptions.create() /* shared bucket */));
journal: sharedJournal,
snapshotStore: new ObjectStorageSnapshotStore(objectStorageSnapshotStoreOptions),
}

Long-lived deployments accumulate retired replicas in vector clocks:

vc { A: 1000, B: 500, C: 200, RETIRED-D: 50, RETIRED-E: 30 }

The retired replicas’ components are inert but take space + slow comparisons.

Vector-clock garbage collection is out of scope for v1 — there is no pruning hook, and a snapshot stores the full, unpruned vector clock. VC entries grow with the set of replicas ever seen: fine for a stable cluster, but a node-churn-heavy deployment will eventually want compaction. Keep this in mind before relying on aggressive replica rotation.

Replica A writes event_A at t1.
Replica A snapshots at t2 (sees state with event_A).
snapshot.vc = { A: 1 }
Meanwhile, replica B was concurrently writing event_B at t1.5.
event_B has vc { B: 1 }; not in snapshot.
Replica A reads event_B at t3:
snapshot.vc { A: 1 } vs event_B.vc { B: 1 } → concurrent
invoke resolver, apply.

Concurrent events that arrive after the snapshot are handled at read-time by the resolver — same as without snapshots.

Snapshot writes for replicated ES are slightly heavier than single-writer:

  • Vector clock serialization — typically 50-200 bytes extra per snapshot.
  • Resolver state merging — if the snapshot is taken during concurrent-write reconciliation, the merge runs first.

Negligible in most cases. Bigger snapshots come from the state itself.