CQRS + Event Sourcing — 02 — Event store : interface et implémentations

L'interface abstraite EventStore et ses trois implémentations : PostgreSQL (self-hosted), Firestore (GCP/Firebase), DynamoDB (AWS). Append-only, concurrence optimiste, TypeScript et Python.

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 :

  • TransactWriteItems est 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_READS pour 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.

Réservez un audit gratuit de 30 minutes. Je vous montre concrètement ce qu'on peut automatiser.