Architecture & conception

CQRS et Event Sourcing

Séparation lecture-écriture, MediatR, journal d'événements.

Vérifié en septembre 2026 · .NET 10 · environ 16 min

CQRS et event sourcing sont deux motifs distincts qu'on présente souvent ensemble. CQRS sépare ce qui modifie les données de ce qui les lit ; il sert quand les écrans tordent le modèle métier pour servir l'affichage. L'event sourcing enregistre chaque changement comme un événement et fait de ce journal la source de vérité ; il sert quand l'histoire des changements a une valeur propre — audit, questions posées après coup — et coûte cher partout ailleurs. Le fil rouge reprend la commande et son plafond du cours sur le DDD tactique, simplifiés : les montants sont des decimal, sans value object. Les exemples sont les fichiers d'un même projet console .NET 10 ; ceux qui ont une méthode Executer affichent ce que disent leurs commentaires. Celui de MediatR, seul, n'a été ni compilé contre le paquet ni exécuté.

De CQS à CQRS

Bertrand Meyer a posé le principe de séparation commande-requête, que Martin Fowler résume en 2005 : une requête renvoie un résultat sans changer l'état observable du système, une commande change l'état sans rien renvoyer. Le principe vaut pour chaque méthode. Greg Young, le premier que Fowler ait entendu parler de CQRS, le résume dans un billet du 16 février 2010 : créer deux objets là où il n'y en avait qu'un, l'un pour les commandes, l'autre pour les requêtes. Un motif très simple, ajoute-t-il, qui vaut par les architectures qu'il rend possibles.

Ce qu'il résout se voit sur un modèle unique. Dans ses « CQRS Documents », mis en ligne en novembre 2010, Young en relève les symptômes : des repositories chargés de méthodes de lecture, avec pagination et tri ; des accesseurs qui exposent l'état interne des agrégats pour fabriquer des DTO ; plusieurs agrégats chargés pour remplir un écran. Le code suivant les réunit, et c'est ce qui cloche : l'agrégat et son repository se plient à l'affichage, et vingt agrégats entiers sont chargés pour montrer quatre colonnes.

using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;

namespace Ventes.Melange;

// Un seul modele pour tout : il decide, et il sert aussi les ecrans.
public sealed class Commande(Guid id, Guid client, decimal plafond)
{
    private readonly List<Ligne> _lignes = [];

    public Guid Id { get; } = id;

    public Guid Client { get; } = client;

    public bool EstValidee { get; private set; }

    // Exposes pour fabriquer les DTO : l'ecran lit l'etat interne de l'agregat.
    public IReadOnlyList<Ligne> Lignes => _lignes;

    public decimal Total => _lignes.Sum(l => l.Quantite * l.PrixUnitaire);

    public void AjouterLigne(string reference, int quantite, decimal prixUnitaire)
    {
        if (Total + quantite * prixUnitaire > plafond)
        {
            throw new InvalidOperationException($"Plafond de {plafond} depasse.");
        }

        _lignes.Add(new Ligne(reference, quantite, prixUnitaire));
    }

    public void Valider() => EstValidee = true;
}

public sealed record Ligne(string Reference, int Quantite, decimal PrixUnitaire);

public sealed record ResumeCommande(Guid Id, int NombreDeLignes, decimal Total, bool EstValidee);

// Le repository de l'agregat grossit au rythme des ecrans : pagination, tri, filtres.
public interface ICommandeRepository
{
    Task<Commande?> Charger(Guid id, CancellationToken annulation);

    Task Enregistrer(Commande commande, CancellationToken annulation);

    Task<IReadOnlyList<Commande>> ListerParClient(
        Guid client, int page, int taille, string tri, CancellationToken annulation);
}

public sealed class ServiceCommandes(ICommandeRepository commandes)
{
    public async Task Valider(Guid id, CancellationToken annulation)
    {
        var commande = await commandes.Charger(id, annulation)
            ?? throw new InvalidOperationException("Commande introuvable.");
        commande.Valider();
        await commandes.Enregistrer(commande, annulation);
    }

    // Pour afficher quatre colonnes, vingt agregats entiers, lignes comprises.
    public async Task<IReadOnlyList<ResumeCommande>> CommandesDuClient(
        Guid client, int page, CancellationToken annulation)
    {
        var liste = await commandes.ListerParClient(client, page, 20, "date", annulation);
        return liste
            .Select(c => new ResumeCommande(c.Id, c.Lignes.Count, c.Total, c.EstValidee))
            .ToList();
    }
}

Le côté écriture garde l'agrégat, qui décide et n'expose plus rien ; son repository se réduit à charger par l'identifiant et à enregistrer. Le côté lecture est ce que Young appelle une couche de lecture mince : elle interroge la base directement et projette un résultat taillé pour l'écran, sans passer par le domaine.

using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;

namespace Ventes.Separe;

// Cote ecriture : l'agregat decide, et n'a plus rien a montrer.
public sealed class Commande(Guid id, decimal plafond)
{
    private readonly List<(int Quantite, decimal PrixUnitaire)> _lignes = [];
    private bool _validee;

    public Guid Id { get; } = id;

    public void AjouterLigne(int quantite, decimal prixUnitaire)
    {
        if (_validee)
        {
            throw new InvalidOperationException("Une commande validee ne se modifie plus.");
        }

        if (_lignes.Sum(l => l.Quantite * l.PrixUnitaire) + quantite * prixUnitaire > plafond)
        {
            throw new InvalidOperationException($"Plafond de {plafond} depasse.");
        }

        _lignes.Add((quantite, prixUnitaire));
    }

    public void Valider() => _validee = true;
}

// Le repository ne sert que l'ecriture : charger par l'identifiant, enregistrer.
public interface ICommandes
{
    Task<Commande?> Charger(Guid id, CancellationToken annulation);

    Task Enregistrer(Commande commande, CancellationToken annulation);
}

// Une commande, a l'imperatif : une intention, que le domaine peut refuser.
public sealed record ValiderCommande(Guid Commande);

public sealed class ValiderCommandeGestionnaire(ICommandes commandes)
{
    public async Task Executer(ValiderCommande commande, CancellationToken annulation)
    {
        var agregat = await commandes.Charger(commande.Commande, annulation)
            ?? throw new InvalidOperationException("Commande introuvable.");
        agregat.Valider();
        await commandes.Enregistrer(agregat, annulation);
    }
}

// Cote lecture : une requete, et un resultat taille pour l'ecran.
public sealed record CommandesDuClient(Guid Client, int Page);

public sealed record ResumeCommande(Guid Id, int NombreDeLignes, decimal Total, bool EstValidee);

// Implementee par une requete SQL qui lit et projette ces quatre colonnes,
// sans charger aucun agregat : SELECT ... GROUP BY ... ORDER BY ... LIMIT 20.
public interface ILectureCommandes
{
    Task<IReadOnlyList<ResumeCommande>> Executer(
        CommandesDuClient requete, CancellationToken annulation);
}

Une requête ne change rien : elle peut être répétée, mise en cache ou servie par une réplique. Le cours sur le DDD tactique appelait déjà cela la forme la plus modeste de CQRS : écrire passe par l'agrégat, lire n'y est pas obligé.

Ce que CQRS n'impose pas

Young l'écrit dans son billet : CQRS n'est ni la cohérence à terme, ni les événements, ni la messagerie, ni même des modèles séparés pour lire et écrire, ni l'event sourcing. Fowler, en 2011, laisse les deux modèles partager la même base, qui sert alors de lien entre eux, ou utiliser deux bases, celle des requêtes devenant une base de reporting tenue à jour. En 2017, il précise que CQRS, strictement, se passe d'événements. La définition minimale de Young s'arrête aux objets ; la séparation va ensuite plus loin par degrés, et chacun a son prix.

DegréCe qui est séparéCe que ça coûte
Objets séparésles classes ; une base, une transactionpresque rien
Modèles séparés, même basevues ou tables dénormalisées pour la lectureune projection de plus à maintenir
Bases séparéesle stockage, synchronisé par des événements ou une réplicationla cohérence à terme, une infrastructure de plus

Aucun de ces degrés ne demande de bibliothèque spécialisée. Le répartiteur suivant tient en deux méthodes : il trouve le gestionnaire dans le conteneur de Microsoft.Extensions.DependencyInjection et l'appelle.

using System;
using System.Collections.Generic;
using System.Linq;
using System.Reflection;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;

namespace Ventes.Repartition;

public interface ICommande;

public interface IRequete<TResultat>;

public interface IGestionnaireCommande<in TCommande>
    where TCommande : ICommande
{
    Task Executer(TCommande commande, CancellationToken annulation);
}

public interface IGestionnaireRequete<in TRequete, TResultat>
    where TRequete : IRequete<TResultat>
{
    Task<TResultat> Executer(TRequete requete, CancellationToken annulation);
}

// Tout le « bus » : trouver le gestionnaire dans le conteneur, et l'appeler.
public sealed class Repartiteur(IServiceProvider services)
{
    public async Task Executer<TCommande>(
        TCommande commande, CancellationToken annulation = default)
        where TCommande : ICommande
    {
        var gestionnaire = services.GetRequiredService<IGestionnaireCommande<TCommande>>();

        // Le point unique ou poser ce qui vaut pour toutes les commandes :
        // journalisation, validation, transaction.
        Console.WriteLine($"-> {typeof(TCommande).Name}");
        await gestionnaire.Executer(commande, annulation);
    }

    public Task<TResultat> Demander<TResultat>(
        IRequete<TResultat> requete, CancellationToken annulation = default)
    {
        // Le type exact de la requete n'est connu qu'a l'execution :
        // on construit celui de son gestionnaire, puis on l'appelle.
        var type = typeof(IGestionnaireRequete<,>)
            .MakeGenericType(requete.GetType(), typeof(TResultat));
        var gestionnaire = services.GetRequiredService(type);
        var executer = type.GetMethod("Executer")!;
        return (Task<TResultat>)executer.Invoke(
            gestionnaire, BindingFlags.DoNotWrapExceptions, null, [requete, annulation], null)!;
    }
}

// La base, partagee par les deux cotes : c'est elle qui les fait communiquer.
public sealed class BaseEnMemoire
{
    public Dictionary<Guid, CommandeEnBase> Commandes { get; } = [];
}

public sealed class CommandeEnBase(decimal total)
{
    public decimal Total { get; } = total;

    public bool Validee { get; set; }
}

public sealed record ValiderCommande(Guid Commande) : ICommande;

public sealed record CommandesAValider : IRequete<IReadOnlyList<ResumeCommande>>;

public sealed record ResumeCommande(Guid Id, decimal Total);

public sealed class ValiderCommandeGestionnaire(BaseEnMemoire donnees)
    : IGestionnaireCommande<ValiderCommande>
{
    public Task Executer(ValiderCommande commande, CancellationToken annulation)
    {
        var enregistrement = donnees.Commandes[commande.Commande];
        if (enregistrement.Total == 0)
        {
            throw new InvalidOperationException("Une commande vide ne se valide pas.");
        }

        enregistrement.Validee = true;
        return Task.CompletedTask;
    }
}

public sealed class CommandesAValiderGestionnaire(BaseEnMemoire donnees)
    : IGestionnaireRequete<CommandesAValider, IReadOnlyList<ResumeCommande>>
{
    public Task<IReadOnlyList<ResumeCommande>> Executer(
        CommandesAValider requete, CancellationToken annulation)
    {
        IReadOnlyList<ResumeCommande> resultat = donnees.Commandes
            .Where(c => !c.Value.Validee)
            .Select(c => new ResumeCommande(c.Key, c.Value.Total))
            .ToList();
        return Task.FromResult(resultat);
    }
}

static class DemonstrationRepartition
{
    public static async Task Executer()
    {
        using var services = new ServiceCollection()
            .AddSingleton<BaseEnMemoire>()
            .AddTransient<IGestionnaireCommande<ValiderCommande>, ValiderCommandeGestionnaire>()
            .AddTransient<
                IGestionnaireRequete<CommandesAValider, IReadOnlyList<ResumeCommande>>,
                CommandesAValiderGestionnaire>()
            .AddTransient<Repartiteur>()
            .BuildServiceProvider();

        var donnees = services.GetRequiredService<BaseEnMemoire>();
        var premiere = Guid.NewGuid();
        donnees.Commandes[premiere] = new CommandeEnBase(300m);
        donnees.Commandes[Guid.NewGuid()] = new CommandeEnBase(120m);

        var repartiteur = services.GetRequiredService<Repartiteur>();
        await repartiteur.Executer(new ValiderCommande(premiere)); // -> ValiderCommande

        var aValider = await repartiteur.Demander(new CommandesAValider());
        Console.WriteLine($"{aValider.Count} a valider, total {aValider[0].Total}");
        // 1 a valider, total 120
    }
}

Pour une commande, le type est connu à la compilation, et le conteneur suffit. Pour une requête, l'appelant ne fournit que le type du résultat : celui du gestionnaire ne se construit qu'à l'exécution, d'où la réflexion. L'option DoNotWrapExceptions y est nécessaire : sans elle, une exception levée par le gestionnaire avant qu'il rende sa tâche arriverait enveloppée dans une TargetInvocationException. Le répartiteur reste facultatif : un contrôleur qui reçoit son gestionnaire par injection et l'appelle fait encore du CQRS.

MediatR, et sa licence

MediatR, de Jimmy Bogard, se décrit comme une implémentation simple du motif médiateur, une messagerie interne au processus. C'est la mécanique du répartiteur, en plus complet : requêtes avec ou sans réponse, notifications à plusieurs gestionnaires, flux, et comportements de pipeline, IPipelineBehavior, qui enveloppent chaque requête pour la journaliser ou la valider. Ses interfaces ne distinguent pas une commande d'une requête : les deux sont des IRequest, et la séparation reste une affaire de nommage.

using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using Microsoft.Extensions.DependencyInjection;
using Ventes.Repartition; // BaseEnMemoire et CommandeEnBase, de l'exemple precedent

namespace Ventes.AvecMediatR;

// Les memes messages, marques par les interfaces de MediatR.
public sealed record ValiderCommande(Guid Commande) : IRequest;

public sealed record CommandesAValider : IRequest<IReadOnlyList<ResumeCommande>>;

public sealed record ResumeCommande(Guid Id, decimal Total);

public sealed class ValiderCommandeGestionnaire(BaseEnMemoire donnees)
    : IRequestHandler<ValiderCommande>
{
    public Task Handle(ValiderCommande commande, CancellationToken annulation)
    {
        var enregistrement = donnees.Commandes[commande.Commande];
        if (enregistrement.Total == 0)
        {
            throw new InvalidOperationException("Une commande vide ne se valide pas.");
        }

        enregistrement.Validee = true;
        return Task.CompletedTask;
    }
}

public sealed class CommandesAValiderGestionnaire(BaseEnMemoire donnees)
    : IRequestHandler<CommandesAValider, IReadOnlyList<ResumeCommande>>
{
    public Task<IReadOnlyList<ResumeCommande>> Handle(
        CommandesAValider requete, CancellationToken annulation)
    {
        IReadOnlyList<ResumeCommande> resultat = donnees.Commandes
            .Where(c => !c.Value.Validee)
            .Select(c => new ResumeCommande(c.Key, c.Value.Total))
            .ToList();
        return Task.FromResult(resultat);
    }
}

static class DemonstrationMediatR
{
    public static async Task Executer()
    {
        // Les gestionnaires sont trouves par balayage de l'assembly. La cle de
        // licence se lit dans MEDIATR_LICENSE_KEY quand le code ne la fixe pas.
        // AddLogging vient du paquet Microsoft.Extensions.Logging, que MediatR
        // n'apporte pas : il ne depend que de Logging.Abstractions.
        using var services = new ServiceCollection()
            .AddSingleton<BaseEnMemoire>()
            .AddLogging()
            .AddMediatR(cfg => cfg.RegisterServicesFromAssemblyContaining<ValiderCommande>())
            .BuildServiceProvider();

        var donnees = services.GetRequiredService<BaseEnMemoire>();
        var premiere = Guid.NewGuid();
        donnees.Commandes[premiere] = new CommandeEnBase(300m);
        donnees.Commandes[Guid.NewGuid()] = new CommandeEnBase(120m);

        var mediateur = services.GetRequiredService<ISender>();
        await mediateur.Send(new ValiderCommande(premiere));
        var aValider = await mediateur.Send(new CommandesAValider());
        Console.WriteLine($"{aValider.Count} a valider, total {aValider[0].Total}");
    }
}

Cet exemple suit l'API de la 14.2.0, du 2 juillet 2026, d'après le README et le wiki. L'appel à AddLogging n'est pas décoratif : en 14.0.0, un conteneur sans journalisation faisait échouer la résolution du médiateur, et depuis la 14.1.0 un message explicite demande cet appel. La licence, elle, a changé en 2025. Jusqu'à la 12.5.0 du 1er avril 2025, MediatR était sous Apache 2.0 ; le lendemain, Bogard annonçait son passage au commercial. Depuis la 13.0.0 du 2 juillet 2025, le fichier de licence du dépôt laisse le choix entre la Reciprocal Public License 1.5 et le contrat commercial de Lucky Penny Software. Ce contrat, lu le 26 septembre 2026 dans sa version 2.0, réserve la licence communautaire, gratuite selon l'annonce de Bogard, aux organisations qui remplissent toutes ses conditions à la fois : entre autres, un chiffre d'affaires annuel brut sous 5 millions de dollars, agrégé avec celui de la maison mère, et un plafond de capitaux extérieurs. Il laisse libre tout usage hors production. Selon l'annonce, la bibliothèque fonctionne sans clé et écrit des avertissements dans les journaux ; le README donne la clé par cfg.LicenseKey ou par la variable MEDIATR_LICENSE_KEY.

Ce qui suit est une lecture de ces textes, pas un avis juridique. La RPL 1.5 compte comme déploiement tout usage interne à une entreprise, hors recherche et usage personnel, et oblige qui déploie des extensions à en publier le source ; le fichier de licence de Lucky Penny en tire que le logiciel construit avec la bibliothèque se publie sous ces termes, faute de contrat commercial. Une entreprise qui ne remplit pas l'une des conditions et ne publie pas son code a donc, en production, notamment ces voies : payer, rester sur la 12.5.0, que le contrat laisse utilisable sous Apache 2.0 indéfiniment mais sans support ni correctif de sécurité, ou garder son propre répartiteur.

Event sourcing : le journal comme source de vérité

Fowler le définit en 2005 : capturer tous les changements de l'état d'une application sous la forme d'une suite d'événements. Young précise dans son billet ce qu'il entend par là : stocker l'état courant comme une série d'événements, et reconstruire l'état en rejouant cette série. L'état n'est plus enregistré, il se déduit. Un grand livre comptable fonctionne ainsi, et c'est l'exemple de Young : le solde n'est qu'une commodité, qu'on peut à tout moment recalculer en additionnant les écritures depuis l'ouverture du compte.

L'événement est le pendant de la commande. Une commande demande, à l'impératif, et peut être refusée ; un événement constate, au passé, et ne se refuse plus — le domaine, écrit Young, n'a pas de machine à remonter le temps. D'où la forme de l'agrégat : chaque méthode de commande vérifie les règles, puis produit un événement, et une seule méthode, Appliquer, change l'état, que l'événement soit neuf ou rejoué.

using System;
using System.Collections.Generic;
using System.Linq;

namespace Ventes.Evenements;

// Des faits, au passe, immuables. Chacun nomme le flux auquel il appartient.
public interface IEvenement
{
    Guid Commande { get; }
}

public sealed record CommandeOuverte(Guid Commande, Guid Client, decimal Plafond) : IEvenement;

public sealed record LigneAjoutee(
    Guid Commande, string Reference, int Quantite, decimal PrixUnitaire) : IEvenement;

public sealed record LigneRetiree(Guid Commande, string Reference) : IEvenement;

public sealed record CommandeValidee(Guid Commande) : IEvenement;

// La seconde partie de Commande est dans l'exemple des instantanes.
public sealed partial class Commande
{
    private readonly Dictionary<string, (int Quantite, decimal PrixUnitaire)> _lignes = [];
    private readonly List<IEvenement> _nouveaux = [];
    private decimal _plafond;
    private bool _validee;

    private Commande() { }

    public Guid Id { get; private set; }

    // Le nombre d'evenements deja enregistres : la version sur laquelle on decide.
    public int Version { get; private set; }

    public decimal Total => _lignes.Values.Sum(l => l.Quantite * l.PrixUnitaire);

    public IReadOnlyList<IEvenement> NouveauxEvenements => _nouveaux;

    public static Commande Ouvrir(Guid id, Guid client, decimal plafond)
    {
        ArgumentOutOfRangeException.ThrowIfNegativeOrZero(plafond);
        var commande = new Commande();
        commande.Produire(new CommandeOuverte(id, client, plafond));
        return commande;
    }

    // Rejouer n'est pas redecider : aucune regle, aucun effet de bord.
    public static Commande Reconstituer(IEnumerable<IEvenement> flux)
    {
        var commande = new Commande();
        foreach (var evenement in flux)
        {
            commande.Appliquer(evenement);
            commande.Version++;
        }

        return commande;
    }

    // Les methodes de commande decident : elles verifient, puis produisent un fait.
    public void AjouterLigne(string reference, int quantite, decimal prixUnitaire)
    {
        VerifierModifiable();
        ArgumentOutOfRangeException.ThrowIfNegativeOrZero(quantite);
        ArgumentOutOfRangeException.ThrowIfNegative(prixUnitaire);
        if (_lignes.ContainsKey(reference))
        {
            throw new InvalidOperationException($"{reference} figure deja dans la commande.");
        }

        if (Total + quantite * prixUnitaire > _plafond)
        {
            throw new InvalidOperationException($"Plafond de {_plafond} depasse.");
        }

        Produire(new LigneAjoutee(Id, reference, quantite, prixUnitaire));
    }

    // Rien ne s'efface du journal : retirer une ligne est un fait de plus.
    public void RetirerLigne(string reference)
    {
        VerifierModifiable();
        if (!_lignes.ContainsKey(reference))
        {
            throw new InvalidOperationException($"{reference} ne figure pas dans la commande.");
        }

        Produire(new LigneRetiree(Id, reference));
    }

    public void Valider()
    {
        VerifierModifiable();
        if (_lignes.Count == 0)
        {
            throw new InvalidOperationException("Une commande vide ne se valide pas.");
        }

        Produire(new CommandeValidee(Id));
    }

    private void VerifierModifiable()
    {
        if (_validee)
        {
            throw new InvalidOperationException("Une commande validee ne se modifie plus.");
        }
    }

    private void Produire(IEvenement evenement)
    {
        Appliquer(evenement);
        _nouveaux.Add(evenement);
    }

    // Le seul endroit qui change l'etat, que l'evenement soit neuf ou rejoue.
    private void Appliquer(IEvenement evenement)
    {
        switch (evenement)
        {
            case CommandeOuverte e:
                Id = e.Commande;
                _plafond = e.Plafond;
                break;
            case LigneAjoutee e:
                _lignes[e.Reference] = (e.Quantite, e.PrixUnitaire);
                break;
            case LigneRetiree e:
                _lignes.Remove(e.Reference);
                break;
            case CommandeValidee:
                _validee = true;
                break;
        }
    }
}

static class DemonstrationRejeu
{
    public static void Executer()
    {
        var commande = Commande.Ouvrir(Guid.NewGuid(), Guid.NewGuid(), 500m);
        commande.AjouterLigne("ECRAN-27", 1, 300m);
        commande.AjouterLigne("CLAVIER", 2, 40m);
        commande.RetirerLigne("CLAVIER");

        var journal = commande.NouveauxEvenements.ToList();
        Console.WriteLine(string.Join(", ", journal.Select(e => e.GetType().Name)));
        // CommandeOuverte, LigneAjoutee, LigneAjoutee, LigneRetiree

        var rejouee = Commande.Reconstituer(journal);
        Console.WriteLine($"{rejouee.Total} en version {rejouee.Version}"); // 300 en version 4

        try
        {
            rejouee.AjouterLigne("DOCK", 1, 250m);
        }
        catch (InvalidOperationException ex)
        {
            Console.WriteLine(ex.Message); // Plafond de 500 depasse.
        }
    }
}

Reconstituer n'appelle aucune règle. Les faits ont été vérifiés quand ils ont été décidés ; les revérifier au rejeu laisserait une règle nouvelle réécrire le passé. Rejouer n'est pas non plus le moment d'agir : un e-mail part d'un événement neuf, jamais d'une reconstitution, et Fowler consacre une partie de son article aux passerelles vers les systèmes externes, qui doivent savoir qu'on rejoue. Rien, enfin, ne s'efface. Retirer le clavier ajoute un LigneRetiree : Young parle de transaction d'annulation, qui ramène à l'état voulu en gardant la trace du précédent.

Le magasin d'événements et la version attendue

Young décrit un magasin minimal sur une base relationnelle : une table des événements (identifiant de l'agrégat, données sérialisées, version), une table des agrégats avec leur version courante, et deux opérations seulement : lire les événements d'un agrégat dans l'ordre de leurs versions, et en ajouter. L'ajout porte la version sur laquelle l'appelant a décidé. Le magasin suivant ajoute sans rien vérifier, et c'est ce qui cloche : deux ajouts de 150, chacun légitime au regard du total qu'il a lu, donnent une commande de 600 pour un plafond de 500.

using System;
using System.Collections.Generic;
using System.Threading;

namespace Ventes.Evenements;

// Un magasin qui ajoute ce qu'on lui donne, sans demander sur quelle version on a decide.
public sealed class MagasinSansVersion
{
    private readonly Lock _verrou = new();
    private readonly Dictionary<Guid, List<IEvenement>> _flux = [];

    public IReadOnlyList<IEvenement> Lire(Guid flux)
    {
        lock (_verrou)
        {
            return _flux.TryGetValue(flux, out var evenements) ? [.. evenements] : [];
        }
    }

    public void Ajouter(Guid flux, IReadOnlyList<IEvenement> evenements)
    {
        lock (_verrou)
        {
            if (!_flux.TryGetValue(flux, out var existants))
            {
                _flux[flux] = existants = [];
            }

            existants.AddRange(evenements);
        }
    }
}

static class DemonstrationSansVersion
{
    public static void Executer()
    {
        var magasin = new MagasinSansVersion();
        var id = Guid.NewGuid();
        var commande = Commande.Ouvrir(id, Guid.NewGuid(), 500m);
        commande.AjouterLigne("ECRAN-27", 1, 300m);
        magasin.Ajouter(id, commande.NouveauxEvenements);

        // Deux requetes simultanees chargent la meme commande...
        var a = Commande.Reconstituer(magasin.Lire(id));
        var b = Commande.Reconstituer(magasin.Lire(id));

        // ... et chacune verifie le plafond sur ce qu'elle a lu : 300 + 150 <= 500.
        a.AjouterLigne("DOCK", 1, 150m);
        b.AjouterLigne("CASQUE", 1, 150m);
        magasin.Ajouter(id, a.NouveauxEvenements);
        magasin.Ajouter(id, b.NouveauxEvenements);

        Console.WriteLine(Commande.Reconstituer(magasin.Lire(id)).Total); // 600
    }
}

Chaque requête a vérifié la règle sur l'état qu'elle avait lu, et celui de la seconde était périmé quand elle a écrit ; le rejeu, qui ne revérifie rien, accepte le résultat. Le verrou n'y change rien : il protège la liste, pas la décision prise avant. Le magasin juste compare la version attendue à la version réelle et refuse l'ajout en cas d'écart :

using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;

namespace Ventes.Evenements;

public sealed class ConflitDeVersion(int attendue, int actuelle)
    : Exception($"Version attendue {attendue}, version actuelle {actuelle}.");

// Deux operations : lire un flux, y ajouter des evenements a une version donnee.
public sealed class MagasinEnMemoire
{
    private readonly Lock _verrou = new();
    private readonly Dictionary<Guid, List<IEvenement>> _flux = [];

    // Tous les flux, dans l'ordre d'ecriture : ce que lisent les projections.
    private readonly List<IEvenement> _journal = [];

    public IReadOnlyList<IEvenement> Lire(Guid flux, int apresVersion = 0)
    {
        lock (_verrou)
        {
            return _flux.TryGetValue(flux, out var evenements)
                ? [.. evenements.Skip(apresVersion)]
                : [];
        }
    }

    public IReadOnlyList<IEvenement> LireDepuis(int position)
    {
        lock (_verrou)
        {
            return [.. _journal.Skip(position)];
        }
    }

    public void Ajouter(Guid flux, int versionAttendue, IReadOnlyList<IEvenement> evenements)
    {
        lock (_verrou)
        {
            var actuelle = _flux.TryGetValue(flux, out var existants) ? existants.Count : 0;

            // Verifier et ajouter d'un seul geste : c'est tout ce que le magasin garantit.
            if (actuelle != versionAttendue)
            {
                throw new ConflitDeVersion(versionAttendue, actuelle);
            }

            if (existants is null)
            {
                _flux[flux] = existants = [];
            }

            existants.AddRange(evenements);
            _journal.AddRange(evenements);
        }
    }
}

static class DemonstrationVersionAttendue
{
    public static void Executer()
    {
        var magasin = new MagasinEnMemoire();
        var id = Guid.NewGuid();
        var commande = Commande.Ouvrir(id, Guid.NewGuid(), 500m);
        commande.AjouterLigne("ECRAN-27", 1, 300m);
        magasin.Ajouter(id, commande.Version, commande.NouveauxEvenements);

        var a = Commande.Reconstituer(magasin.Lire(id));
        var b = Commande.Reconstituer(magasin.Lire(id));
        a.AjouterLigne("DOCK", 1, 150m);
        b.AjouterLigne("CASQUE", 1, 150m);
        magasin.Ajouter(id, a.Version, a.NouveauxEvenements);

        try
        {
            magasin.Ajouter(id, b.Version, b.NouveauxEvenements);
        }
        catch (ConflitDeVersion ex)
        {
            Console.WriteLine(ex.Message); // Version attendue 2, version actuelle 3.

            // On recharge et on redecide sur l'etat reel : la regle refuse, cette fois.
            var rechargee = Commande.Reconstituer(magasin.Lire(id));
            try
            {
                rechargee.AjouterLigne("CASQUE", 1, 150m);
            }
            catch (InvalidOperationException refus)
            {
                Console.WriteLine(refus.Message); // Plafond de 500 depasse.
            }
        }

        Console.WriteLine(Commande.Reconstituer(magasin.Lire(id)).Total); // 450
    }
}

La seconde requête reçoit un conflit, recharge l'état réel, redécide, et la règle refuse. C'est la concurrence optimiste : rien n'est verrouillé pendant la décision, seul l'ajout est atomique, et un conflit coûte une relecture. Les versions ne sont uniques et consécutives qu'à l'intérieur d'un flux, parce que l'agrégat est l'unité de cohérence.

Projections et instantanés

Le journal ne répond qu'à une question : quels événements pour cet agrégat. Young le souligne, on ne peut pas lui demander tous les utilisateurs prénommés Greg, faute d'état courant. C'est CQRS qui rend l'event sourcing praticable : l'écriture ne charge que par identifiant, et la lecture se construit à partir des événements. Une projection lit le journal dans l'ordre, retient sa position et remplit un modèle taillé pour un écran.

using System;
using System.Collections.Generic;

namespace Ventes.Evenements;

public sealed record ResumeCommande(int NombreDeLignes, decimal Total, bool EstValidee);

// Un modele de lecture, taille pour un ecran, rempli par les evenements.
public sealed class ProjectionResumes
{
    private readonly Dictionary<Guid, ResumeCommande> _resumes = [];

    // La projection garde ce dont elle a besoin : LigneRetiree ne porte pas de montant.
    private readonly Dictionary<(Guid, string), decimal> _montants = [];

    // Ce qui a deja ete lu du journal : on reprend la, sans rien traiter deux fois.
    public int Position { get; private set; }

    public ResumeCommande? Lire(Guid commande) => _resumes.GetValueOrDefault(commande);

    public void Rattraper(MagasinEnMemoire magasin)
    {
        foreach (var evenement in magasin.LireDepuis(Position))
        {
            Traiter(evenement);
            Position++;
        }
    }

    private void Traiter(IEvenement evenement)
    {
        var id = evenement.Commande;
        switch (evenement)
        {
            case CommandeOuverte:
                _resumes[id] = new ResumeCommande(0, 0m, false);
                break;
            case LigneAjoutee e:
                var montant = e.Quantite * e.PrixUnitaire;
                _montants[(id, e.Reference)] = montant;
                _resumes[id] = _resumes[id] with
                {
                    NombreDeLignes = _resumes[id].NombreDeLignes + 1,
                    Total = _resumes[id].Total + montant,
                };
                break;
            case LigneRetiree e:
                _montants.Remove((id, e.Reference), out var retire);
                _resumes[id] = _resumes[id] with
                {
                    NombreDeLignes = _resumes[id].NombreDeLignes - 1,
                    Total = _resumes[id].Total - retire,
                };
                break;
            case CommandeValidee:
                _resumes[id] = _resumes[id] with { EstValidee = true };
                break;
        }
    }
}

// Une question que personne n'avait posee : quelles references sont retirees ?
public sealed class ProjectionRetraits
{
    private readonly Dictionary<string, int> _retraits = [];

    public int Position { get; private set; }

    public IReadOnlyDictionary<string, int> Retraits => _retraits;

    public void Rattraper(MagasinEnMemoire magasin)
    {
        foreach (var evenement in magasin.LireDepuis(Position))
        {
            if (evenement is LigneRetiree e)
            {
                _retraits[e.Reference] = _retraits.GetValueOrDefault(e.Reference) + 1;
            }

            Position++;
        }
    }
}

static class DemonstrationProjections
{
    public static void Executer()
    {
        var magasin = new MagasinEnMemoire();
        var resumes = new ProjectionResumes();

        var id = Guid.NewGuid();
        var commande = Commande.Ouvrir(id, Guid.NewGuid(), 500m);
        commande.AjouterLigne("ECRAN-27", 1, 300m);
        commande.AjouterLigne("CLAVIER", 2, 40m);
        commande.RetirerLigne("CLAVIER");
        magasin.Ajouter(id, commande.Version, commande.NouveauxEvenements);

        resumes.Rattraper(magasin);
        Console.WriteLine(resumes.Lire(id));
        // ResumeCommande { NombreDeLignes = 1, Total = 300, EstValidee = False }

        var suite = Commande.Reconstituer(magasin.Lire(id));
        suite.Valider();
        magasin.Ajouter(id, suite.Version, suite.NouveauxEvenements);

        // Tant que la projection n'a pas rattrape le journal, elle montre le passe.
        Console.WriteLine(resumes.Lire(id)!.EstValidee); // False
        resumes.Rattraper(magasin);
        Console.WriteLine(resumes.Lire(id)!.EstValidee); // True

        // Une projection ecrite apres coup repond pour tout l'historique.
        var retraits = new ProjectionRetraits();
        retraits.Rattraper(magasin);
        Console.WriteLine(string.Join(", ", retraits.Retraits)); // [CLAVIER, 1]
    }
}

Trois propriétés en découlent. La projection est jetable : effacée, elle se reconstruit en relisant le journal. Elle est en retard tant qu'elle ne l'a pas rattrapé, ce que montre le False : c'est la cohérence à terme, que l'interface doit assumer dès que la mise à jour est asynchrone. Et une question nouvelle reçoit une réponse sur tout le passé : chez Young, un rapport sur les articles retirés des paniers remonte ainsi à des années dès sa livraison. En production, la position s'enregistre dans la même transaction que le modèle de lecture ; sinon, une reprise traite un événement deux fois, ou le perd. Le journal ordonné sert du même coup de boîte d'envoi : Young décrit ce magasin employé comme une file, suivie par numéro de séquence, et la double écriture que le cours sur le DDD tactique règle par l'outbox disparaît.

Un agrégat chargé de milliers d'événements coûte cher à reconstituer. L'instantané, le « rolling snapshot » de Young, photographie l'état à une version donnée : on part de lui, et on ne rejoue que la suite.

using System;
using System.Collections.Generic;
using System.Linq;

namespace Ventes.Evenements;

public sealed record LigneEnInstantane(string Reference, int Quantite, decimal PrixUnitaire);

// L'etat a une version donnee, range a part du journal : une optimisation, pas une verite.
public sealed record Instantane(
    Guid Commande,
    int Version,
    decimal Plafond,
    bool Validee,
    IReadOnlyList<LigneEnInstantane> Lignes);

public sealed partial class Commande
{
    public Instantane Photographier()
    {
        // Seul un etat enregistre se photographie : sinon, la version mentirait.
        if (_nouveaux.Count > 0)
        {
            throw new InvalidOperationException("Des evenements ne sont pas enregistres.");
        }

        var lignes = _lignes.Select(
            l => new LigneEnInstantane(l.Key, l.Value.Quantite, l.Value.PrixUnitaire));
        return new Instantane(Id, Version, _plafond, _validee, [.. lignes]);
    }

    // Partir de l'instantane, puis rejouer seulement ce qui l'a suivi.
    public static Commande Reconstituer(Instantane instantane, IEnumerable<IEvenement> suite)
    {
        var commande = new Commande
        {
            Id = instantane.Commande,
            Version = instantane.Version,
            _plafond = instantane.Plafond,
            _validee = instantane.Validee,
        };
        foreach (var ligne in instantane.Lignes)
        {
            commande._lignes[ligne.Reference] = (ligne.Quantite, ligne.PrixUnitaire);
        }

        foreach (var evenement in suite)
        {
            commande.Appliquer(evenement);
            commande.Version++;
        }

        return commande;
    }
}

static class DemonstrationInstantane
{
    public static void Executer()
    {
        var magasin = new MagasinEnMemoire();
        var id = Guid.NewGuid();
        var commande = Commande.Ouvrir(id, Guid.NewGuid(), 500m);
        commande.AjouterLigne("ECRAN-27", 1, 300m);
        commande.AjouterLigne("CLAVIER", 2, 40m);
        magasin.Ajouter(id, commande.Version, commande.NouveauxEvenements);

        // Pris a la version 3, par un processus a part, sans rien bloquer.
        var instantane = Commande.Reconstituer(magasin.Lire(id)).Photographier();

        var suite = Commande.Reconstituer(magasin.Lire(id));
        suite.RetirerLigne("CLAVIER");
        suite.AjouterLigne("DOCK", 1, 150m);
        magasin.Ajouter(id, suite.Version, suite.NouveauxEvenements);

        var apres = magasin.Lire(id, instantane.Version);
        var rapide = Commande.Reconstituer(instantane, apres);
        var complete = Commande.Reconstituer(magasin.Lire(id));
        Console.WriteLine($"{apres.Count} evenements rejoues, total {rapide.Total}");
        // 2 evenements rejoues, total 450
        Console.WriteLine(rapide.Total == complete.Total && rapide.Version == complete.Version);
        // True
    }
}

Young range les instantanés dans une table à part, associés à leur version et pris en arrière-plan ; placés dans le flux, ils entreraient en conflit de version avec les écritures d'un agrégat très sollicité. Il les tient pour une heuristique et conseille de développer sans, quitte à les ajouter plus tard. Avec le sérialiseur par défaut, changer la forme de l'agrégat oblige à les jeter et à les recalculer ; Young propose un Memento, qui versionne leur schéma à part : c'est le rôle du record Instantane.

Versionner les événements

Un événement enregistré ne se modifie pas : des consommateurs l'ont lu, des projections en dépendent, et le journal perdrait sa valeur de trace. Le code, lui, change. Young a posé la règle à DDD eXchange 2017, selon le compte rendu d'InfoQ du 17 juillet : une nouvelle version d'un événement doit se convertir depuis l'ancienne, faute de quoi c'est un nouvel événement. Son livre sur la question, « Versioning in an Event Sourced System », est resté inachevé sur Leanpub. La conversion se fait à la lecture, l'upcasting, et le journal reste tel qu'il a été écrit.

using System;
using System.Text.Json;
using System.Text.Json.Nodes;

namespace Ventes.Evenements.Versions;

// Ce que le magasin conserve vraiment : un type, une version de schema, du JSON.
public sealed record Enveloppe(string Type, int Version, string Donnees);

// Version 2 de LigneAjoutee. La version 1 stockait le montant de la ligne (Prix),
// pas le prix unitaire. Devise est venue plus tard, facultative.
public sealed record LigneAjoutee(
    string Reference, int Quantite, decimal PrixUnitaire, string Devise = "EUR");

public static class LecteurLigneAjoutee
{
    public static LigneAjoutee Lire(Enveloppe enveloppe)
    {
        var donnees = JsonNode.Parse(enveloppe.Donnees)!.AsObject();

        // Convertir la version 1 en version 2 a la lecture ; le journal, lui, ne change pas.
        if (enveloppe.Version == 1)
        {
            var prix = donnees["Prix"]!.GetValue<decimal>();
            var quantite = donnees["Quantite"]!.GetValue<int>();
            donnees.Remove("Prix");
            donnees["PrixUnitaire"] = prix / quantite;
        }

        // Un champ absent prend la valeur par defaut du parametre : c'est le schema faible.
        return donnees.Deserialize<LigneAjoutee>()!;
    }
}

static class DemonstrationVersions
{
    public static void Executer()
    {
        Enveloppe[] journal =
        [
            new("LigneAjoutee", 1, """{"Reference":"CABLE","Quantite":4,"Prix":60}"""),
            new("LigneAjoutee", 2, """{"Reference":"DOCK","Quantite":1,"PrixUnitaire":150}"""),
            new("LigneAjoutee", 2,
                """{"Reference":"ADAPTATEUR","Quantite":2,"PrixUnitaire":12,"Devise":"USD"}"""),
        ];

        foreach (var enveloppe in journal)
        {
            var ligne = LecteurLigneAjoutee.Lire(enveloppe);
            Console.WriteLine(
                $"{ligne.Reference} : {ligne.Quantite} x {ligne.PrixUnitaire} {ligne.Devise}");
        }
        // CABLE : 4 x 15 EUR
        // DOCK : 1 x 150 EUR
        // ADAPTATEUR : 2 x 12 USD
    }
}

La version 1 stockait le montant de la ligne. Garder le nom Prix en changeant son sens ferait lire aux anciens événements un prix unitaire quatre fois trop élevé, sans la moindre erreur ; le numéro de version dans l'enveloppe rend l'écart explicite, et le lecteur convertit. La devise relève d'un autre mécanisme, le schéma faible : un champ ajouté, facultatif, que les anciens événements n'ont pas et que System.Text.Json remplit avec la valeur par défaut du paramètre. Ce n'est juste que si cette valeur dit vrai pour le passé : EUR, pour une boutique qui n'a vendu qu'en euros.

Quand ne pas le faire

Fowler met en garde dès 2011 : CQRS est un saut mental, à n'entreprendre que si le bénéfice le justifie, et à réserver à certaines parties d'un système, un bounded context, jamais au système entier. La plupart des cas qu'il a rencontrés ont mal tourné, le motif ajoutant une complexité risquée là où un modèle unique suffisait. Young attache la valeur du journal au même critère que le DDD : élevée là où l'entreprise tire un avantage concurrentiel, elle peut être négative ailleurs. Le cours sur le DDD stratégique appelle cette zone le domaine cœur.

Les coûts de l'event sourcing sont connus d'avance :

  • la cohérence à terme, dès que la lecture est asynchrone ;
  • un schéma d'événements à maintenir à perpétuité, puisque rien ne se réécrit sans copier tout le journal ;
  • des effets externes à tenir strictement à l'écart du rejeu ;
  • un journal immuable face à l'obligation d'effacer des données personnelles : où les ranger se décide avant le premier événement ;
  • une équipe qui doit penser en faits plutôt qu'en lignes de table.

Le signal pour commencer est un besoin, pas une mode. Des écrans qui déforment le modèle appellent une séparation de la lecture, qui ne coûte presque rien. Une exigence d'audit, ou des questions sur l'histoire des données, appellent le journal. Un back-office qui édite des fiches n'a besoin ni de l'un ni de l'autre.

Ce cours vous a servi ? Offrir un café Signaler une erreur