Architecture & conception

Microservices et messagerie

Outbox, idempotence, saga, RabbitMQ et Service Bus.

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

Un microservice possède ses données : aucun autre service n'écrit dans sa base, et c'est ce qui permet de le déployer seul. Sa taille, selon le guide Azure que cite le cours sur le DDD stratégique, va de l'agrégat au bounded context. Ce qu'une transaction faisait d'un bloc se répartit alors entre plusieurs bases, et les services se préviennent par messages, à travers un broker qui les découple dans le temps : la fidélité peut être arrêtée pendant que les ventes continuent. Ce découplage a un prix : un message peut se perdre ou arriver deux fois, une écriture et sa publication peuvent diverger, et un traitement qui touche trois services ne s'annule plus par un rollback. Les exemples des quatre premières sections sont les fichiers d'un projet console .NET 10 avec EF Core 10.0.12 sur SQLite ; ils affichent ce que disent leurs commentaires. Ceux de RabbitMQ sont compilés contre RabbitMQ.Client 7.0.0, celle du cache local (la courante est la 7.2.2, du 5 août 2026), sans serveur pour les exécuter ; celui d'Azure Service Bus n'a été ni compilé contre le paquet ni exécuté.

Au plus une fois, au moins une fois

Un broker garde un message tant que son consommateur ne l'a pas acquitté, et le moment de l'acquittement fixe la garantie. Acquitter à la réception, avant de traiter, c'est livrer au plus une fois : si le consommateur tombe entre les deux, le broker a déjà oublié le message, et son effet n'aura jamais lieu. Acquitter après le traitement, c'est livrer au moins une fois : si le consommateur tombe après l'effet mais avant l'acquittement, le broker relivre, et l'effet a lieu deux fois. La file simulée suivante fait tomber son consommateur une fois, avant ou après l'effet du premier message.

using System;
using System.Collections.Generic;

namespace Messagerie.Livraison;

// Une file reduite a l'essentiel : un message livre reste « en vol » tant qu'il
// n'est pas acquitte, et revient en tete si le consommateur disparait.
public sealed class FileSimulee
{
    private readonly LinkedList<string> _prets = new();
    private readonly Dictionary<int, string> _enVol = [];
    private int _numero;

    public void Publier(string message) => _prets.AddLast(message);

    public (int Numero, string Message)? Recevoir()
    {
        if (_prets.First is not { } premier)
        {
            return null;
        }

        _prets.RemoveFirst();
        _enVol[++_numero] = premier.Value;
        return (_numero, premier.Value);
    }

    public void Acquitter(int numero) => _enVol.Remove(numero);

    // La connexion du consommateur tombe : ce qui n'est pas acquitte est relivre.
    public void PerdreLeConsommateur()
    {
        foreach (var message in _enVol.Values)
        {
            _prets.AddFirst(message);
        }

        _enVol.Clear();
    }
}

static class DemonstrationLivraison
{
    // Le consommateur tombe une fois, pendant le premier message recu.
    static List<string> Consommer(bool acquitterAvant, bool plantageApresEffet)
    {
        var file = new FileSimulee();
        file.Publier("facture 1");
        file.Publier("facture 2");
        var effets = new List<string>();
        var plantageAFaire = true;

        while (file.Recevoir() is (var numero, var message))
        {
            if (acquitterAvant)
            {
                file.Acquitter(numero);
            }

            var tombe = plantageAFaire;
            plantageAFaire = false;
            if (tombe && !plantageApresEffet)
            {
                file.PerdreLeConsommateur();
                continue;
            }

            effets.Add(message); // l'effet : un e-mail envoye, une ligne ecrite

            if (tombe)
            {
                file.PerdreLeConsommateur();
                continue;
            }

            if (!acquitterAvant)
            {
                file.Acquitter(numero);
            }
        }

        return effets;
    }

    public static void Executer()
    {
        foreach (var avant in new[] { true, false })
        {
            foreach (var apresEffet in new[] { false, true })
            {
                Console.WriteLine(
                    $"Acquitter {(avant ? "avant" : "apres")}, plantage "
                    + $"{(apresEffet ? "apres" : "avant")} l'effet : "
                    + string.Join(", ", Consommer(avant, apresEffet)));
            }
        }
        // Acquitter avant, plantage avant l'effet : facture 2
        // Acquitter avant, plantage apres l'effet : facture 1, facture 2
        // Acquitter apres, plantage avant l'effet : facture 1, facture 2
        // Acquitter apres, plantage apres l'effet : facture 1, facture 1, facture 2
    }
}

Chaque ordre a sa fenêtre : acquitter avant perd la facture 1 quand la panne précède l'effet, acquitter après la traite deux fois quand la panne le suit. Aucun ordre ne ferme les deux fenêtres, parce que l'effet et l'acquittement touchent deux systèmes, la base du consommateur et le broker, qu'aucune transaction commune ne lie. « Exactement une fois » n'est donc pas une option du transport, c'est un résultat qui se construit : une livraison au moins une fois, et un consommateur qui reconnaît ce qu'il a déjà traité, pour que chaque effet n'ait lieu qu'une fois. Perdre un message est rarement acceptable ; la suite choisit donc le doublon, puis le neutralise.

La double écriture et la boîte d'envoi

Le même problème se pose chez l'émetteur. Valider une commande, c'est l'enregistrer et publier CommandeValidee : deux écritures, l'une dans la base, l'autre dans le broker. Le modèle suivant sert aux deux versions : une commande, une table des messages sortants, et un bus qui tombe sur demande.

using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;

namespace Ventes.Outbox;

public sealed record CommandeValidee(Guid Commande, Guid Client, decimal Total);

public sealed class Commande
{
    public Guid Id { get; init; }

    public Guid Client { get; init; }

    public decimal Total { get; init; }
}

// Une ligne par message a publier. Numero donne l'ordre d'ecriture ;
// MessageId suit le message jusque chez le consommateur.
public sealed class MessageSortant
{
    public long Numero { get; init; }

    public Guid MessageId { get; init; }

    public required string Type { get; init; }

    public required string Contenu { get; init; }

    public DateTime? PublieLe { get; set; }
}

public sealed class VentesContexte(DbContextOptions<VentesContexte> options) : DbContext(options)
{
    public DbSet<Commande> Commandes => Set<Commande>();

    public DbSet<MessageSortant> BoiteEnvoi => Set<MessageSortant>();

    protected override void OnModelCreating(ModelBuilder modele) =>
        modele.Entity<MessageSortant>().HasKey(m => m.Numero);
}

// Le broker, vu du producteur : il tombe quand on le lui demande.
public sealed class BusCapricieux
{
    public bool EnPanne { get; set; }

    public List<Guid> Recus { get; } = [];

    public Task Publier(Guid messageId, string type, string contenu)
    {
        if (EnPanne)
        {
            throw new InvalidOperationException("Broker injoignable.");
        }

        Recus.Add(messageId);
        return Task.CompletedTask;
    }
}

La version naïve enregistre puis publie, et c'est ce qui cloche : si le broker est injoignable, ou si le processus s'arrête entre les deux, la commande existe et aucun autre service ne l'apprendra.

using System;
using System.Text.Json;
using System.Threading.Tasks;
using Microsoft.Data.Sqlite;
using Microsoft.EntityFrameworkCore;

namespace Ventes.Outbox;

public sealed class ValiderCommandeDoubleEcriture(VentesContexte contexte, BusCapricieux bus)
{
    public async Task Executer(Guid commande, Guid client, decimal total)
    {
        contexte.Commandes.Add(new Commande { Id = commande, Client = client, Total = total });
        await contexte.SaveChangesAsync();

        // Seconde ecriture, dans un autre systeme : rien ne la lie a la premiere.
        var evenement = new CommandeValidee(commande, client, total);
        var contenu = JsonSerializer.Serialize(evenement);
        await bus.Publier(Guid.NewGuid(), nameof(CommandeValidee), contenu);
    }
}

static class DemonstrationDoubleEcriture
{
    public static async Task Executer()
    {
        using var connexion = new SqliteConnection("DataSource=:memory:");
        connexion.Open();
        var options = new DbContextOptionsBuilder<VentesContexte>().UseSqlite(connexion).Options;
        using var contexte = new VentesContexte(options);
        contexte.Database.EnsureCreated();

        var bus = new BusCapricieux { EnPanne = true };
        try
        {
            await new ValiderCommandeDoubleEcriture(contexte, bus)
                .Executer(Guid.NewGuid(), Guid.NewGuid(), 300m);
        }
        catch (InvalidOperationException ex)
        {
            Console.WriteLine(ex.Message); // Broker injoignable.
        }

        var commandes = await contexte.Commandes.CountAsync();
        Console.WriteLine($"{commandes} commande, {bus.Recus.Count} message");
        // 1 commande, 0 message
    }
}
using System;
using System.Linq;
using System.Text.Json;
using System.Threading.Tasks;
using Microsoft.Data.Sqlite;
using Microsoft.EntityFrameworkCore;

namespace Ventes.Outbox;

public sealed class ValiderCommande(VentesContexte contexte)
{
    public async Task Executer(Guid commande, Guid client, decimal total)
    {
        contexte.Commandes.Add(new Commande { Id = commande, Client = client, Total = total });
        contexte.BoiteEnvoi.Add(new MessageSortant
        {
            MessageId = Guid.NewGuid(),
            Type = nameof(CommandeValidee),
            Contenu = JsonSerializer.Serialize(new CommandeValidee(commande, client, total)),
        });

        // Les deux INSERT partent dans la meme transaction : tout ou rien.
        await contexte.SaveChangesAsync();
    }
}

// Le relais : lit la boite d'envoi dans l'ordre, publie, marque.
public sealed class RelaisBoiteEnvoi(VentesContexte contexte, BusCapricieux bus)
{
    public async Task PublierEnAttente()
    {
        var enAttente = await contexte.BoiteEnvoi
            .Where(m => m.PublieLe == null)
            .OrderBy(m => m.Numero)
            .Take(100)
            .ToListAsync();

        foreach (var message in enAttente)
        {
            await bus.Publier(message.MessageId, message.Type, message.Contenu);

            // Un arret juste ici : publie mais pas marque, donc republie au tour suivant.
            message.PublieLe = DateTime.UtcNow;
            await contexte.SaveChangesAsync();
        }
    }
}

static class DemonstrationBoiteEnvoi
{
    public static async Task Executer()
    {
        using var connexion = new SqliteConnection("DataSource=:memory:");
        connexion.Open();
        var options = new DbContextOptionsBuilder<VentesContexte>().UseSqlite(connexion).Options;
        using var contexte = new VentesContexte(options);
        contexte.Database.EnsureCreated();

        var bus = new BusCapricieux { EnPanne = true };
        var relais = new RelaisBoiteEnvoi(contexte, bus);
        await new ValiderCommande(contexte).Executer(Guid.NewGuid(), Guid.NewGuid(), 300m);

        try
        {
            await relais.PublierEnAttente();
        }
        catch (InvalidOperationException ex)
        {
            Console.WriteLine(ex.Message); // Broker injoignable.
        }

        var enAttente = await contexte.BoiteEnvoi.CountAsync(m => m.PublieLe == null);
        var commandes = await contexte.Commandes.CountAsync();
        Console.WriteLine($"{commandes} commande, {enAttente} en attente");
        // 1 commande, 1 en attente

        bus.EnPanne = false;
        await relais.PublierEnAttente();
        Console.WriteLine($"{bus.Recus.Count} message publie"); // 1 message publie
    }
}

Inverser l'ordre ne sauve rien : publier d'abord annonce une commande que la base peut encore refuser. La boîte d'envoi, ou outbox, que le cours sur le DDD tactique présente en quelques lignes, retire la seconde écriture de la requête. Le message devient une ligne de la même base, ajoutée au même SaveChanges, et le cours sur EF Core l'a montré : les commandes d'un SaveChanges partent dans une transaction. La commande et son message existent ensemble, ou pas du tout. Un relais lit ensuite les lignes non publiées dans l'ordre de Numero, publie et marque ; si le broker tombe, la ligne attend le tour suivant.

Le relais a son propre trou, et le commentaire le nomme : publié mais pas encore marqué, le message sera republié. La boîte d'envoi garantit donc une publication au moins une fois, et le MessageId, fixé à l'écriture de la ligne, reste le même d'une publication à l'autre : c'est lui que le consommateur reconnaîtra. En production, le relais tourne dans un BackgroundService, avec une portée par passe et l'arrêt propre que décrit le cours « Résilience et performance » ; à plusieurs instances, chacune verrouille les lignes qu'elle prend, sinon deux relais peuvent publier le même lot ; l'ordre de publication, lui, se perd. Même avec un seul relais, sur un SGBD concurrent, l'ordre des identifiants n'est pas forcément celui des validations. Avec l'event sourcing, le journal ordonné en tient lieu (cours « CQRS et Event Sourcing »).

MassTransit et NServiceBus offrent tout faits une boîte d'envoi transactionnelle, la déduplication à la réception et des sagas persistées, et leur licence pèse dans le choix. MassTransit est commercial depuis sa version 9 : le contrat de Massient, mis à jour le 17 septembre 2026, vise la v9 et les suivantes et laisse les versions 8 et antérieures sous leur licence open source d'origine, hors du support qu'il prévoit ; la dernière version sur NuGet est la 9.2.2, du 14 septembre 2026. NServiceBus, de Particular Software, est gratuit en développement ; sa page de tarifs propose une édition Community gratuite, limitée à 3 endpoints et 10 000 messages par jour, sans dire explicitement si elle vaut en production, et un programme pour les petites entreprises. Ces lectures datent du 26 septembre 2026 et ne valent pas avis juridique.

Le consommateur idempotent

Un consommateur est idempotent quand traiter deux fois le même message laisse le même état que le traiter une fois. Le service de fidélité crédite un point par tranche de dix euros d'une commande validée ; il a sa propre base, avec une table de déduplication.

using System;
using Microsoft.EntityFrameworkCore;

namespace Fidelite;

// Le contrat du message, tel que le service de fidelite le lit.
public sealed record CommandeValidee(Guid Commande, Guid Client, decimal Total);

public sealed class Compte
{
    public Guid Client { get; init; }

    public int Points { get; set; }
}

// La table de deduplication : un message traite y laisse son identifiant.
public sealed class MessageTraite
{
    public Guid MessageId { get; init; }

    public DateTime TraiteLe { get; init; }
}

public sealed class FideliteContexte(DbContextOptions<FideliteContexte> options)
    : DbContext(options)
{
    public DbSet<Compte> Comptes => Set<Compte>();

    public DbSet<MessageTraite> MessagesTraites => Set<MessageTraite>();

    protected override void OnModelCreating(ModelBuilder modele)
    {
        modele.Entity<Compte>().HasKey(c => c.Client);
        modele.Entity<MessageTraite>().HasKey(m => m.MessageId);
    }
}

Le consommateur naïf ajoute les points à chaque livraison, et c'est ce qui cloche : le message que le relais a republié après une panne crédite la même commande deux fois.

using System;
using System.Threading.Tasks;
using Microsoft.Data.Sqlite;
using Microsoft.EntityFrameworkCore;

namespace Fidelite;

public sealed class CrediterPointsNaif(FideliteContexte contexte)
{
    public async Task Traiter(Guid messageId, CommandeValidee evenement)
    {
        var compte = await contexte.Comptes.SingleAsync(c => c.Client == evenement.Client);
        compte.Points += (int)(evenement.Total / 10);
        await contexte.SaveChangesAsync();
    }
}

static class DemonstrationDoublon
{
    public static async Task Executer()
    {
        using var connexion = new SqliteConnection("DataSource=:memory:");
        connexion.Open();
        var options = new DbContextOptionsBuilder<FideliteContexte>().UseSqlite(connexion).Options;
        using var contexte = new FideliteContexte(options);
        contexte.Database.EnsureCreated();

        var client = Guid.NewGuid();
        contexte.Comptes.Add(new Compte { Client = client });
        await contexte.SaveChangesAsync();

        // Le meme message, livre deux fois : le relais l'a republie apres une panne.
        var messageId = Guid.NewGuid();
        var evenement = new CommandeValidee(Guid.NewGuid(), client, 300m);
        var consommateur = new CrediterPointsNaif(contexte);
        await consommateur.Traiter(messageId, evenement);
        await consommateur.Traiter(messageId, evenement);

        Console.WriteLine((await contexte.Comptes.SingleAsync()).Points); // 60
    }
}
using System;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.Data.Sqlite;
using Microsoft.EntityFrameworkCore;

namespace Fidelite;

public sealed class CrediterPoints(FideliteContexte contexte)
{
    public async Task Traiter(Guid messageId, CommandeValidee evenement)
    {
        if (await contexte.MessagesTraites.AnyAsync(m => m.MessageId == messageId))
        {
            return; // deja traite : on acquitte sans rien refaire
        }

        // La trace et l'effet, dans la meme transaction : tout ou rien.
        await using var transaction = await contexte.Database.BeginTransactionAsync();

        // La cle primaire arrete ici un doublon concurrent du meme message.
        var trace = new MessageTraite { MessageId = messageId, TraiteLe = DateTime.UtcNow };
        contexte.MessagesTraites.Add(trace);
        await contexte.SaveChangesAsync();

        // Une mise a jour relative, Points = Points + n : deux messages distincts
        // sur le meme compte s'additionnent au lieu de s'ecraser.
        var points = (int)(evenement.Total / 10);
        await contexte.Comptes
            .Where(c => c.Client == evenement.Client)
            .ExecuteUpdateAsync(s => s.SetProperty(c => c.Points, c => c.Points + points));

        await transaction.CommitAsync();
    }
}

static class DemonstrationIdempotence
{
    public static async Task Executer()
    {
        using var connexion = new SqliteConnection("DataSource=:memory:");
        connexion.Open();
        var options = new DbContextOptionsBuilder<FideliteContexte>().UseSqlite(connexion).Options;
        using var contexte = new FideliteContexte(options);
        contexte.Database.EnsureCreated();

        var client = Guid.NewGuid();
        contexte.Comptes.Add(new Compte { Client = client });
        await contexte.SaveChangesAsync();

        var messageId = Guid.NewGuid();
        var evenement = new CommandeValidee(Guid.NewGuid(), client, 300m);
        var consommateur = new CrediterPoints(contexte);
        await consommateur.Traiter(messageId, evenement);
        await consommateur.Traiter(messageId, evenement);

        // Relire en base : ExecuteUpdate ne touche pas les entites deja suivies.
        Console.WriteLine(await contexte.Comptes.Select(c => c.Points).SingleAsync()); // 30
    }
}

Tout tient à la transaction : la ligne de MessagesTraites et les points sont validés ensemble. Enregistrer la trace à part, avant ou après, rouvrirait la fenêtre de la première section. La vérification initiale n'est qu'un raccourci ; c'est la clé primaire qui garantit. Deux livraisons simultanées du même message peuvent passer toutes deux la vérification : la seconde à enregistrer viole alors la contrainte, SaveChanges lève une DbUpdateException, dont l'exception interne dit, avec SQLite, SQLite Error 19: 'UNIQUE constraint failed: MessagesTraites.MessageId' ; sa transaction est annulée, et le message, non acquitté, revient ; à la livraison suivante, la vérification le trouve.

La clé primaire ne protège que du même message. Deux messages distincts sur le même compte, traités en même temps, passent tous deux ; si chacun lisait le compte pour écrire son total, EF Core enverrait UPDATE "Comptes" SET "Points" = @p0, une valeur absolue, et deux commandes de 300 € donneraient 30 points au lieu de 60. La version juste écrit donc une mise à jour relative, Points = Points + n, par ExecuteUpdateAsync dans la transaction ouverte à la main. Un jeton de concurrence, que le cours sur EF Core évoque à propos d'ExecuteUpdate, ou la sérialisation des messages d'un client, comme les sessions de Service Bus, conviennent aussi.

La table se purge, en gardant chaque identifiant plus longtemps que le broker ne peut relivrer et que le relais ne peut republier une ligne restée non marquée. Une opération idempotente par nature s'en passe : passer un statut à « validée » se répète sans dommage, ajouter des points non. Une clé métier, un crédit par commande sous clé primaire, résiste même à deux messages distincts annonçant la même commande.

Saga : orchestration ou chorégraphie

Hector Garcia-Molina et Kenneth Salem ont nommé la saga en 1987, pour les transactions longues d'un système de bases de données : une transaction qui s'écrit comme une suite de transactions plus courtes, avec la garantie que toutes aboutissent ou que des transactions de compensation corrigent l'exécution partielle. Les microservices reprennent l'idée entre services : chaque étape est une transaction locale dans la base d'un service, et un échec déclenche la compensation des étapes faites, dans l'ordre inverse. Une compensation n'efface rien, elle corrige par un effet de sens contraire : libérer le stock réservé, rembourser le paiement.

Dans l'orchestration, un coordinateur tient l'état de la saga, ordonne chaque étape et décide des compensations. Le suivant est une machine à états : il reçoit une réponse, rend les ordres à envoyer, et son état s'enregistre entre deux messages, car une saga peut durer des heures.

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

namespace Ventes.Saga;

// Ce que l'orchestrateur ordonne, et ce que les participants lui repondent.
public sealed record ReserverStock(Guid Commande);
public sealed record LibererStock(Guid Commande);
public sealed record DebiterPaiement(Guid Commande, decimal Montant);
public sealed record ConfirmerCommande(Guid Commande);
public sealed record AnnulerCommande(Guid Commande, string Motif);

public sealed record StockReserve(Guid Commande);
public sealed record StockIndisponible(Guid Commande);
public sealed record PaiementAccepte(Guid Commande);
public sealed record PaiementRefuse(Guid Commande);

public enum Etape { StockDemande, PaiementDemande, Confirmee, Annulee }

// L'etat d'une saga, enregistre entre deux messages : elle peut durer des heures.
public sealed class SagaCommande(Guid commande, decimal montant)
{
    public Etape Etape { get; private set; } = Etape.StockDemande;

    public IReadOnlyList<object> Demarrer() => [new ReserverStock(commande)];

    public IReadOnlyList<object> Recevoir(object reponse) => (Etape, reponse) switch
    {
        (Etape.StockDemande, StockReserve) =>
            Passer(Etape.PaiementDemande, new DebiterPaiement(commande, montant)),
        (Etape.StockDemande, StockIndisponible) =>
            Passer(Etape.Annulee, new AnnulerCommande(commande, "stock")),
        (Etape.PaiementDemande, PaiementAccepte) =>
            Passer(Etape.Confirmee, new ConfirmerCommande(commande)),

        // Le paiement echoue apres la reservation : on compense ce qui a ete fait.
        (Etape.PaiementDemande, PaiementRefuse) =>
            Passer(Etape.Annulee,
                new LibererStock(commande), new AnnulerCommande(commande, "paiement")),

        // Un doublon, ou une reponse arrivee trop tard : l'etape l'a deja depassee.
        _ => [],
    };

    private object[] Passer(Etape suivante, params object[] ordres)
    {
        Etape = suivante;
        return ordres;
    }
}

static class DemonstrationSaga
{
    public static void Executer()
    {
        var id = Guid.NewGuid();
        var saga = new SagaCommande(id, 300m);
        object[] reponses = [new StockReserve(id), new StockReserve(id), new PaiementRefuse(id)];

        var ordres = saga.Demarrer().Concat(reponses.SelectMany(saga.Recevoir)).ToList();
        Console.WriteLine(string.Join(", ", ordres.Select(o => o.GetType().Name)));
        // ReserverStock, DebiterPaiement, LibererStock, AnnulerCommande
        Console.WriteLine(saga.Etape); // Annulee
    }
}

Le second StockReserve, un doublon, ne produit rien : filtrer sur le couple étape et réponse rend la saga idempotente sans table. Le refus du paiement arrive après la réservation, d'où la compensation LibererStock. Dans la chorégraphie, personne ne coordonne : le stock réserve en recevant CommandeValidee et publie StockReserve, le paiement débite en recevant StockReserve, et le stock libère de lui-même en recevant PaiementRefuse. Le guide des patrons d'architecture Azure résume l'arbitrage. La chorégraphie n'a ni coordinateur ni point de défaillance unique et convient aux flux simples, mais le flux ne se lit nulle part et les participants risquent des dépendances circulaires. L'orchestration convient aux flux complexes, au prix d'un coordinateur à écrire, qui devient un point de défaillance. Le même guide rappelle ce qu'une saga ne donne pas : l'isolation. Entre deux étapes, les autres voient un état intermédiaire, un stock réservé pour une commande qui sera annulée.

RabbitMQ : échanges, files, liaisons

Avec AMQP 0.9.1, son protocole historique, RabbitMQ ne publie jamais directement dans une file. Le producteur publie dans un échange, avec une clé de routage ; l'échange copie le message dans chaque file qu'une liaison lui rattache et dont le motif correspond. Un échange direct compare la clé à l'identique, un topic par motif (commande.*, commande.#), un fanout ignore la clé et copie partout. Chaque service a sa file : ajouter un consommateur, c'est ajouter une file et une liaison, sans toucher au producteur, qui ne déclare que son échange. La connexion et le canal se créent une fois : la documentation de RabbitMQ suppose des connexions longues, pas une par opération. L'exemple suit l'API entièrement asynchrone de RabbitMQ.Client 7, où IChannel remplace l'ancien IModel.

using System;
using System.Text.Json;
using System.Threading.Tasks;
using RabbitMQ.Client;

namespace Ventes.Rabbit;

public sealed record CommandeValidee(Guid Commande, Guid Client, decimal Total);

// Cote ventes : le producteur ne connait que son echange, pas les files des autres.
public sealed class PublieurVentes(IChannel canal)
{
    // Une connexion par application et un canal par publieur, crees au demarrage :
    // ils vivent aussi longtemps que le service.
    public static async Task<PublieurVentes> Creer(IConnection connexion)
    {
        // Confirmations suivies : l'await rend la main quand le broker a pris le message.
        var canal = await connexion.CreateChannelAsync(new CreateChannelOptions(
            publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true));

        // Topic : il route selon la cle, par motif. Redeclarer a l'identique ne change rien.
        await canal.ExchangeDeclareAsync("ventes", ExchangeType.Topic, durable: true);
        return new PublieurVentes(canal);
    }

    public async Task Publier(Guid messageId, CommandeValidee evenement)
    {
        var proprietes = new BasicProperties
        {
            MessageId = messageId.ToString(), // celui de la boite d'envoi
            ContentType = "application/json",
            Persistent = true,
        };

        // mandatory : une PublishException plutot qu'un message que rien ne route.
        await canal.BasicPublishAsync("ventes", "commande.validee", mandatory: true,
            basicProperties: proprietes, body: JsonSerializer.SerializeToUtf8Bytes(evenement));
    }
}

Sans confirmations, le producteur ne sait pas si le broker a pris le message. Avec les confirmations suivies, BasicPublishAsync attend l'accusé du broker et, selon la documentation de la 7.0.0, lève une PublishException si le message est refusé, ou renvoyé faute de file : c'est ce que le relais doit attendre avant de marquer sa ligne. Le service de fidélité déclare le reste : sa file, ses lettres mortes, sa liaison. Trois réglages y décident des garanties.

using System;
using System.Collections.Generic;
using System.Text.Json;
using System.Threading.Tasks;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;

namespace Ventes.Rabbit;

// Cote fidelite : le consommateur declare sa file, ses lettres mortes et sa liaison.
public static class Consommation
{
    private static async Task DeclarerTopologie(IChannel canal)
    {
        // L'echange des ventes, declare a l'identique : l'ordre de demarrage n'importe plus.
        await canal.ExchangeDeclareAsync("ventes", ExchangeType.Topic, durable: true);

        // Les lettres mortes : un echange, et une file ou les ranger.
        await canal.ExchangeDeclareAsync("fidelite.mortes", ExchangeType.Fanout, durable: true);
        await canal.QueueDeclareAsync("fidelite.mortes", durable: true, exclusive: false,
            autoDelete: false);
        await canal.QueueBindAsync("fidelite.mortes", "fidelite.mortes", routingKey: "");

        // La file : quorum, cinq relivraisons en echec au plus, puis les lettres mortes.
        await canal.QueueDeclareAsync("fidelite.commandes", durable: true, exclusive: false,
            autoDelete: false, arguments: new Dictionary<string, object?>
            {
                ["x-queue-type"] = "quorum",
                ["x-delivery-limit"] = 5,
                ["x-dead-letter-exchange"] = "fidelite.mortes",
            });

        // La liaison : ce que la file recoit de l'echange.
        await canal.QueueBindAsync("fidelite.commandes", "ventes", routingKey: "commande.validee");
    }

    public static async Task<string> Demarrer(
        IChannel canal, Func<Guid, CommandeValidee, Task> traiter)
    {
        await DeclarerTopologie(canal);

        // Pas plus de dix messages en vol a la fois pour ce consommateur.
        await canal.BasicQosAsync(prefetchSize: 0, prefetchCount: 10, global: false);

        var consommateur = new AsyncEventingBasicConsumer(canal);
        consommateur.ReceivedAsync += async (_, livraison) =>
        {
            var evenement = Lire(livraison.Body);
            if (evenement is null
                || !Guid.TryParse(livraison.BasicProperties.MessageId, out var id))
            {
                // Illisible : le relivrer n'y changera rien. Directement aux lettres mortes.
                await canal.BasicRejectAsync(livraison.DeliveryTag, requeue: false);
                return;
            }

            try
            {
                await traiter(id, evenement); // le consommateur idempotent
                await canal.BasicAckAsync(livraison.DeliveryTag, multiple: false);
            }
            catch (Exception)
            {
                // Peut-etre passager : remis en file. Depuis RabbitMQ 4.3, reject compte
                // comme un echec pour x-delivery-limit ; nack remettrait en file sans compter.
                await canal.BasicRejectAsync(livraison.DeliveryTag, requeue: true);
            }
        };

        // autoAck: false, sinon le broker oublie le message des qu'il l'a envoye.
        return await canal.BasicConsumeAsync(
            "fidelite.commandes", autoAck: false, consumer: consommateur);
    }

    private static CommandeValidee? Lire(ReadOnlyMemory<byte> corps)
    {
        try
        {
            return JsonSerializer.Deserialize<CommandeValidee>(corps.Span);
        }
        catch (JsonException)
        {
            return null;
        }
    }
}

autoAck: false est l'acquittement après traitement de la première section ; autoAck: true serait la livraison au plus une fois. prefetchCount borne les messages livrés et non acquittés. Un message rejeté sans remise en file part vers l'échange de lettres mortes de sa file ; la documentation de RabbitMQ 4.3 y ajoute l'expiration d'un TTL, le dépassement d'une longueur maximale et, pour une file quorum, celui de la limite de livraisons. Cette limite vaut 20 par défaut depuis RabbitMQ 4.0, et un message n'est écarté qu'une fois relivré plus de fois qu'elle : avec 5, il peut être livré six fois. Depuis la 4.3, elle ne compte que les échecs : un basic.reject ou une connexion perdue l'incrémentent, un basic.nack remet en file sans compter. Un BasicNackAsync avec remise en file relivrerait donc sans fin un message qui échoue toujours, d'où BasicRejectAsync. Une file quorum transmet par défaut ses lettres mortes au plus une fois ; l'au moins une fois demande dead-letter-strategy à at-least-once et overflow à reject-publish. La documentation recommande de poser ces réglages par une politique plutôt que par des arguments x-, qui imposent de recréer la file ; l'exemple les écrit dans le code pour se lire seul.

Azure Service Bus

Azure Service Bus est un broker géré, sans échange à déclarer. Les consommateurs d'une file se partagent ses messages. Une rubrique est la variante publier-s'abonner : chaque abonnement reçoit sa copie de chaque message, se lit comme une file, et ses règles filtrent ce qu'il retient, ici par le sujet du message. Le paquet Azure.Messaging.ServiceBus, en 7.21.0 du 24 septembre 2026 sur NuGet, n'est pas disponible sur cette machine : l'exemple suit la référence de l'API sur learn.microsoft.com.

using System;
using System.Text.Json;
using System.Threading.Tasks;
using Azure.Identity;
using Azure.Messaging.ServiceBus;
using Azure.Messaging.ServiceBus.Administration;

namespace Ventes.ServiceBus;

public sealed record CommandeValidee(Guid Commande, Guid Client, decimal Total);

public static class Messagerie
{
    private const string Espace = "ventes.servicebus.windows.net";

    // Une fois, a l'installation : l'appel echoue si l'entite existe deja.
    public static async Task CreerEntites()
    {
        var admin = new ServiceBusAdministrationClient(Espace, new DefaultAzureCredential());

        // La rubrique ecarte tout MessageId deja vu dans l'heure.
        await admin.CreateTopicAsync(new CreateTopicOptions("ventes")
        {
            RequiresDuplicateDetection = true,
            DuplicateDetectionHistoryTimeWindow = TimeSpan.FromHours(1),
        });

        // L'abonnement de la fidelite : sa copie des messages, filtree, par sessions.
        await admin.CreateSubscriptionAsync(
            new CreateSubscriptionOptions("ventes", "fidelite")
            {
                RequiresSession = true,
                MaxDeliveryCount = 5,
            },
            new CreateRuleOptions(
                "validees", new CorrelationRuleFilter { Subject = nameof(CommandeValidee) }));
    }

    // L'emetteur vient de client.CreateSender("ventes"), appele une fois au demarrage :
    // client et emetteur vivent aussi longtemps que l'application.
    public static Task Publier(
        ServiceBusSender emetteur, Guid messageId, CommandeValidee evenement) =>
        emetteur.SendMessageAsync(new ServiceBusMessage(JsonSerializer.Serialize(evenement))
        {
            MessageId = messageId.ToString(),        // la cle de la detection des doublons
            SessionId = evenement.Client.ToString(), // l'ordre est garanti client par client
            Subject = nameof(CommandeValidee),
            ContentType = "application/json",
        });

    public static async Task<ServiceBusSessionProcessor> Ecouter(
        ServiceBusClient client, Func<Guid, CommandeValidee, Task> traiter)
    {
        var options = new ServiceBusSessionProcessorOptions
        {
            AutoCompleteMessages = false,
            MaxConcurrentSessions = 8,
        };
        var processeur = client.CreateSessionProcessor("ventes", "fidelite", options);

        processeur.ProcessMessageAsync += async args =>
        {
            var evenement = Lire(args.Message.Body);
            if (evenement is null || !Guid.TryParse(args.Message.MessageId, out var id))
            {
                await args.DeadLetterMessageAsync(
                    args.Message, "Illisible", "Contenu ou MessageId invalide.");
                return;
            }

            try
            {
                await traiter(id, evenement);
                await args.CompleteMessageAsync(args.Message);
            }
            catch (Exception)
            {
                // Rendu a la file ; au-dela de MaxDeliveryCount, la file des lettres mortes.
                await args.AbandonMessageAsync(args.Message);
            }
        };
        processeur.ProcessErrorAsync += args =>
        {
            Console.Error.WriteLine($"{args.ErrorSource} : {args.Exception.Message}");
            return Task.CompletedTask;
        };

        await processeur.StartProcessingAsync();
        return processeur;
    }

    private static CommandeValidee? Lire(BinaryData corps)
    {
        try
        {
            return corps.ToObjectFromJson<CommandeValidee>();
        }
        catch (JsonException)
        {
            return null;
        }
    }
}

La réception par défaut verrouille le message (peek-lock) : il reste dans l'abonnement jusqu'à ce que le consommateur le termine, le rende ou le range aux lettres mortes. Chaque abandon, ou verrou expiré, incrémente son compteur de livraisons ; au-delà de MaxDeliveryCount, 10 par défaut, le service le déplace dans la file des lettres mortes, une sous-file de chaque file et de chaque abonnement, où il n'expire pas. Les sessions garantissent l'ordre : un récepteur accepte une session et reçoit en exclusivité tous les messages de même SessionId. Ici, les commandes d'un même client sont traitées dans l'ordre où le relais les envoie, pas forcément celui de leur validation ; chaque message destiné à l'abonnement doit donc porter un SessionId.

La détection des doublons écarte, pendant une fenêtre de 10 minutes par défaut (de 20 secondes à 7 jours), tout message dont le MessageId a déjà été vu : l'envoi réussit, le doublon est jeté. Elle protège l'envoi, le relais qui republie après une panne, pourvu que le MessageId vienne de la boîte d'envoi plutôt que d'un Guid.NewGuid() par tentative. Elle ne protège pas la réception : un verrou expiré relivre un message déjà traité, et le consommateur idempotent reste nécessaire. Le niveau Basic n'offre ni rubriques, ni sessions, ni détection des doublons : l'exemple demande au moins Standard.

NotionRabbitMQAzure Service Bus
Aiguillageéchanges et liaisons, déclarés par les applicationsrubrique, abonnements et leurs règles
Règlement d'un messageack, reject, nackcomplete, abandon, dead-letter
Lettres mortesun échange à déclarer, par fileune sous-file intégrée à chaque file et abonnement
Limite de livraisonsfiles quorum, 20 par défaut depuis 4.0MaxDeliveryCount, 10 par défaut
Doublons à l'envoiaucune détection pour les files classiques et quorumpar MessageId, sur une fenêtre
Ce cours vous a servi ? Offrir un café Signaler une erreur