CQRS + Event Sourcing — 03 — Projections : synchrones et asynchrones

Construire des read models depuis les events : projections synchrones (in-process) vs asynchrones (Pub/Sub, EventBridge, Firestore Triggers). Idempotence, eventual consistency. TypeScript et Python.

03 — Projections : synchrones et asynchrones

Ce que tu vas apprendre

  • Ce qu'est une projection et comment elle se distingue de l'état de l'agrégat
  • Projections synchrones : in-process, cohérence forte, adapté à PostgreSQL seul
  • Projections asynchrones : via Pub/Sub / EventBridge / Firestore Triggers, eventual consistency
  • Idempotence : pourquoi elle devient critique en async
  • Le problème du read-your-writes et comment le résoudre
  • Exemples complets TypeScript et Python

Prérequis


Qu'est-ce qu'une projection ?

Une projection est une vue construite en appliquant des events sur un état initial vide. C'est la réponse à la question "à partir des events, construis-moi cette vue particulière."

L'état d'un agrégat est une projection — la reconstitution de son état courant. Mais une projection peut aussi être un read model optimisé pour l'affichage : une table dénormalisée, un document JSON, un agrégat de statistiques.

Events (append-only) ──► Projection A : OrderSummary (pour la liste des commandes)
                     ──► Projection B : CustomerOrderHistory (pour le profil client)
                     ──► Projection C : RevenueByDay (pour le dashboard)

Trois vues différentes, toutes construites depuis les mêmes events.

Les projections sont la raison pour laquelle l'Event Sourcing est puissant.

Ajouter une nouvelle vue ne nécessite pas de migrer la BDD ou de recalculer les données — il suffit de rejouer les events existants à travers une nouvelle projection.


Les read models PostgreSQL

Les read models vivent dans des tables séparées, optimisées pour la lecture. Là où l'event store est normalisé et append-only, les read models sont dénormalisés et orientés affichage.

sql-- Read model : résumé d'une commande (pour la liste)
CREATE TABLE order_summaries (
  order_id      VARCHAR(255) PRIMARY KEY,
  customer_id   VARCHAR(255) NOT NULL,
  customer_name VARCHAR(255),
  status        VARCHAR(50)  NOT NULL,
  total_cents   INTEGER      NOT NULL,
  item_count    INTEGER      NOT NULL DEFAULT 0,
  created_at    TIMESTAMPTZ  NOT NULL,
  updated_at    TIMESTAMPTZ  NOT NULL
);

-- Read model : historique client (pour le profil)
CREATE TABLE customer_order_history (
  id            UUID         PRIMARY KEY DEFAULT gen_random_uuid(),
  customer_id   VARCHAR(255) NOT NULL,
  order_id      VARCHAR(255) NOT NULL,
  status        VARCHAR(50)  NOT NULL,
  total_cents   INTEGER      NOT NULL,
  event_date    TIMESTAMPTZ  NOT NULL,
  UNIQUE (customer_id, order_id)
);

CREATE INDEX coh_customer_id ON customer_order_history (customer_id, event_date DESC);

-- Read model : revenus par jour (pour le dashboard)
CREATE TABLE revenue_by_day (
  date          DATE    PRIMARY KEY,
  total_cents   INTEGER NOT NULL DEFAULT 0,
  order_count   INTEGER NOT NULL DEFAULT 0
);

Implémenter une projection

Une projection est une fonction qui prend un event et met à jour le read model correspondant.

typescript// Interface générique d'une projection
interface Projection {
  handles: string[]; // Quels event_type cette projection traite
  apply(event: StoredEvent): Promise<void>;
}

// Projection : OrderSummary
class OrderSummaryProjection implements Projection {
  readonly handles = [
    "OrderCreated",
    "OrderConfirmed",
    "OrderShipped",
    "OrderCancelled",
  ];

  constructor(private readonly pool: Pool) {}

  async apply(event: StoredEvent): Promise<void> {
    switch (event.eventType) {
      case "OrderCreated":
        await this.onOrderCreated(event);
        break;
      case "OrderConfirmed":
        await this.onOrderConfirmed(event);
        break;
      case "OrderShipped":
        await this.onOrderShipped(event);
        break;
      case "OrderCancelled":
        await this.onOrderCancelled(event);
        break;
    }
  }

  private async onOrderCreated(event: StoredEvent): Promise<void> {
    const { orderId, customerId, items, totalCents } = event.data as any;

    await this.pool.query(
      `INSERT INTO order_summaries
         (order_id, customer_id, status, total_cents, item_count, created_at, updated_at)
       VALUES ($1, $2, 'pending', $3, $4, $5, $5)
       ON CONFLICT (order_id) DO NOTHING`,
      [orderId, customerId, totalCents, items.length, event.createdAt]
    );
  }

  private async onOrderConfirmed(event: StoredEvent): Promise<void> {
    const { orderId } = event.data as any;

    await this.pool.query(
      `UPDATE order_summaries
       SET status = 'confirmed', updated_at = $2
       WHERE order_id = $1`,
      [orderId, event.createdAt]
    );
  }

  private async onOrderShipped(event: StoredEvent): Promise<void> {
    const { orderId } = event.data as any;

    await this.pool.query(
      `UPDATE order_summaries
       SET status = 'shipped', updated_at = $2
       WHERE order_id = $1`,
      [orderId, event.createdAt]
    );
  }

  private async onOrderCancelled(event: StoredEvent): Promise<void> {
    const { orderId } = event.data as any;

    await this.pool.query(
      `UPDATE order_summaries
       SET status = 'cancelled', updated_at = $2
       WHERE order_id = $1`,
      [orderId, event.createdAt]
    );
  }
}

En Python :

pythonimport asyncpg
from typing import List

class OrderSummaryProjection:
    handles = ["OrderCreated", "OrderConfirmed", "OrderShipped", "OrderCancelled"]

    def __init__(self, pool: asyncpg.Pool) -> None:
        self._pool = pool

    async def apply(self, event: 'StoredEvent') -> None:
        handlers = {
            "OrderCreated": self._on_order_created,
            "OrderConfirmed": self._on_order_confirmed,
            "OrderShipped": self._on_order_shipped,
            "OrderCancelled": self._on_order_cancelled,
        }
        handler = handlers.get(event.event_type)
        if handler:
            await handler(event)

    async def _on_order_created(self, event: 'StoredEvent') -> None:
        data = event.data
        async with self._pool.acquire() as conn:
            await conn.execute(
                """
                INSERT INTO order_summaries
                  (order_id, customer_id, status, total_cents, item_count, created_at, updated_at)
                VALUES ($1, $2, 'pending', $3, $4, $5, $5)
                ON CONFLICT (order_id) DO NOTHING
                """,
                data["order_id"],
                data["customer_id"],
                data["total_cents"],
                len(data["items"]),
                event.created_at,
            )

    async def _on_order_confirmed(self, event: 'StoredEvent') -> None:
        async with self._pool.acquire() as conn:
            await conn.execute(
                "UPDATE order_summaries SET status = 'confirmed', updated_at = $2 WHERE order_id = $1",
                event.data["order_id"],
                event.created_at,
            )

    async def _on_order_shipped(self, event: 'StoredEvent') -> None:
        async with self._pool.acquire() as conn:
            await conn.execute(
                "UPDATE order_summaries SET status = 'shipped', updated_at = $2 WHERE order_id = $1",
                event.data["order_id"],
                event.created_at,
            )

    async def _on_order_cancelled(self, event: 'StoredEvent') -> None:
        async with self._pool.acquire() as conn:
            await conn.execute(
                "UPDATE order_summaries SET status = 'cancelled', updated_at = $2 WHERE order_id = $1",
                event.data["order_id"],
                event.created_at,
            )

Le ProjectionEngine : dispatcher d'events

Un ProjectionEngine centralise la mise à jour de toutes les projections après chaque command.

typescriptclass ProjectionEngine {
  private projections: Projection[] = [];

  register(projection: Projection): void {
    this.projections.push(projection);
  }

  async dispatch(events: StoredEvent[]): Promise<void> {
    for (const event of events) {
      const relevant = this.projections.filter(p =>
        p.handles.includes(event.eventType)
      );

      await Promise.all(relevant.map(p => p.apply(event)));
    }
  }
}

// Usage — dans le command handler, après append
const engine = new ProjectionEngine();
engine.register(new OrderSummaryProjection(pool));
engine.register(new CustomerOrderHistoryProjection(pool));
engine.register(new RevenueByDayProjection(pool));

// Après chaque append réussi
const newEvents = await eventStore.load(streamId, expectedVersion + 1);
await engine.dispatch(newEvents.events);

Plusieurs projections, mêmes events

Schéma projections multiple depuis un event store

Trois vues, un seul event

OrderCreated peut alimenter trois projections différentes simultanément :

OrderSummary — met à jour la ligne dans order_summaries avec le statut pending.

CustomerOrderHistory — insère une entrée dans customer_order_history pour que le profil client soit à jour.

RevenueByDay — incrémente le compteur de commandes du jour dans revenue_by_day.

Chaque projection est indépendante. En ajouter une nouvelle ne modifie pas les autres.

typescript// Projection : revenus par jour
class RevenueByDayProjection implements Projection {
  readonly handles = ["OrderConfirmed"]; // Seulement les commandes confirmées

  constructor(private readonly pool: Pool) {}

  async apply(event: StoredEvent): Promise<void> {
    const { totalCents } = event.data as any;
    const date = new Date(event.createdAt).toISOString().split("T")[0]; // "2026-06-01"

    await this.pool.query(
      `INSERT INTO revenue_by_day (date, total_cents, order_count)
       VALUES ($1, $2, 1)
       ON CONFLICT (date) DO UPDATE
       SET total_cents  = revenue_by_day.total_cents + EXCLUDED.total_cents,
           order_count  = revenue_by_day.order_count + 1`,
      [date, totalCents]
    );
  }
}

Lire depuis le read model

Les queries lisent directement dans les tables de read models — pas besoin de toucher l'event store.

typescriptclass OrderQueryHandler {
  constructor(private readonly pool: Pool) {}

  async getOrderSummaries(
    customerId: string,
    status?: string
  ): Promise<OrderSummaryDTO[]> {
    const { rows } = await this.pool.query(
      `SELECT order_id, customer_id, status, total_cents, item_count, created_at
       FROM order_summaries
       WHERE customer_id = $1
         AND ($2::varchar IS NULL OR status = $2)
       ORDER BY created_at DESC
       LIMIT 50`,
      [customerId, status ?? null]
    );

    return rows.map(r => ({
      orderId: r.order_id,
      customerId: r.customer_id,
      status: r.status,
      totalEuros: r.total_cents / 100,
      itemCount: r.item_count,
      createdAt: r.created_at,
    }));
  }

  async getRevenueSummary(from: string, to: string): Promise<RevenueSummaryDTO> {
    const { rows } = await this.pool.query(
      `SELECT SUM(total_cents) AS total, SUM(order_count) AS orders
       FROM revenue_by_day
       WHERE date BETWEEN $1 AND $2`,
      [from, to]
    );

    return {
      totalEuros: (rows[0].total ?? 0) / 100,
      orderCount: rows[0].orders ?? 0,
    };
  }
}

Reconstruire une projection

C'est là que l'Event Sourcing devient vraiment puissant. Si tu ajoutes une nouvelle projection — ou si tu corriges un bug dans une projection existante — tu peux la reconstruire depuis zéro en rejouant tous les events.

typescriptasync function rebuildProjection(
  eventStore: PostgresEventStore,
  projection: Projection,
  pool: Pool
): Promise<void> {
  // 1. Vider la table de la projection
  const tableMap: Record<string, string> = {
    OrderSummaryProjection: "order_summaries",
    CustomerOrderHistoryProjection: "customer_order_history",
    RevenueByDayProjection: "revenue_by_day",
  };
  const tableName = tableMap[projection.constructor.name];
  await pool.query(`TRUNCATE TABLE ${tableName}`);

  // 2. Lire tous les events pertinents par batches
  let offset = 0;
  const batchSize = 500;

  while (true) {
    const { rows } = await pool.query<StoredEvent>(
      `SELECT * FROM events
       WHERE event_type = ANY($1)
       ORDER BY created_at ASC, version ASC
       LIMIT $2 OFFSET $3`,
      [projection.handles, batchSize, offset]
    );

    if (rows.length === 0) break;

    for (const event of rows) {
      await projection.apply(event);
    }

    offset += rows.length;
    if (rows.length < batchSize) break;
  }

  console.log(`Projection reconstruite depuis ${offset} events`);
}

Projections synchrones vs asynchrones

Ce qu'on a vu jusqu'ici — le ProjectionEngine appelé dans le même process que le command handler — est une projection synchrone. Elle garantit que le read model est à jour dès que la commande répond.

C'est le bon choix pour un PostgreSQL self-hosted. Ce n'est pas toujours possible avec une infrastructure cloud.

Synchrone : cohérence forte

HTTP Request
    │
    ▼
Command Handler
    │ append(events)
    ▼
Event Store ──────────────────────────────┐
    │                                     │
    │ dispatch(events) [même process]     │
    ▼                                     │
ProjectionEngine                          │
    │ apply(event)                        │
    ▼                                     │
Read Model (même BDD ou même instance)   │
    │                                     │
    └──────────── HTTP Response ◄─────────┘

Commande réussit → Read model à jour immédiatement

Forces : simple, pas de latence entre écriture et lecture. Limite : le command handler attend que toutes les projections terminent — si une projection est lente, la commande est lente.

Asynchrone : eventual consistency

Avec Pub/Sub (GCP), EventBridge (AWS), ou les Firestore Triggers (Firebase), les projections tournent dans un processus séparé, déclenché par un message.

HTTP Request
    │
    ▼
Command Handler
    │ append(events)
    ▼
Event Store
    │ publish(event) → Pub/Sub / EventBridge / DynamoDB Streams
    │
    └──► HTTP Response ◄── immédiat, le read model n'est PAS encore mis à jour

        [quelques ms à quelques secondes plus tard]
        Pub/Sub → Cloud Function / Lambda
            │ apply(event)
            ▼
        Read Model mis à jour

Forces : le handler ne bloque pas sur les projections, scale indépendamment. Contrainte majeure : le read model est éventuellement cohérent — une lecture immédiatement après une commande peut retourner un état obsolète.


Implémentation asynchrone avec Pub/Sub (GCP)

Côté publisher — après l'append

typescriptimport { PubSub } from "@google-cloud/pubsub";

const pubsub = new PubSub();
const TOPIC = "domain-events";

async function publishEvents(events: StoredEvent[]): Promise<void> {
  const topic = pubsub.topic(TOPIC);

  await Promise.all(
    events.map(event =>
      topic.publishMessage({
        data: Buffer.from(JSON.stringify(event)),
        attributes: {
          eventType: event.eventType,
          streamId: event.streamId,
          version: String(event.version),
        },
      })
    )
  );
}

// Dans le command handler, après append
await eventStore.append(streamId, rawEvents, expectedVersion);
const { events: newEvents } = await eventStore.loadFrom(streamId, expectedVersion + 1);
await publishEvents(newEvents); // Fire-and-forget vers Pub/Sub

Côté subscriber — Cloud Function

typescript// Cloud Function déclenchée par Pub/Sub
import { CloudEvent } from "@google-cloud/functions-framework";

export async function handleDomainEvent(event: CloudEvent<Buffer>): Promise<void> {
  const storedEvent: StoredEvent = JSON.parse(
    Buffer.from(event.data as string, "base64").toString()
  );

  const projection = new OrderSummaryProjection(pool);

  if (projection.handles.includes(storedEvent.eventType)) {
    await projection.apply(storedEvent);
    // La projection doit être idempotente — Pub/Sub garantit at-least-once
  }
}
python# Cloud Function Python
import json
import base64
import functions_framework

@functions_framework.cloud_event
async def handle_domain_event(cloud_event):
    raw = base64.b64decode(cloud_event.data["message"]["data"]).decode()
    stored_event_data = json.loads(raw)

    stored_event = StoredEvent(**stored_event_data)
    projection = OrderSummaryProjection(pool)

    if stored_event.event_type in projection.handles:
        await projection.apply(stored_event)

L'idempotence : critique en async

Pub/Sub, EventBridge et DynamoDB Streams garantissent la livraison at-least-once — un event peut arriver deux fois en cas de retry. La projection doit produire le même résultat qu'elle reçoive l'event une ou plusieurs fois.

typescript// Sans idempotence — doublon possible
private async onOrderCreated(event: StoredEvent): Promise<void> {
  await this.pool.query(
    `INSERT INTO order_summaries (order_id, ...) VALUES ($1, ...)`,
    [event.data.orderId, ...]
  );
  // Si l'event arrive deux fois → erreur de clé dupliquée ou doublon silencieux
}

// Avec idempotence — ON CONFLICT DO NOTHING
private async onOrderCreated(event: StoredEvent): Promise<void> {
  await this.pool.query(
    `INSERT INTO order_summaries (order_id, ...)
     VALUES ($1, ...)
     ON CONFLICT (order_id) DO NOTHING`, // ← clé de l'idempotence
    [event.data.orderId, ...]
  );
}

// Avec idempotence — ON CONFLICT DO UPDATE (si on veut mettre à jour)
private async onOrderConfirmed(event: StoredEvent): Promise<void> {
  await this.pool.query(
    `UPDATE order_summaries
     SET status = 'confirmed', updated_at = $2
     WHERE order_id = $1
       AND status = 'pending'`, // ← n'applique que si pas déjà confirmé
    [event.data.orderId, event.createdAt]
  );
}

Règle : toute projection async doit être idempotente. Teste-la en envoyant le même event deux fois.


Le problème du read-your-writes

Avec des projections asynchrones, un utilisateur qui vient de passer une commande et recharge immédiatement la liste peut ne pas la voir — le read model n'est pas encore à jour.

Trois stratégies :

1. Retourner l'état depuis la command (pragmatique)

typescript// Le handler retourne les données directement — pas besoin de lire le read model
async createOrder(cmd: CreateOrderCommand): Promise<OrderDTO> {
  // ... append events ...
  return {
    orderId,
    status: "pending",
    totalEuros: totalCents / 100,
    // Construit depuis les données de la command — pas depuis la BDD
  };
}

2. Version header — le client attend la bonne version

typescript// Le handler retourne la version écrite
// Le client passe cette version dans la requête suivante
// Le read model répond seulement s'il a atteint cette version

async function getOrder(orderId: string, minVersion?: number): Promise<OrderSummaryDTO> {
  if (minVersion !== undefined) {
    // Attendre que le read model atteigne la version demandée
    // (polling avec timeout)
    for (let i = 0; i < 10; i++) {
      const order = await queryOrder(orderId);
      if (order && order.version >= minVersion) return order;
      await sleep(100);
    }
    throw new Error("Read model non à jour après timeout");
  }
  return queryOrder(orderId);
}

3. Lire depuis l'event store pour les lectures critiques

typescript// Pour les cas où la fraîcheur est critique (validation, paiement)
// Lire depuis l'event store et reconstituer l'état
async getOrderForPayment(orderId: string): Promise<Order> {
  const { events } = await eventStore.load(`order-${orderId}`);
  return OrderAggregate.rehydrate(events); // État garanti à jour
}

// Pour les affichages non-critiques (liste, dashboard)
async listOrders(customerId: string): Promise<OrderSummaryDTO[]> {
  return queryFromReadModel(customerId); // Éventuellement cohérent — acceptable
}

Résumé

Synchrone Asynchrone (Pub/Sub / EventBridge)
Cohérence Forte — read model immédiatement à jour Éventuelle — délai de quelques ms à s
Performance Command attend les projections Command répond immédiatement
Idempotence Recommandée Obligatoire
Read-your-writes Pas de problème Stratégie nécessaire
Infra Self-hosted PostgreSQL GCP, AWS, Firebase
Scalabilité Limitée par le process Indépendante

L'event store reste la source de vérité. Les read models sont des caches — synchrones ou asynchrones selon l'infrastructure, détruits et reconstruits sans perte de données.

Étape suivante : 04 — Replay et snapshots — reconstituer l'état d'un agrégat et optimiser le chargement avec des snapshots.


Sources

  • Young, G. (2010). CQRS Documents. cqrs.files.wordpress.com.
  • Fowler, M. (2005). Event Sourcing. martinfowler.com.
  • Richardson, C. (2019). Microservices Patterns. Manning.
  • Google Cloud. Pub/Sub — Subscriber guide. cloud.google.com/pubsub/docs.
  • Amazon. DynamoDB Streams and AWS Lambda. 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.