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.
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.v1etorder.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.v1abonne à un seul événement ; -
order.order.*.v1abonne à 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
topicpour 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
- RabbitMQ, tutoriel Work Queues : distribution des messages à tour de rôle entre les consommateurs d'une queue.
-
RabbitMQ, tutoriel
Topics
: clés de routage, jokers
*et#. - RabbitMQ, Alternate Exchanges : recueillir les messages non routés.
- RabbitMQ, Streams : un journal rejouable dans RabbitMQ.
-
RabbitMQ,
rabbitmqctl
: commande
list_bindings. -
Symfony,
Messenger, transport AMQP
: options
exchange,queues,binding_keys,auto_setup. -
Symfony,
Messenger, retries et échecs
:
retry_strategyetfailure_transport. - Symfony, code du transport AMQP : déclaration des queues et routage des retries.
- Enterprise Integration Patterns, Publish-Subscribe Channel et Competing Consumers.