Le problème
Reprenons notre boutique en ligne. Le service
order enregistre une commande et publie
order.order.placed.v1 par un relais d'outbox.
Le service stock lit sa queue
stock.events et réserve la marchandise de
chaque ligne.
Avec des acquittements, RabbitMQ garantit une livraison « au
moins une fois ». Il ne garantit pas « une seule fois ». Le
même événement peut donc arriver plusieurs fois chez
stock, pour des raisons ordinaires :
- le relais republie : il a publié la ligne d'outbox, puis s'est arrêté avant de la marquer comme publiée. Au redémarrage, il la publie de nouveau ;
- le broker relivre : le worker a validé sa transaction, puis a perdu sa connexion avant d'acquitter. RabbitMQ remet alors le message dans la queue ;
- quelqu'un rejoue : une période d'événements est republiée après un incident. Ou un message est repris à la main alors qu'une autre copie a déjà été traitée.
Sans précaution, chaque livraison produit son effet. Le stock est réservé deux fois, le client reçoit deux avis d'expédition. Les retries n'y peuvent rien : ils relancent un traitement, ils ne savent pas qu'il a déjà réussi.
Le drapeau redelivered de RabbitMQ ne suffit
pas à repérer les doublons. Un événement republié par le
relais arrive comme un message neuf, sans ce drapeau.
Côté Symfony, un message relivré par RabbitMQ n'est pas traité directement. Un middleware par défaut de Messenger le renvoie dans le circuit de retry : le doublon revient donc plus tard.
Les options
L'idempotence naturelle
- Ce qu'elle règle : un effet qui pose un état cible (« la commande est expédiée ») peut être rejoué sans dommage. Un upsert sur une clé en est l'exemple type.
- Ce qu'elle coûte : elle ne couvre pas les effets cumulatifs. Réserver deux unités, décrémenter un compteur ou envoyer un e-mail ne sont pas des états cibles.
La déduplication par clé métier
- Ce qu'elle règle : une contrainte d'unicité porte sur l'effet lui-même, par exemple une seule réservation par commande et par produit.
- Ce qu'elle coûte : il faut trouver une clé pour chaque effet, et il n'en existe pas toujours. La règle de déduplication se réécrit alors dans chaque handler.
Une table des événements traités
- Ce qu'elle règle : chaque service consommateur note l'identifiant des événements qu'il a traités, dans sa propre base. Un doublon est reconnu avant tout effet, quel que soit le handler.
- Ce qu'elle coûte : une écriture de plus par message, une table par service consommateur et une purge à organiser.
La déduplication par le broker ou par Messenger
- Ce qu'elle règle : écarter les doublons avant qu'ils n'atteignent le consommateur.
- Ce qu'elle coûte : RabbitMQ n'en propose pas pour les queues classiques ou quorum. Les streams dédupliquent à la publication, mais notre topologie n'en utilise pas. Un plugin communautaire, hors distribution officielle, ne retient les identifiants que dans un cache limité (exchange) ou tant que l'original attend dans la queue.
Symfony 7.3 a ajouté un DeduplicateStamp à
Messenger. Il pose un verrou à l'envoi et le relâche à la
fin du traitement : une copie arrivée après coup passe. Dans
tous les cas, la seule trace fiable du traitement est la
base du consommateur, validée avec l'effet.
Notre recommandation
Nous recommandons une table
processed_event dans la base du consommateur,
remplie avant tout effet métier par
INSERT ... ON CONFLICT DO NOTHING, dans la
même transaction que cet effet.
L'identifiant vient du producteur
L'event_id est un UUID v7 posé par le
producteur quand il écrit sa ligne d'outbox. Il voyage dans
l'enveloppe JSON du message, à côté de
event_name, occurred_at et
payload. Le consommateur ne le génère jamais :
un identifiant créé à la réception serait neuf à chaque
livraison.
La même transaction que l'effet
Le middleware doctrine_transaction de Messenger
ouvre une transaction avant les handlers. S'ils réussissent,
il appelle flush() puis commit().
Sur une exception, il annule tout : la ligne de
processed_event disparaît avec les
réservations, et le retry retraite l'événement.
Si le commit passe mais que l'acquittement se perd, la livraison suivante trouve la ligne et s'arrête.
Éviter l'erreur plutôt que l'attraper
Un INSERT simple lèverait une violation
d'unicité sur un doublon. L'attraper ne suffit pas : après
une erreur, PostgreSQL refuse toute autre instruction de la
transaction (code 25P02) jusqu'au rollback.
ON CONFLICT DO NOTHING ne lève pas d'erreur, et
executeStatement() de Doctrine DBAL renvoie le
nombre de lignes insérées : zéro pour un doublon.
Si deux workers reçoivent le même événement au même moment,
la clé primaire les départage. Le second
INSERT attend la fin de la transaction du
premier. Si elle est validée, il n'insère rien ; si elle est
annulée, il insère et traite.
Dans un middleware, pas dans chaque handler
Messenger ne passe au handler que le message : l'event_id
reste sur l'enveloppe Messenger, dans un stamp. Un
middleware voit cette enveloppe et s'exécute une fois par
message, avant tous ses handlers. Insérée dans chaque
handler, la ligne gênerait un second handler du même
événement : il la trouverait et sortirait sans rien faire.
Les compromis assumés
Une écriture de plus par message. Chaque événement traité ajoute une ligne et une entrée d'index dans la base du consommateur. Sur une queue à fort trafic, ce coût se surveille.
Une table par service consommateur, et une purge.
stock, shipping et
billing ont chacun leur table
processed_event. Sans purge, elle grossit sans
fin ; avec une purge trop courte, les doublons reviennent
(étape 8).
La clé est l'event_id à l'échelle du service.
Si un service lie le même événement à deux de ses
transports, chacun avec ses handlers
(fromTransport), la seconde copie est écartée.
Seuls les effets en base sont protégés. La transaction couvre ce que le handler écrit dans PostgreSQL. Un e-mail envoyé ou un appel HTTP n'est pas annulé par un rollback : si le commit échoue ensuite, le retry le refera.
Ces effets demandent leur propre clé d'idempotence, ou un message distinct envoyé après le commit.
L'ordre n'est pas traité. La table dit si un événement a déjà été traité, pas s'il arrive trop tard. Pour une copie d'état, le contrôle de version de l'étape 6 couvre ce cas.
Le signal qui doit faire réexaminer ce choix
: l'insertion dans processed_event apparaît
dans les mesures comme le goulot du traitement. Ou bien la
plupart des effets d'un service vivent hors de sa base, et
la table ne les protège plus.
Mise en œuvre
Les exemples utilisent Symfony 7.4, Doctrine ORM 3 avec DBAL
4, PostgreSQL et le transport AMQP de Messenger. Les classes
citées existent aussi en Symfony 8. Le transport
events et la queue
stock.events sont ceux de l'article
RabbitMQ : une seule topic exchange et une queue par
service consommateur.
1. Garder l'event_id dans un stamp
Le sérialiseur du transport lit l'enveloppe JSON. Il
construit le DTO à partir de payload et range
event_id dans un stamp.
// src/Messenger/EventIdStamp.php (service stock)
namespace App\Messenger;
use Symfony\Component\Messenger\Stamp\StampInterface;
final class EventIdStamp implements StampInterface
{
public function __construct(
public readonly string $eventId,
) {
}
}
// src/Messenger/EventSerializer.php (service stock, extrait)
use Symfony\Component\Messenger\Exception\MessageDecodingFailedException;
// ...
public function decode(array $encodedEnvelope): Envelope
{
try {
$body = json_decode($encodedEnvelope['body'], true, flags: \JSON_THROW_ON_ERROR);
// toEvent() associe event_name à la classe du DTO (OrderPlacedV1...).
$event = $this->toEvent($body['event_name'], $body['payload']);
} catch (\Throwable $e) {
// Le receiver AMQP rejette alors le message, sans le remettre dans la queue.
throw new MessageDecodingFailedException($e->getMessage(), 0, $e);
}
// decodeStamps() refait le RedeliveryStamp depuis l'en-tête écrit par encode().
$stamps = $this->decodeStamps($encodedEnvelope['headers'] ?? []);
// Les autres champs de l'enveloppe JSON vont dans un second stamp, omis ici.
return new Envelope($event, [new EventIdStamp($body['event_id']), ...$stamps]);
}
Lors d'un retry, Messenger réencode le message avec ce
sérialiseur. encode() doit donc réécrire
event_id depuis le stamp, pour que la copie
garde son identifiant.
encode() et decode() doivent aussi
transporter le compteur de tentatives dans un en-tête, et
decode() en refaire un
RedeliveryStamp. Sans lui, le compteur de retry
repart de zéro et le message n'atteint jamais le failure
transport. Celui-ci garde ensuite l'enveloppe Messenger
entière, EventIdStamp compris, avec son
sérialiseur par défaut.
2. Déclarer la table
L'entité sert au schéma : elle déclare la table pour les
migrations. Les insertions passent par DBAL, car le DQL ne
connaît pas INSERT, et un
flush() en doublon échouerait.
// src/Entity/ProcessedEvent.php (service stock)
namespace App\Entity;
use Doctrine\DBAL\Types\Types;
use Doctrine\ORM\Mapping as ORM;
use Symfony\Bridge\Doctrine\Types\UuidType;
use Symfony\Component\Uid\Uuid;
#[ORM\Entity]
#[ORM\Table(name: 'processed_event')]
#[ORM\Index(columns: ['processed_at'])]
class ProcessedEvent
{
#[ORM\Id]
#[ORM\Column(name: 'event_id', type: UuidType::NAME)]
private Uuid $eventId;
#[ORM\Column(name: 'processed_at', type: Types::DATETIMETZ_IMMUTABLE)]
private \DateTimeImmutable $processedAt;
}
php bin/console make:migration génère la
migration de la table : event_id de type
UUID natif en clé primaire, et un index sur
processed_at pour la purge.
3. Écrire le middleware
// src/Messenger/ProcessedEventMiddleware.php (service stock)
namespace App\Messenger;
use Doctrine\DBAL\Connection;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Middleware\MiddlewareInterface;
use Symfony\Component\Messenger\Middleware\StackInterface;
use Symfony\Component\Messenger\Stamp\ReceivedStamp;
final class ProcessedEventMiddleware implements MiddlewareInterface
{
public function __construct(
private readonly Connection $connection,
) {
}
public function handle(Envelope $envelope, StackInterface $stack): Envelope
{
$eventId = $envelope->last(EventIdStamp::class)?->eventId;
// Message envoyé par ce service, ou sans event_id : rien à dédupliquer.
if (null === $envelope->last(ReceivedStamp::class) || null === $eventId) {
return $stack->next()->handle($envelope, $stack);
}
if (!$this->connection->isTransactionActive()) {
throw new \LogicException('doctrine_transaction doit précéder ProcessedEventMiddleware.');
}
$inserted = $this->connection->executeStatement(
'INSERT INTO processed_event (event_id, processed_at) VALUES (:eventId, now())
ON CONFLICT (event_id) DO NOTHING',
['eventId' => $eventId],
);
if (0 === (int) $inserted) {
// Déjà traité : aucun handler n'est appelé, le worker acquitte le message.
return $envelope;
}
return $stack->next()->handle($envelope, $stack);
}
}
Le middleware s'exécute aussi à l'envoi des messages. Le
test sur ReceivedStamp le limite aux messages
reçus d'un transport.
4. Configurer le bus
# config/packages/messenger.yaml (service stock, extrait)
framework:
messenger:
buses:
messenger.bus.default:
middleware:
# Ouvre la transaction ; flush() et commit si tout réussit, rollback sinon.
- doctrine_transaction
# Écarte les événements déjà traités, dans cette transaction.
- App\Messenger\ProcessedEventMiddleware
transports:
events:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
serializer: App\Messenger\EventSerializer
# options exchange et queues : voir l'article sur la topologie
L'ordre de la liste compte : Messenger exécute ces middleware dans l'ordre déclaré, avant d'appeler les handlers. Le worker remet au bus par défaut un message reçu sans nom de bus. C'est le cas d'une enveloppe JSON venue d'un autre service.
5. Écrire le handler
Le handler ne contient que l'effet métier. Il n'appelle pas
flush() :
doctrine_transaction s'en charge avant le
commit. Reservation est une entité Doctrine
classique, avec un id UUID v7 généré dans son
constructeur.
// src/MessageHandler/ReserveStockOnOrderPlaced.php (service stock)
namespace App\MessageHandler;
use App\Entity\Reservation;
use AppContracts\Order\OrderPlacedV1;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;
use Symfony\Component\Uid\Uuid;
#[AsMessageHandler]
final class ReserveStockOnOrderPlaced
{
public function __construct(
private readonly EntityManagerInterface $entityManager,
) {
}
public function __invoke(OrderPlacedV1 $event): void
{
foreach ($event->lines as $line) {
$reservation = new Reservation();
$reservation->setOrderId(Uuid::fromString($event->orderId));
$reservation->setProductId(Uuid::fromString($line['productId']));
$reservation->setQuantity($line['quantity']); // chaîne décimale, '2.000'
$this->entityManager->persist($reservation);
}
}
}
6. Une copie d'état : ajouter un contrôle de version
Le service shipping garde une copie du statut
de chaque commande, pour savoir si un colis peut partir. Si
le producteur publie un numéro de version de la commande,
incrémenté à chaque changement, l'upsert n'applique que les
versions plus récentes.
-- service shipping : copie locale du statut des commandes
INSERT INTO order_status (order_id, status, version)
VALUES (:orderId, :status, :version)
ON CONFLICT (order_id) DO UPDATE
SET status = EXCLUDED.status, version = EXCLUDED.version
WHERE order_status.version < EXCLUDED.version;
Le middleware écarte déjà les doublons. Le contrôle de version écarte en plus un événement arrivé après un plus récent.
7. Tester la double livraison
Le test passe par le bus, pas par le handler seul : il
vérifie ensemble la transaction, le middleware et le
handler. Le ReceivedStamp fait traiter le
message sur place, comme s'il sortait de la queue.
// tests/Messenger/DoubleDeliveryTest.php (service stock)
namespace App\Tests\Messenger;
use App\Entity\Reservation;
use App\Messenger\EventIdStamp;
use AppContracts\Order\OrderPlacedV1;
use Doctrine\DBAL\Connection;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Bundle\FrameworkBundle\Test\KernelTestCase;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Messenger\Stamp\ReceivedStamp;
use Symfony\Component\Uid\Uuid;
final class DoubleDeliveryTest extends KernelTestCase
{
public function testTheSameEventDeliveredTwiceReservesStockOnce(): void
{
$container = static::getContainer();
$orderId = Uuid::v7();
$eventId = (string) Uuid::v7();
$event = new OrderPlacedV1(
orderId: (string) $orderId,
customerId: (string) Uuid::v7(),
lines: [['productId' => (string) Uuid::v7(), 'quantity' => '2.000']],
);
$bus = $container->get(MessageBusInterface::class);
$envelope = new Envelope($event, [new EventIdStamp($eventId), new ReceivedStamp('events')]);
// Deux livraisons du même événement, comme après un crash avant l'acquittement.
$bus->dispatch($envelope);
$bus->dispatch($envelope);
$reservations = $container->get(EntityManagerInterface::class)
->getRepository(Reservation::class)
->count(['orderId' => $orderId]);
self::assertSame(1, $reservations);
$processed = $container->get(Connection::class)->fetchOne(
'SELECT count(*) FROM processed_event WHERE event_id = :eventId',
['eventId' => $eventId],
);
self::assertSame(1, (int) $processed);
}
}
Un second test vaut la peine : faire échouer le handler,
puis vérifier que processed_event ne contient
pas l'event_id. C'est ce qui permet au retry de
retraiter l'événement.
8. Purger selon l'horizon de rejeu
-- tâche planifiée du service stock
DELETE FROM processed_event WHERE processed_at < :horizon;
Une ligne peut partir quand plus aucune copie de son événement ne peut arriver. Cet horizon dépend de la durée pendant laquelle les producteurs peuvent rejouer leur outbox, et du séjour possible d'un message dans le failure transport. Purger en deçà, c'est laisser un rejeu retraiter des événements déjà appliqués.
Checklist
-
Un
event_idstable, posé par le producteur dans l'enveloppe JSON, jamais généré à la réception. -
Une table
processed_event(event_idUUID en clé primaire,processed_at) dans la base de chaque service consommateur. - L'insertion avant tout effet métier, dans la même transaction que lui.
-
doctrine_transactiondéclaré avant le middleware de déduplication. -
INSERT ... ON CONFLICT DO NOTHING, et aucun handler appelé si aucune ligne n'est insérée. - Les effets hors base (e-mail, appel HTTP) protégés par leur propre clé.
- Un contrôle de version pour les copies d'état sensibles à l'ordre.
- Un test qui livre deux fois le même message et vérifie un seul effet.
- Une purge alignée sur l'horizon de rejeu.
Sources
- Symfony, Messenger : Writing Idempotent Handlers, Middleware for Doctrine, Middleware, Message Deduplication (Symfony 7.3), Message Serializer For Custom Data Formats.
- Symfony, code de DoctrineTransactionMiddleware et de RejectRedeliveredMessageMiddleware, branche 7.4.
-
PostgreSQL :
INSERT, clause ON CONFLICT,
Index Uniqueness Checks,
Transactions
et
codes d'erreur
(
25P02). - Doctrine : DBAL, executeStatement() et ORM, pas d'INSERT en DQL.
-
RabbitMQ :
Automatic Requeueing,
Reliability Guide
(livraison au moins une fois, drapeau
redelivered) et déduplication des streams. - rabbitmq-message-deduplication, plugin communautaire.