CQRS + Event Sourcing — 06 — Choisir son infrastructure

Guide complet pour choisir l'infrastructure CQRS + Event Sourcing : PostgreSQL seul, PostgreSQL + réplicas, GCP (Pub/Sub + Cloud SQL/Firestore), AWS (DynamoDB + EventBridge), Firebase.

06 — Choisir son infrastructure CQRS + Event Sourcing

Ce que tu vas apprendre

  • Ce qui change selon l'infrastructure (et ce qui ne change pas)
  • Les 5 architectures les plus courantes en production
  • Les trade-offs concrets de chacune
  • La matrice de décision pour choisir

Prérequis


Ce qui ne change jamais

Quelle que soit l'infrastructure, les concepts restent identiques :

  • Les commands expriment une intention de modification
  • Les events sont des faits immuables stockés dans un stream ordonné
  • Les projections construisent des read models depuis les events
  • La concurrence optimiste protège l'intégrité des streams
  • Le replay reconstruit n'importe quel état depuis les events

Ce qui change : les mécanismes techniques qui implémentent ces concepts.

L'infrastructure est un détail d'implémentation. Les patterns CQRS + Event Sourcing sont indépendants de la stack. Ce que tu choisis détermine la cohérence, le coût opérationnel et les limites — pas les concepts.


Architecture 1 : PostgreSQL seul (référence)

┌─────────────────────────────────────────────┐
│                  PostgreSQL                  │
│                                             │
│   Table events      Table order_summaries   │
│   (event store)     (read model)            │
│                                             │
└─────────────────────────────────────────────┘
         ▲                    ▲
         │ append()           │ UPDATE (synchrone)
         │                    │
    Command Handler ──────────┘
         (même process)

Cohérence : forte — command réussit → read model à jour immédiatement.

Infra requise : une instance PostgreSQL. Docker + PostgreSQL en local. Cloud SQL, RDS, ou Supabase en prod.

Concurrence optimiste : UNIQUE INDEX (stream_id, version).

Quand choisir :

  • Équipe de 1 à 5 développeurs
  • Pas de contrainte de montée en charge horizontale
  • Self-hosted (VPS, Kubernetes)
  • Besoin de simplicité opérationnelle
  • Prototype ou produit en phase de démarrage

Limites :

  • Montée en charge limitée par l'instance PostgreSQL
  • Projections synchrones : si une projection est lente, tout est lent
  • Pas de séparation physique entre write side et read side

Architecture 2 : PostgreSQL + réplicas en lecture

┌─────────────────────┐      ┌─────────────────────────┐
│   PostgreSQL Primary │      │   PostgreSQL Replica(s)  │
│                     │      │                          │
│   Table events      │─────►│   Table order_summaries  │
│   (event store)     │      │   (répliquée en lecture) │
│   Table projections │      │                          │
└─────────────────────┘      └─────────────────────────┘
         ▲                                ▲
         │ append()                       │ SELECT (queries)
    Command Handler                  Query Handler

Cohérence : éventuelle — replication lag de quelques ms à quelques secondes selon la charge.

Concurrence optimiste : même mécanisme que l'architecture 1, sur le primary.

Ce qui change :

typescript// Lire depuis le primary pour les cas critiques (read-your-writes)
class OrderQueryHandler {
  constructor(
    private readonly primaryPool: Pool,   // Pour les lectures critiques
    private readonly replicaPool: Pool    // Pour les lectures non-critiques
  ) {}

  // Lecture critique — depuis le primary
  async getOrderForPayment(orderId: string): Promise<OrderDTO | null> {
    const { rows } = await this.primaryPool.query(
      "SELECT * FROM order_summaries WHERE order_id = $1",
      [orderId]
    );
    return rows[0] ?? null;
  }

  // Lecture non-critique — depuis le replica (éventuelle)
  async listCustomerOrders(customerId: string): Promise<OrderDTO[]> {
    const { rows } = await this.replicaPool.query(
      "SELECT * FROM order_summaries WHERE customer_id = $1 ORDER BY created_at DESC",
      [customerId]
    );
    return rows;
  }
}

Stratégie read-your-writes : retourner l'état depuis la command plutôt que de lire immédiatement le read model.

typescript// Handler retourne les données — pas besoin de relire le read model
async createOrder(cmd: CreateOrderCommand): Promise<CreateOrderResult> {
  const orderId = await this.commandService.createOrder(cmd);
  // On retourne l'état construit depuis la command — pas depuis la BDD
  return { orderId, status: "pending", totalEuros: total / 100 };
}

Quand choisir :

  • Lecture ≫ écriture (catalogue produits, e-commerce)
  • Besoin de scale la lecture indépendamment de l'écriture
  • PostgreSQL géré (Cloud SQL, Amazon RDS) avec réplication activée

Limites :

  • Replication lag — accepter l'eventual consistency côté lecture
  • Coût : deux instances au minimum
  • Complexité : deux pools de connexions, logique de routing

Architecture 3 : GCP — Cloud SQL + Pub/Sub + Cloud Functions

                        Cloud SQL (Primary)
                         ┌──────────────┐
Command Handler          │  Table events│
      │ append()         └──────────────┘
      │                         │
      │ publish(event)           │ (optionnel : réplication vers Cloud SQL Replica)
      ▼                          │
  Cloud Pub/Sub ◄────────────────┘
      │
      ├──► Cloud Function A → PostgreSQL read model
      ├──► Cloud Function B → Firestore read model
      └──► Cloud Function C → BigQuery (analytics)

Cohérence : éventuelle — délai de quelques ms à secondes selon la charge Pub/Sub.

Event store : Cloud SQL (PostgreSQL managé) avec le schéma standard.

Projection :

typescript// Publisher — dans le command handler, après append
import { PubSub } from "@google-cloud/pubsub";

const pubsub = new PubSub({ projectId: "mon-projet" });

async function publishToTopic(events: StoredEvent[]): Promise<void> {
  const topic = pubsub.topic("domain-events");

  await Promise.all(
    events.map(event =>
      topic.publishMessage({
        data: Buffer.from(JSON.stringify(event)),
        attributes: {
          eventType: event.eventType,
          streamId: event.streamId,
        },
        orderingKey: event.streamId, // Garantit l'ordre par stream
      })
    )
  );
}
typescript// Subscriber — Cloud Function
import * as functions from "@google-cloud/functions-framework";

functions.cloudEvent("handleDomainEvent", async (cloudEvent) => {
  const data = Buffer.from(
    (cloudEvent.data as any).message.data, "base64"
  ).toString();

  const event: StoredEvent = JSON.parse(data);
  const projection = new OrderSummaryProjection(cloudSqlPool);

  if (projection.handles.includes(event.eventType)) {
    await projection.apply(event); // Doit être idempotente
  }
});

Points de vigilance GCP :

  • orderingKey sur Pub/Sub garantit l'ordre des messages par clé — utiliser streamId comme clé
  • Pub/Sub garantit at-least-once — les projections doivent être idempotentes
  • Cloud Functions peuvent démarrer à froid (cold start) — préférer Cloud Run pour les projections critiques
  • Pub/Sub Dead Letter Topic pour les messages qui échouent en permanence

Quand choisir :

  • Déjà sur GCP
  • Besoin de déclencher plusieurs services depuis les mêmes events
  • Analytics temps réel (BigQuery via Pub/Sub)
  • Serverless — pas d'infrastructure à gérer

Architecture 4 : Firebase (Firestore + Cloud Functions)

Command Handler
      │
      │ append() — transaction Firestore
      ▼
  Firestore (event_streams collection)
      │
      │ onWrite trigger automatique
      ▼
  Cloud Function
      │ apply(event)
      ▼
  Firestore (read models collection)

Cohérence : éventuelle — les triggers Firestore s'exécutent de manière asynchrone.

Event store : Firestore avec la structure subcollection définie dans l'article 02.

Projection via Firestore Triggers :

typescriptimport * as functions from "firebase-functions/v2/firestore";
import { getFirestore } from "firebase-admin/firestore";

// Se déclenche à chaque création d'un document dans la sous-collection events
export const onEventCreated = functions.onDocumentCreated(
  "event_streams/{streamId}/events/{version}",
  async (event) => {
    const data = event.data?.data();
    if (!data) return;

    const storedEvent: StoredEvent = {
      id: event.data!.id,
      streamId: event.params.streamId,
      eventType: data.eventType,
      data: data.data,
      metadata: data.metadata ?? {},
      version: data.version,
      createdAt: data.createdAt,
    };

    const db = getFirestore();

    if (storedEvent.eventType === "OrderCreated") {
      // Mettre à jour le read model dans Firestore
      await db.collection("order_summaries").doc(storedEvent.data.orderId as string).set({
        orderId: storedEvent.data.orderId,
        customerId: storedEvent.data.customerId,
        status: "pending",
        totalCents: storedEvent.data.totalCents,
        itemCount: (storedEvent.data.items as any[]).length,
        createdAt: storedEvent.createdAt,
        updatedAt: storedEvent.createdAt,
      }, { merge: false });
    }

    if (storedEvent.eventType === "OrderConfirmed") {
      await db.collection("order_summaries").doc(storedEvent.data.orderId as string).update({
        status: "confirmed",
        updatedAt: storedEvent.createdAt,
      });
    }
  }
);

Limitations Firebase spécifiques :

  • Firestore n'est pas un vrai event store — pas de numéro de séquence global, pas de loadFrom efficace sans index spécifique
  • Les triggers Cloud Functions s'exécutent au plus une fois par event en théorie, mais les retries en cas d'échec peuvent provoquer des doublons
  • Pas de transaction entre l'écriture dans l'event store et la publication du trigger — si la Cloud Function échoue définitivement, le read model peut diverger
  • Firestore facture à la lecture/écriture : un stream de 10 000 events coûte 10 000 opérations de lecture

Quand choisir :

  • Application mobile avec données temps réel (les listeners Firestore onSnapshot sont intégrés)
  • Équipe déjà familière avec Firebase
  • Pas besoin d'un historique long (< 1 000 events par agrégat)
  • Prototype rapide

Architecture 5 : AWS — DynamoDB + EventBridge + Lambda

Command Handler
      │ TransactWriteItems (concurrence optimiste)
      ▼
  DynamoDB (table Events)
      │
      │ DynamoDB Streams (changement capturé automatiquement)
      ▼
  EventBridge (bus d'events)
      │
      ├──► Lambda A → DynamoDB read model
      ├──► Lambda B → Elasticsearch
      └──► Lambda C → SQS → autre service

Cohérence : éventuelle — DynamoDB Streams + EventBridge introduisent un délai typique de quelques ms à quelques secondes.

Event store : DynamoDB avec la clé composite (stream_id, version) définie dans l'article 02.

DynamoDB Streams vers EventBridge :

typescript// Lambda déclenchée par DynamoDB Streams
import { DynamoDBStreamEvent } from "aws-lambda";
import { EventBridgeClient, PutEventsCommand } from "@aws-sdk/client-eventbridge";

const eb = new EventBridgeClient({});

export async function handler(event: DynamoDBStreamEvent): Promise<void> {
  const entries = event.Records
    .filter(r => r.eventName === "INSERT" && r.dynamodb?.NewImage)
    .map(r => {
      const item = unmarshall(r.dynamodb!.NewImage!);
      return {
        Source: "myapp.events",
        DetailType: item.event_type,
        Detail: JSON.stringify({
          streamId: item.stream_id,
          eventType: item.event_type,
          data: item.data,
          metadata: item.metadata ?? {},
          version: item.version,
          createdAt: item.created_at,
        }),
        EventBusName: "MyEventBus",
      };
    });

  if (entries.length > 0) {
    await eb.send(new PutEventsCommand({ Entries: entries }));
  }
}
typescript// Lambda de projection déclenchée par EventBridge
import { EventBridgeEvent } from "aws-lambda";

export async function handleOrderCreated(
  event: EventBridgeEvent<"OrderCreated", StoredEvent>
): Promise<void> {
  const stored = event.detail;

  // Upsert dans DynamoDB read model — idempotent
  await dynamoClient.send(new PutItemCommand({
    TableName: "OrderSummaries",
    Item: marshall({
      order_id: stored.data.orderId,
      customer_id: stored.data.customerId,
      status: "pending",
      total_cents: stored.data.totalCents,
      created_at: stored.createdAt,
    }),
    // Idempotence : ne pas écraser si déjà plus récent
    ConditionExpression: "attribute_not_exists(order_id) OR #v < :newVersion",
    ExpressionAttributeNames: { "#v": "version" },
    ExpressionAttributeValues: marshall({ ":newVersion": stored.version }),
  }));
}

Points de vigilance AWS :

  • DynamoDB Streams garantit at-least-once delivery — idempotence obligatoire dans les Lambdas
  • EventBridge garantit la livraison mais pas l'ordre entre streams — utiliser l'ordering key DynamoDB ou un timestamp pour dédupliquer
  • Lambda cold starts — préférer les Lambdas à mémoire réservée ou passer à ECS pour les projections critiques
  • Dead Letter Queue (DLQ) obligatoire pour capturer les events qui échouent définitivement

Quand choisir :

  • Déjà sur AWS
  • Besoin de scale très élevé (millions d'events/jour)
  • Architecture event-driven multi-services (microservices AWS natifs)
  • Budget pour l'infrastructure managée

Matrice de décision complète

Matrice de décision infrastructure CQRS

Questions à se poser

1. Quelle est ta contrainte d'infrastructure ? Déjà sur GCP → Architecture 3 ou 4. Déjà sur AWS → Architecture 5. Self-hosted ou liberté de choix → Architecture 1 ou 2.

2. Quel volume d'events par jour ? < 100 000 → Architecture 1 suffit. 100 000 à 10M → Architecture 2 ou 3.

10M → Architecture 5.

3. Quelle tolérance à l'eventual consistency ? Zéro tolérance (paiements, stocks) → Architecture 1 ou lecture depuis l'event store. Tolérance faible → Architecture 2 avec read-your-writes. Tolérance normale → Architectures 3, 4, 5.

Critère PG seul PG + réplicas GCP Pub/Sub Firebase AWS DynamoDB
Cohérence Forte Éventuelle Éventuelle Éventuelle Éventuelle
Idempotence Recommandée Recommandée Obligatoire Obligatoire Obligatoire
Complexité opé. Faible Moyenne Haute Moyenne Haute
Scalabilité Moyenne Haute Très haute Haute Très haute
Self-hosted Oui Oui Non Non Non
Cold start Non Non Oui (Functions) Oui Oui (Lambda)
Replay facile Oui Oui Moyen Difficile Moyen
Coût démarrage Faible Moyen Moyen Faible Moyen

Ce que le replay change selon l'infra

Le replay — reconstruire une projection depuis zéro — est la fonctionnalité qui différencie le plus les architectures.

PostgreSQL : trivial. SELECT * FROM events ORDER BY created_at, version et rejouer tout.

GCP Pub/Sub : Pub/Sub ne stocke pas les messages indéfiniment (7 jours max). Pour rejouer, lire depuis Cloud SQL (l'event store) et republier. La projection doit avoir un mode "batch rebuild".

Firebase Firestore : rejouer signifie lire toutes les sous-collections events de tous les streams — coûteux en lectures facturées. Envisager de stocker une copie dans Cloud Storage pour les rejeux.

DynamoDB : rejouer via une Scan de la table Events (coûteux) ou via DynamoDB Streams si la fenêtre de 24h est suffisante. Pour les rejeux historiques, exporter vers S3 d'abord.


Résumé de la série

Article Contenu
00 — Introduction Concepts, philosophie, quand utiliser
01 — Commands et events Handlers, domain events, validation
02 — Event store : interface et implémentations PostgreSQL, Firestore, DynamoDB
03 — Projections : synchrones et asynchrones In-process, Pub/Sub, idempotence, read-your-writes
04 — Replay et snapshots Rehydration, snapshots, versioning
05 — Projet réel Système complet TypeScript + Python
06 — Choisir son infrastructure PG, GCP, AWS, Firebase — matrice de décision

Sources

  • Young, G. (2010). CQRS Documents. cqrs.files.wordpress.com.
  • Richardson, C. (2019). Microservices Patterns. Manning.
  • Google Cloud. Pub/Sub ordering. cloud.google.com/pubsub/docs/ordering.
  • Google Cloud. Firestore transactions. cloud.google.com/firestore/docs/transactions.
  • Amazon. DynamoDB Streams. docs.aws.amazon.com/dynamodb/streams.
  • Amazon. EventBridge — Targets. docs.aws.amazon.com/eventbridge.
  • Firebase. Cloud Functions triggers. firebase.google.com/docs/functions.
  • Kleppmann, M. (2017). Designing Data-Intensive Applications. O'Reilly. (Chapitres 11-12 sur les streams et les systèmes distribués)

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