04 — Replay d'events et snapshots
Ce que tu vas apprendre
- Comment un agrégat se reconstitue depuis ses events (rehydration)
- Implémenter un agrégat avec la méthode
applypar event - Les snapshots : quand et comment les utiliser pour éviter de charger 10 000 events
- Le versioning des events : gérer l'évolution du schéma
- Exemples complets TypeScript et Python
Prérequis
La rehydration d'un agrégat
Dans l'Event Sourcing, l'état courant d'un agrégat n'est pas stocké directement. Il est reconstitué en rejouant la séquence des events depuis le début du stream.
Stream "order-abc" :
version 0 → OrderCreated { total: 150, status: "pending" }
version 1 → OrderConfirmed { confirmedAt: "2026-06-01T10:00" }
version 2 → OrderShipped { trackingId: "TRACK-123" }
État reconstitué :
{ status: "shipped", total: 150, trackingId: "TRACK-123", ... }
La reconstitution applique chaque event dans l'ordre sur un état initial vide.
L'agrégat avec méthode `apply`
Le pattern standard : l'agrégat a une méthode apply pour chaque type d'event. La méthode statique rehydrate charge les events et les applique.
typescript// L'état interne de l'agrégat
interface OrderState {
orderId: string;
customerId: string;
status: "pending" | "confirmed" | "shipped" | "delivered" | "cancelled";
items: Array<{ productId: string; quantity: number; unitPriceCents: number }>;
totalCents: number;
version: number;
}
class OrderAggregate {
private state: OrderState;
private uncommittedEvents: DomainEvent[] = [];
private constructor(state: OrderState) {
this.state = state;
}
// Reconstitution depuis les events stockés
static rehydrate(events: StoredEvent[]): OrderAggregate {
const initial: OrderState = {
orderId: "",
customerId: "",
status: "pending",
items: [],
totalCents: 0,
version: -1,
};
const aggregate = new OrderAggregate(initial);
for (const event of events) {
aggregate.applyStored(event);
}
return aggregate;
}
// Appliquer un event stocké (lecture depuis BDD)
private applyStored(event: StoredEvent): void {
this.applyEvent(event.eventType, event.data);
this.state.version = event.version;
}
// Appliquer un event (logique partagée entre rehydrate et les commands)
private applyEvent(type: string, data: Record<string, unknown>): void {
switch (type) {
case "OrderCreated":
this.state.orderId = data.orderId as string;
this.state.customerId = data.customerId as string;
this.state.items = data.items as any[];
this.state.totalCents = data.totalCents as number;
this.state.status = "pending";
break;
case "OrderConfirmed":
this.state.status = "confirmed";
break;
case "OrderShipped":
this.state.status = "shipped";
break;
case "OrderDelivered":
this.state.status = "delivered";
break;
case "OrderCancelled":
this.state.status = "cancelled";
break;
}
}
// Commands — vérifient les invariants, produisent des events
confirm(): void {
if (this.state.status !== "pending") {
throw new Error(`Impossible de confirmer une commande "${this.state.status}"`);
}
const event: OrderConfirmedEvent = {
eventType: "OrderConfirmed",
orderId: this.state.orderId,
confirmedAt: new Date().toISOString(),
};
this.applyEvent("OrderConfirmed", event as any);
this.uncommittedEvents.push(event);
}
cancel(reason: string): void {
if (["delivered", "cancelled"].includes(this.state.status)) {
throw new Error(`Impossible d'annuler une commande "${this.state.status}"`);
}
const event: OrderCancelledEvent = {
eventType: "OrderCancelled",
orderId: this.state.orderId,
reason,
cancelledAt: new Date().toISOString(),
};
this.applyEvent("OrderCancelled", event as any);
this.uncommittedEvents.push(event);
}
// Accesseurs
get status(): string { return this.state.status; }
get version(): number { return this.state.version; }
get pendingEvents(): DomainEvent[] { return [...this.uncommittedEvents]; }
clearEvents(): void { this.uncommittedEvents = []; }
}
En Python :
pythonfrom typing import List
from dataclasses import dataclass, field
@dataclass
class OrderState:
order_id: str = ""
customer_id: str = ""
status: str = "pending"
items: list = field(default_factory=list)
total_cents: int = 0
version: int = -1
class OrderAggregate:
def __init__(self, state: OrderState) -> None:
self._state = state
self._uncommitted_events: list = []
@classmethod
def rehydrate(cls, events: List['StoredEvent']) -> 'OrderAggregate':
aggregate = cls(OrderState())
for event in events:
aggregate._apply_stored(event)
return aggregate
def _apply_stored(self, event: 'StoredEvent') -> None:
self._apply_event(event.event_type, event.data)
self._state.version = event.version
def _apply_event(self, event_type: str, data: dict) -> None:
if event_type == "OrderCreated":
self._state.order_id = data["order_id"]
self._state.customer_id = data["customer_id"]
self._state.items = data["items"]
self._state.total_cents = data["total_cents"]
self._state.status = "pending"
elif event_type == "OrderConfirmed":
self._state.status = "confirmed"
elif event_type == "OrderShipped":
self._state.status = "shipped"
elif event_type == "OrderDelivered":
self._state.status = "delivered"
elif event_type == "OrderCancelled":
self._state.status = "cancelled"
def confirm(self) -> None:
if self._state.status != "pending":
raise ValueError(f"Impossible de confirmer une commande '{self._state.status}'")
from datetime import datetime, timezone
event = {
"event_type": "OrderConfirmed",
"order_id": self._state.order_id,
"confirmed_at": datetime.now(timezone.utc).isoformat(),
}
self._apply_event("OrderConfirmed", event)
self._uncommitted_events.append(event)
def cancel(self, reason: str) -> None:
if self._state.status in ("delivered", "cancelled"):
raise ValueError(f"Impossible d'annuler une commande '{self._state.status}'")
from datetime import datetime, timezone
event = {
"event_type": "OrderCancelled",
"order_id": self._state.order_id,
"reason": reason,
"cancelled_at": datetime.now(timezone.utc).isoformat(),
}
self._apply_event("OrderCancelled", event)
self._uncommitted_events.append(event)
@property
def status(self) -> str:
return self._state.status
@property
def version(self) -> int:
return self._state.version
@property
def pending_events(self) -> list:
return list(self._uncommitted_events)
def clear_events(self) -> None:
self._uncommitted_events.clear()
Le problème : streams longs
Pour un agrégat avec 3-10 events, le replay est instantané. Mais un compte bancaire actif depuis 5 ans peut avoir 50 000 transactions. Rejouer 50 000 events à chaque lecture est prohibitif.
La solution : les snapshots.
Les snapshots
Un snapshot est une capture de l'état de l'agrégat à un moment donné. Au lieu de rejouer tous les events depuis le début, on charge le dernier snapshot et on rejoue uniquement les events postérieurs.
sql-- Table des snapshots
CREATE TABLE snapshots (
stream_id VARCHAR(255) PRIMARY KEY,
state JSONB NOT NULL, -- L'état sérialisé de l'agrégat
version INTEGER NOT NULL, -- La version du dernier event inclus
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
typescriptclass SnapshotStore {
constructor(private readonly pool: Pool) {}
async save(streamId: string, state: unknown, version: number): Promise<void> {
await this.pool.query(
`INSERT INTO snapshots (stream_id, state, version)
VALUES ($1, $2, $3)
ON CONFLICT (stream_id) DO UPDATE
SET state = EXCLUDED.state,
version = EXCLUDED.version,
created_at = NOW()`,
[streamId, JSON.stringify(state), version]
);
}
async load(streamId: string): Promise<{ state: unknown; version: number } | null> {
const { rows } = await this.pool.query(
"SELECT state, version FROM snapshots WHERE stream_id = $1",
[streamId]
);
if (rows.length === 0) return null;
return { state: rows[0].state, version: rows[0].version };
}
}
// Chargement avec snapshot
class AggregateLoader {
constructor(
private readonly eventStore: PostgresEventStore,
private readonly snapshotStore: SnapshotStore
) {}
async load(streamId: string): Promise<{ aggregate: OrderAggregate; version: number }> {
// 1. Charger le dernier snapshot
const snapshot = await this.snapshotStore.load(streamId);
if (snapshot) {
// 2a. Snapshot trouvé — rejouer uniquement les events postérieurs
const { events, version } = await this.eventStore.loadFrom(
streamId,
snapshot.version + 1
);
const aggregate = OrderAggregate.fromSnapshot(snapshot.state as OrderState);
for (const event of events) {
aggregate.applyStored(event);
}
return { aggregate, version };
} else {
// 2b. Pas de snapshot — rejouer depuis le début
const { events, version } = await this.eventStore.load(streamId);
if (events.length === 0) {
throw new Error(`Stream introuvable : ${streamId}`);
}
const aggregate = OrderAggregate.rehydrate(events);
return { aggregate, version };
}
}
}
Snapshot tous les N events — une règle simple est de prendre un snapshot tous les 50 ou 100 events. Pour la plupart des agrégats métier, le stream ne dépasse pas 20-30 events, donc les snapshots ne sont pas nécessaires au démarrage.
Versioning des events
Les events sont immuables une fois écrits, mais leur schéma peut évoluer. Comment gérer un OrderCreatedEvent dont la structure a changé ?
La technique standard : l'upcasting. Quand tu charges un event de l'ancienne version, tu le transformes à la volée vers le schéma courant.
typescript// Event v1 (original)
// { orderId, customerId, total }
// Event v2 (nouveau champ shippingAddress ajouté)
// { orderId, customerId, total, shippingAddress }
function upcaster(event: StoredEvent): StoredEvent {
if (event.eventType === "OrderCreated" && !("shippingAddress" in event.data)) {
// Event v1 — ajouter la valeur par défaut
return {
...event,
data: {
...event.data,
shippingAddress: null, // Valeur par défaut pour les anciens events
},
};
}
return event; // Event déjà v2 — pas de transformation
}
// Dans l'EventStore.load — appliquer l'upcaster à tous les events chargés
async load(streamId: string): Promise<LoadResult> {
const raw = await this.loadRaw(streamId);
return {
events: raw.events.map(upcaster),
version: raw.version,
};
}
En Python :
pythondef upcast_event(event: StoredEvent) -> StoredEvent:
if event.event_type == "OrderCreated" and "shipping_address" not in event.data:
return StoredEvent(
**{**vars(event), "data": {**event.data, "shipping_address": None}}
)
return event
Résumé
| Concept | Quand l'utiliser | Implémentation |
|---|---|---|
| Rehydration | Toujours — reconstruction de l'état | apply par event_type |
| Snapshot | Streams longs (> 50-100 events) | Table snapshots + load partiel |
| Upcasting | Évolution du schéma des events | Transformation au chargement |
Étape suivante : 05 — Projet réel — tout ensemble dans un système de commandes complet.
Sources
- Young, G. (2010). CQRS Documents. cqrs.files.wordpress.com.
- Vernon, V. (2013). Implementing Domain-Driven Design. Addison-Wesley.
- Richardson, C. (2019). Microservices Patterns. Manning.
- Fowler, M. (2012). Patterns of Enterprise Application Architecture. Addison-Wesley.