Le problème

Prenons une boutique en ligne découpée en services. Le service order enregistre une commande et publie l'événement order.order.placed.v1. Deux services doivent y réagir : stock réserve la marchandise, shipping prépare l'expédition.

Chacun doit recevoir sa propre copie de l'événement. La topologie RabbitMQ qui porte ce flux se dégrade souvent de l'une de ces trois façons.

Une queue partagée entre plusieurs services. stock et shipping consomment la même queue order-placed. RabbitMQ répartit alors les messages entre les consommateurs connectés. Chaque message n'est traité que par un seul service : stock ne réserve que pour une partie des commandes, shipping ne prépare que les autres.

Une queue par cas d'usage. Pour corriger, on crée stock.reserve-on-order-placed, puis shipping.prepare-on-order-placed, puis une queue à chaque nouveau besoin. Le nom de chaque queue décrit une réaction métier. La topologie grossit à chaque fonctionnalité et l'émetteur se retrouve couplé aux usages de ses événements.

Une exchange par domaine. order, stock et shipping publient chacun dans leur propre exchange. Un service qui écoute plusieurs domaines multiplie les bindings croisés. Sans ouvrir l'interface de gestion, plus personne ne sait dire quel message arrive dans quelle queue.

Les options

Une queue partagée par événement

  • Ce qu'elle règle : rien pour la diffusion. Elle convient quand plusieurs instances d'un même service se partagent le travail.
  • Ce qu'elle coûte : chaque message n'atteint qu'un consommateur. Brancher un nouveau service sur la queue lui fait « voler » une partie des messages des autres.

Une queue par cas d'usage

  • Ce qu'elle règle : chaque réaction reçoit sa copie du message.
  • Ce qu'elle coûte : une queue, un binding et souvent une politique de retry par fonctionnalité. Le nombre de queues suit le nombre de fonctionnalités, pas le nombre de services.

Une exchange par domaine

  • Ce qu'elle règle : chaque domaine publie dans un espace qui lui est propre.
  • Ce qu'elle coûte : un consommateur se lie à plusieurs exchanges, et les bindings se dispersent. La clé de routage répète de toute façon le domaine : l'exchange n'apporte pas d'information de plus.

Une seule topic exchange, une queue par service consommateur

  • Ce qu'elle règle : chaque service intéressé reçoit sa copie. Les producteurs ont une seule porte d'entrée, et chaque service possède une queue à son nom.
  • Ce qu'elle coûte : une convention de nommage stricte, et une vue d'ensemble des abonnements à construire (voir les compromis plus bas).

Notre recommandation

Une seule exchange topic pour tout le système, une queue par service consommateur, et des binding keys qui disent ce que chacun reçoit.

Une topic exchange copie chaque message dans toutes les queues dont une binding key correspond à sa clé de routage. Le producteur publie une fois, et RabbitMQ remet une copie à chaque service intéressé. Dans une queue, les instances d'un même service se partagent ensuite le travail : c'est là que la répartition de charge a sa place.

Une exchange, une queue par service : chaque service reçoit sa copie de l'événement order service producteur publie order.order.placed.v1 app.events topic exchange stock.events binding : order.order.placed.v1 consommée par le service stock shipping.events binding : order.order.*.v1 consommée par le service shipping À éviter : une queue partagée répartit les messages au lieu de les diffuser order-placed queue partagée stock shipping stock traite les messages 1, 3, 5... shipping traite les messages 2, 4, 6...

La clé de routage est le nom de l'événement

Nous suivons la convention <contexte>.<entité>.<fait au passé>.v<N> :

  • order.order.placed.v1 et order.order.cancelled.v1 ;
  • stock.reservation.confirmed.v1 ;
  • shipping.shipment.dispatched.v1.

Le contexte est le service qui possède le fait. Le verbe au passé rappelle qu'un événement décrit ce qui s'est produit, pas ce qu'il faut faire. La version permet de publier une v2 à côté de la v1 pendant une transition.

Une queue par service, nommée d'après lui

Les queues s'appellent stock.events, shipping.events, billing.events. Chaque service déclare sa queue, choisit ses binding keys et sa politique de retry. Ajouter un consommateur ne demande aucune modification au producteur.

Des binding keys précises

Dans une binding key, * remplace exactement un mot et # remplace zéro ou plusieurs mots :

  • order.order.placed.v1 abonne à un seul événement ;
  • order.order.*.v1 abonne à tous les faits de la version 1 de la commande ;
  • # abonne à tout le trafic.

Évitez #. Le service recevrait des événements qu'il ne sait pas traiter. Et chaque nouvel événement du système arriverait dans sa queue sans que personne ne l'ait décidé.

Le joker * a le même effet, à plus petite échelle : chaque nouvel événement order.order.<fait>.v1 arrivera chez shipping. Réservez-le à une famille d'événements que le service traite en entier.

Événement, commande, requête

Seuls les événements passent par la topic exchange :

  • un événement est un fait passé, diffusé à qui veut l'entendre ;
  • une commande, au sens CQRS, est une demande adressée à un service précis (« réserve ce stock »), en point à point ;
  • une requête attend une réponse immédiate : elle reste synchrone, en HTTP.

Une commande diffusée par la topic exchange pourrait être exécutée par deux services, ou par aucun si personne ne s'est abonné.

Quand découper la queue d'un service

Une queue par service est un point de départ. Nous ne la découpons que sur un signal concret :

  • des politiques de retry ou de durée de vie (TTL) qui divergent entre deux familles d'événements ;
  • des messages lourds qui retardent les messages critiques placés derrière eux (head-of-line blocking), constaté sur les délais de traitement ;
  • une contrainte d'ordre sur un sous-ensemble des événements, qui impose un seul worker sur la queue qui les porte ;
  • le besoin de faire monter en charge un traitement séparément des autres.

Les queues issues d'un découpage restent préfixées par le service : stock.events.priority, stock.events.bulk. On sait toujours à qui elles appartiennent.

Les compromis assumés

Un seul espace de noms. Toutes les clés de routage passent par la même exchange. Une faute de frappe dans une clé passe inaperçue : le message part, aucune queue ne le reçoit, et RabbitMQ l'écarte. Nous l'encadrons par la convention de nommage, par des constantes partagées et, si besoin, par une alternate exchange qui recueille les messages non routés.

Les abonnements vivent chez les consommateurs. Savoir qui consomme quoi demande de lire la configuration de chaque service, ou l'interface de gestion de RabbitMQ. Prévoyez un catalogue : une page de documentation tenue à jour, ou un export généré depuis la configuration.

Une queue mélange les événements d'un service. Tant que leurs traitements se ressemblent, c'est un avantage : une seule queue à surveiller par service. Le jour où ils divergent, il faut découper, avec les critères vus plus haut.

Le signal qui doit faire réexaminer ce choix : un retard mesuré sur des messages critiques, causé par d'autres messages de la même queue. Ou un besoin régulier de rejouer l'historique, qui relève plutôt d'un journal : Kafka, ou un stream RabbitMQ lu par un autre client que Messenger.

Mise en œuvre

Les exemples utilisent Symfony 7.4 et le transport AMQP de Messenger (symfony/amqp-messenger). Les options sont les mêmes en Symfony 8.

1. Partager les noms, pas les queues

Le nom de l'exchange et les clés de routage sont partagés par tous les services, dans un package de contrats. Les noms de queues et les binding keys restent chez chaque consommateur.

// app-contracts/src/Routing.php
namespace AppContracts;

final class Routing
{
    public const EXCHANGE = 'app.events';

    public const ORDER_PLACED_V1 = 'order.order.placed.v1';
    public const ORDER_CANCELLED_V1 = 'order.order.cancelled.v1';
    public const STOCK_RESERVATION_CONFIRMED_V1 = 'stock.reservation.confirmed.v1';
}

2. Côté producteur : l'exchange, sans queue

Le service order déclare la topic exchange et aucune queue. Avec auto_setup, actif par défaut, Messenger déclare et lie toutes les queues configurées, y compris côté émetteur. queues: [] évite de créer une queue que personne ne consommerait.

# config/packages/messenger.yaml (service order)
framework:
  messenger:
    transports:
      events:
        dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
        options:
          exchange:
            name: !php/const AppContracts\Routing::EXCHANGE
            type: topic
          # Pas de queue côté producteur : chaque consommateur possède la sienne.
          queues: []
    routing:
      AppContracts\Order\OrderPlacedV1: events
      AppContracts\Order\OrderCancelledV1: events

Le tag YAML !php/const lit les constantes du package de contrats : un nom mal orthographié devient une erreur au démarrage, pas un message perdu.

La clé de routage se pose sur l'enveloppe Messenger avec un AmqpStamp. Sans lui, Messenger utilise default_publish_routing_key, ou aucune clé : une topic exchange ne routerait alors le message vers aucune queue.

// src/Messenger/EventPublisher.php (service order)
namespace App\Messenger;

use Symfony\Component\Messenger\Bridge\Amqp\Transport\AmqpStamp;
use Symfony\Component\Messenger\MessageBusInterface;

final class EventPublisher
{
    public function __construct(
        private readonly MessageBusInterface $bus,
    ) {
    }

    public function publish(object $event, string $routingKey): void
    {
        $this->bus->dispatch($event, [new AmqpStamp($routingKey)]);
    }
}

L'appel s'écrit $publisher->publish($event, Routing::ORDER_PLACED_V1). En production, nous ne publions pas depuis la requête HTTP. Un relais lit une table outbox, écrite dans la même transaction que la commande : aucun événement ne se perd en route.

3. Côté consommateur : sa queue, ses binding keys

Le service stock déclare la même exchange et sa propre queue, avec les événements qu'il traite.

# config/packages/messenger.yaml (service stock)
framework:
  messenger:
    failure_transport: failed
    transports:
      events:
        dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
        options:
          exchange:
            name: !php/const AppContracts\Routing::EXCHANGE
            type: topic
          queues:
            stock.events:
              binding_keys:
                - !php/const AppContracts\Routing::ORDER_PLACED_V1
                - !php/const AppContracts\Routing::ORDER_CANCELLED_V1
        retry_strategy:
          max_retries: 3
          delay: 1000
          multiplier: 2
      failed: 'doctrine://default?queue_name=failed'

Le service shipping fait de même avec la queue shipping.events et la binding key order.order.*.v1. Chaque handler reçoit le DTO de l'événement :

// src/MessageHandler/ReserveStockOnOrderPlaced.php (service stock)
namespace App\MessageHandler;

use AppContracts\Order\OrderPlacedV1;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;

#[AsMessageHandler]
final class ReserveStockOnOrderPlaced
{
    public function __invoke(OrderPlacedV1 $event): void
    {
        // Réserve le stock de chaque ligne de la commande $event->orderId.
    }
}

Le worker se lance avec php bin/console messenger:consume events. Un événement lié sans handler lève une exception côté consommateur : ne liez que ce que le service traite réellement.

Ces exemples gardent le sérialiseur par défaut de Messenger, qui encode la classe PHP du message. Producteur et consommateur partagent alors ces classes, par le package app-contracts. En production, nous publions une enveloppe JSON, décodée par un sérialiseur dédié (option serializer du transport) : le format ne dépend plus du code PHP.

4. Retry et quarantaine, service par service

Avec ce transport, un retry ne repasse pas par la topic exchange. Messenger republie le message dans une queue de délai, qui le renvoie ensuite à la seule queue en échec. Un échec dans stock ne redistribue donc pas l'événement à shipping.

Après max_retries nouvelles tentatives, le message part dans le failure transport du service. Les commandes messenger:failed:show et messenger:failed:retry permettent de l'examiner et de le rejouer. Sans failure transport, le message serait rejeté et perdu, sauf dead letter exchange configurée sur la queue.

5. Créer et vérifier la topologie

Une queue ne reçoit que les messages publiés après sa création, et le service order n'en déclare aucune. Nous lançons donc php bin/console messenger:setup-transports au déploiement de chaque consommateur, avant que le producteur n'émette les événements qu'il attend. Nous vérifions ensuite les bindings :

rabbitmqctl list_bindings source_name destination_name routing_key

La sortie doit montrer app.events liée à stock.events et shipping.events, avec les binding keys attendues, et aucune binding key #.

Checklist

  • Une seule exchange de type topic pour les événements.
  • Une queue par service consommateur, nommée d'après lui (stock.events).
  • Des clés de routage versionnées : <contexte>.<entité>.<fait au passé>.v<N>.
  • Le nom de l'exchange et les clés de routage dans des constantes partagées.
  • Aucune binding key #.
  • Un handler pour chaque événement que les binding keys laissent passer, jokers compris.
  • Un failure transport (ou une DLQ) par service, et une procédure pour le traiter.
  • Les commandes et les requêtes hors de l'exchange d'événements.
  • Un découpage de queue seulement sur un signal mesuré.

Sources