02 — Event store : interface et implémentations
Ce que tu vas apprendre
- L'interface abstraite
EventStore— indépendante de l'infrastructure - Les trois propriétés invariables d'un event store : append-only, streams, concurrence optimiste
- Implémentation PostgreSQL :
UNIQUE INDEX, transactions,JSONB - Implémentation Firestore (GCP/Firebase) : transactions Firestore, subcollections
- Implémentation DynamoDB (AWS) : écriture conditionnelle, partition + sort key
- Comment choisir selon ton contexte
Prérequis
L'interface abstraite
Avant toute implémentation, définir le contrat. Quelle que soit l'infrastructure, un event store fait trois choses :
typescriptexport interface StoredEvent {
id: string;
streamId: string;
eventType: string;
data: Record<string, unknown>;
metadata: Record<string, unknown>;
version: number;
createdAt: string; // ISO 8601
}
export interface LoadResult {
events: StoredEvent[];
version: number; // -1 si le stream est vide
}
export interface RawEvent {
eventType: string;
data: Record<string, unknown>;
metadata?: Record<string, unknown>;
}
export interface EventStore {
/**
* Appende des events à un stream.
* expectedVersion = -1 pour un nouveau stream.
* Lance ConcurrencyError si la version courante ≠ expectedVersion.
*/
append(streamId: string, events: RawEvent[], expectedVersion: number): Promise<void>;
/** Charge tous les events d'un stream dans l'ordre. */
load(streamId: string): Promise<LoadResult>;
/** Charge les events depuis une version donnée (incluse). */
loadFrom(streamId: string, fromVersion: number): Promise<LoadResult>;
}
export class ConcurrencyError extends Error {
constructor(message: string) {
super(message);
this.name = "ConcurrencyError";
}
}
pythonfrom abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import List, Optional
@dataclass
class StoredEvent:
id: str
stream_id: str
event_type: str
data: dict
metadata: dict
version: int
created_at: str
@dataclass
class LoadResult:
events: List[StoredEvent]
version: int # -1 si vide
@dataclass
class RawEvent:
event_type: str
data: dict
metadata: dict = field(default_factory=dict)
class ConcurrencyError(Exception):
pass
class EventStore(ABC):
@abstractmethod
async def append(
self,
stream_id: str,
events: List[RawEvent],
expected_version: int,
) -> None: ...
@abstractmethod
async def load(self, stream_id: str) -> LoadResult: ...
@abstractmethod
async def load_from(self, stream_id: str, from_version: int) -> LoadResult: ...
Les trois implémentations ci-dessous respectent cette interface. Le reste du code (handlers, projections) ne dépend que de l'interface — pas de l'implémentation concrète.
Les 3 propriétés invariables
Quelle que soit l'implémentation, un event store garantit :
| Append-only | Streams ordonnés | Concurrence optimiste |
|---|---|---|
| Jamais de UPDATE/DELETE | Version croissante par stream | Rejet si version inattendue |
Append-only : un event enregistré ne change jamais. Pour corriger, on émet un nouvel event correctif.
Streams ordonnés : chaque event dans un stream a un numéro de version strictement croissant. L'ordre est garanti.
Concurrence optimiste : si deux handlers lisent le stream à la version N et tentent d'écrire simultanément à N+1, un seul réussit. L'autre reçoit un ConcurrencyError et doit retenter.
Implémentation 1 : PostgreSQL
Adapté pour : self-hosted (VPS, Kubernetes), services managés (Cloud SQL, Amazon RDS, Supabase).
Schéma
sqlCREATE TABLE events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
stream_id VARCHAR(255) NOT NULL,
event_type VARCHAR(255) NOT NULL,
data JSONB NOT NULL,
metadata JSONB NOT NULL DEFAULT '{}',
version INTEGER NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX events_stream_id_idx ON events (stream_id);
CREATE INDEX events_event_type_idx ON events (event_type);
CREATE UNIQUE INDEX events_stream_version ON events (stream_id, version);
-- L'UNIQUE INDEX est le mécanisme de concurrence optimiste
TypeScript
typescriptimport { Pool } from "pg";
export class PostgresEventStore implements EventStore {
constructor(private readonly pool: Pool) {}
async append(
streamId: string,
events: RawEvent[],
expectedVersion: number
): Promise<void> {
if (!events.length) return;
const client = await this.pool.connect();
try {
await client.query("BEGIN");
const { rows } = await client.query<{ max_version: number | null }>(
"SELECT MAX(version) AS max_version FROM events WHERE stream_id = $1",
[streamId]
);
const current = rows[0].max_version ?? -1;
if (current !== expectedVersion) {
throw new ConcurrencyError(
`Version conflict sur "${streamId}" : attendu ${expectedVersion}, trouvé ${current}`
);
}
for (let i = 0; i < events.length; i++) {
await client.query(
`INSERT INTO events (stream_id, event_type, data, metadata, version)
VALUES ($1, $2, $3, $4, $5)`,
[
streamId,
events[i].eventType,
JSON.stringify(events[i].data),
JSON.stringify(events[i].metadata ?? {}),
expectedVersion + 1 + i,
]
);
}
await client.query("COMMIT");
} catch (err) {
await client.query("ROLLBACK");
if ((err as any).code === "23505") {
throw new ConcurrencyError(`Conflit de concurrence sur "${streamId}"`);
}
throw err;
} finally {
client.release();
}
}
async load(streamId: string): Promise<LoadResult> {
const { rows } = await this.pool.query<StoredEvent>(
`SELECT id, stream_id AS "streamId", event_type AS "eventType",
data, metadata, version, created_at AS "createdAt"
FROM events WHERE stream_id = $1 ORDER BY version ASC`,
[streamId]
);
return { events: rows, version: rows.at(-1)?.version ?? -1 };
}
async loadFrom(streamId: string, fromVersion: number): Promise<LoadResult> {
const { rows } = await this.pool.query<StoredEvent>(
`SELECT id, stream_id AS "streamId", event_type AS "eventType",
data, metadata, version, created_at AS "createdAt"
FROM events WHERE stream_id = $1 AND version >= $2 ORDER BY version ASC`,
[streamId, fromVersion]
);
return { events: rows, version: rows.at(-1)?.version ?? fromVersion - 1 };
}
}
Python
pythonimport json
import asyncpg
class PostgresEventStore(EventStore):
def __init__(self, pool: asyncpg.Pool) -> None:
self._pool = pool
async def append(self, stream_id: str, events: List[RawEvent], expected_version: int) -> None:
if not events:
return
async with self._pool.acquire() as conn:
async with conn.transaction():
current = await conn.fetchval(
"SELECT MAX(version) FROM events WHERE stream_id = $1", stream_id
)
current_version = current if current is not None else -1
if current_version != expected_version:
raise ConcurrencyError(
f"Version conflict sur '{stream_id}' : "
f"attendu {expected_version}, trouvé {current_version}"
)
for i, event in enumerate(events):
await conn.execute(
"INSERT INTO events (stream_id, event_type, data, metadata, version) "
"VALUES ($1, $2, $3, $4, $5)",
stream_id,
event.event_type,
json.dumps(event.data),
json.dumps(event.metadata),
expected_version + 1 + i,
)
async def load(self, stream_id: str) -> LoadResult:
rows = await self._pool.fetch(
"SELECT * FROM events WHERE stream_id = $1 ORDER BY version ASC", stream_id
)
events = [self._map(r) for r in rows]
return LoadResult(events=events, version=events[-1].version if events else -1)
async def load_from(self, stream_id: str, from_version: int) -> LoadResult:
rows = await self._pool.fetch(
"SELECT * FROM events WHERE stream_id = $1 AND version >= $2 ORDER BY version ASC",
stream_id, from_version,
)
events = [self._map(r) for r in rows]
return LoadResult(events=events, version=events[-1].version if events else from_version - 1)
def _map(self, row) -> StoredEvent:
return StoredEvent(
id=str(row["id"]),
stream_id=row["stream_id"],
event_type=row["event_type"],
data=dict(row["data"]),
metadata=dict(row["metadata"]),
version=row["version"],
created_at=str(row["created_at"]),
)
Implémentation 2 : Firestore (GCP / Firebase)
Adapté pour : Google Cloud (Firestore natif), Firebase, projets serverless sans BDD relationnelle.
Structure des données
Firestore ne connaît pas le concept de "version entière incrémentale" nativement. On le simule avec un document de métadonnées par stream :
Collection : event_streams
Document : "order-abc123" ← métadonnées du stream
{ version: 2, updatedAt: ... }
Sous-collection : events
Document : "0" → { eventType, data, metadata, version: 0, createdAt }
Document : "1" → { eventType, data, metadata, version: 1, createdAt }
Document : "2" → { eventType, data, metadata, version: 2, createdAt }
TypeScript (Firebase Admin SDK)
typescriptimport { Firestore, FieldValue } from "firebase-admin/firestore";
export class FirestoreEventStore implements EventStore {
constructor(private readonly db: Firestore) {}
async append(
streamId: string,
events: RawEvent[],
expectedVersion: number
): Promise<void> {
if (!events.length) return;
const streamRef = this.db.collection("event_streams").doc(streamId);
await this.db.runTransaction(async (tx) => {
const streamDoc = await tx.get(streamRef);
const current: number = streamDoc.exists
? (streamDoc.data()!.version as number)
: -1;
if (current !== expectedVersion) {
throw new ConcurrencyError(
`Version conflict sur "${streamId}" : attendu ${expectedVersion}, trouvé ${current}`
);
}
// Écrire chaque event dans la sous-collection
for (let i = 0; i < events.length; i++) {
const version = expectedVersion + 1 + i;
const eventRef = streamRef.collection("events").doc(String(version));
tx.set(eventRef, {
eventType: events[i].eventType,
data: events[i].data,
metadata: events[i].metadata ?? {},
version,
createdAt: new Date().toISOString(),
streamId,
});
}
// Mettre à jour la version du stream
const newVersion = expectedVersion + events.length;
tx.set(streamRef, { version: newVersion, updatedAt: FieldValue.serverTimestamp() });
});
}
async load(streamId: string): Promise<LoadResult> {
const eventsRef = this.db
.collection("event_streams")
.doc(streamId)
.collection("events")
.orderBy("version", "asc");
const snapshot = await eventsRef.get();
const events: StoredEvent[] = snapshot.docs.map(doc => ({
id: doc.id,
streamId,
eventType: doc.data().eventType,
data: doc.data().data,
metadata: doc.data().metadata ?? {},
version: doc.data().version,
createdAt: doc.data().createdAt,
}));
return { events, version: events.at(-1)?.version ?? -1 };
}
async loadFrom(streamId: string, fromVersion: number): Promise<LoadResult> {
const eventsRef = this.db
.collection("event_streams")
.doc(streamId)
.collection("events")
.where("version", ">=", fromVersion)
.orderBy("version", "asc");
const snapshot = await eventsRef.get();
const events: StoredEvent[] = snapshot.docs.map(doc => ({
id: doc.id,
streamId,
eventType: doc.data().eventType,
data: doc.data().data,
metadata: doc.data().metadata ?? {},
version: doc.data().version,
createdAt: doc.data().createdAt,
}));
return { events, version: events.at(-1)?.version ?? fromVersion - 1 };
}
}
Points de vigilance Firestore :
- Une transaction Firestore est limitée à 500 documents et 10 secondes — ne pas appender plus de 499 events en une fois
- Firestore facture à la lecture et à l'écriture — charger un stream de 10 000 events coûte 10 000 lectures
- Firestore ne garantit pas l'ordre global des écritures entre streams — chaque stream est cohérent, pas l'ensemble
Implémentation 3 : DynamoDB (AWS)
Adapté pour : AWS, architectures serverless (Lambda), besoins de montée en charge élevée.
Structure de la table
Table : Events
Partition key : stream_id (String)
Sort key : version (Number) ← garantit l'unicité et l'ordre
Attributs : event_type, data (Map), metadata (Map), created_at (String)
La clé composite (stream_id, version) fait le même travail que l'UNIQUE INDEX PostgreSQL.
TypeScript (AWS SDK v3)
typescriptimport {
DynamoDBClient,
PutItemCommand,
QueryCommand,
TransactWriteItemsCommand,
} from "@aws-sdk/client-dynamodb";
import { marshall, unmarshall } from "@aws-sdk/util-dynamodb";
export class DynamoDBEventStore implements EventStore {
private readonly TABLE = "Events";
constructor(private readonly client: DynamoDBClient) {}
async append(
streamId: string,
events: RawEvent[],
expectedVersion: number
): Promise<void> {
if (!events.length) return;
// DynamoDB TransactWriteItems — max 100 opérations par transaction
if (events.length > 99) {
throw new Error("DynamoDB limite les transactions à 99 opérations");
}
const transactItems = events.map((event, i) => {
const version = expectedVersion + 1 + i;
return {
Put: {
TableName: this.TABLE,
Item: marshall({
stream_id: streamId,
version,
event_type: event.eventType,
data: event.data,
metadata: event.metadata ?? {},
created_at: new Date().toISOString(),
}),
// Écriture conditionnelle — échoue si (stream_id, version) existe déjà
ConditionExpression:
"attribute_not_exists(stream_id) AND attribute_not_exists(#v)",
ExpressionAttributeNames: { "#v": "version" },
},
};
});
try {
await this.client.send(
new TransactWriteItemsCommand({ TransactItems: transactItems })
);
} catch (err: any) {
if (
err.name === "TransactionCanceledException" &&
err.CancellationReasons?.some(
(r: any) => r.Code === "ConditionalCheckFailed"
)
) {
throw new ConcurrencyError(
`Conflit de concurrence sur le stream "${streamId}"`
);
}
throw err;
}
}
async load(streamId: string): Promise<LoadResult> {
const { Items = [] } = await this.client.send(
new QueryCommand({
TableName: this.TABLE,
KeyConditionExpression: "stream_id = :sid",
ExpressionAttributeValues: marshall({ ":sid": streamId }),
ScanIndexForward: true, // Ordre croissant sur version
})
);
const events: StoredEvent[] = Items.map((item) => {
const d = unmarshall(item);
return {
id: `${d.stream_id}#${d.version}`,
streamId: d.stream_id,
eventType: d.event_type,
data: d.data,
metadata: d.metadata ?? {},
version: d.version,
createdAt: d.created_at,
};
});
return { events, version: events.at(-1)?.version ?? -1 };
}
async loadFrom(streamId: string, fromVersion: number): Promise<LoadResult> {
const { Items = [] } = await this.client.send(
new QueryCommand({
TableName: this.TABLE,
KeyConditionExpression: "stream_id = :sid AND #v >= :from",
ExpressionAttributeNames: { "#v": "version" },
ExpressionAttributeValues: marshall({
":sid": streamId,
":from": fromVersion,
}),
ScanIndexForward: true,
})
);
const events: StoredEvent[] = Items.map((item) => {
const d = unmarshall(item);
return {
id: `${d.stream_id}#${d.version}`,
streamId: d.stream_id,
eventType: d.event_type,
data: d.data,
metadata: d.metadata ?? {},
version: d.version,
createdAt: d.created_at,
};
});
return { events, version: events.at(-1)?.version ?? fromVersion - 1 };
}
}
Points de vigilance DynamoDB :
TransactWriteItemsest limité à 100 opérations — un batch de plus de 99 events doit être découpé- DynamoDB n'a pas de transaction globale sur plusieurs tables — les projections doivent être gérées séparément (via DynamoDB Streams ou EventBridge)
- Les lectures consistent en
STRONGLY_CONSISTENT_READSpour garantir que tu lis le dernier état écrit
Tableau comparatif
| PostgreSQL | Firestore | DynamoDB | |
|---|---|---|---|
| Concurrence optimiste | UNIQUE INDEX |
Transaction Firestore sur doc metadata | ConditionExpression sur (stream_id, version) |
| Atomicité multi-events | Transaction SQL — illimité | Transaction Firestore — max 499 documents | TransactWriteItems — max 100 items |
| Ordre garanti | ORDER BY version ASC |
.orderBy("version", "asc") |
ScanIndexForward: true |
| Coût | Fixe (instance) | Au document lu/écrit | À la capacité (RCU/WCU) |
| Intégration projections | In-process ou déclencheur PG | Firestore Triggers → Cloud Functions | DynamoDB Streams → Lambda |
| Limite pratique streams | Pas de limite | 10 000 lectures = 10 000 ops Firestore | Partition hot key au-delà de 3 000 RCU/s |
| Self-hosted possible | Oui | Non (Google) | Non (AWS) |
Résumé
L'interface EventStore est stable — elle ne change pas selon l'infra. Ce qui change :
- Le mécanisme de concurrence optimiste (UNIQUE INDEX vs transaction Firestore vs ConditionExpression)
- Les limites de lot (illimité vs 499 vs 100)
- L'intégration avec les projections (synchrone in-process vs triggers cloud)
Le choix d'implémentation dépend de ton infra existante, pas des patterns CQRS eux-mêmes. L'article 06 — Choisir son infrastructure couvre les critères de décision complets.
Étape suivante : 03 — Projections : synchrones et asynchrones
Sources
- Young, G. (2010). CQRS Documents. cqrs.files.wordpress.com.
- Richardson, C. (2019). Microservices Patterns. Manning. (Chapitre 6)
- Google. Firestore Transactions and batched writes. cloud.google.com/firestore/docs.
- Amazon. DynamoDB Transactions. docs.aws.amazon.com/dynamodb.
- Betts, D., et al. (2012). Exploring CQRS and Event Sourcing. Microsoft patterns & practices.