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.

Deux livraisons du même événement : une seule réservation de stock 1re livraison Message reçu de stock.events Transaction ( doctrine_transaction ) INSERT processed_event 1 ligne Handler : réservation créée COMMIT : tout est validé Crash avant l'acquittement Non acquitté : il revient en retry 2e livraison : même event_id Message reçu de stock.events Transaction ( doctrine_transaction ) INSERT processed_event 0 ligne Handler non appelé COMMIT : rien de nouveau Acquittement : message retiré Une seule réservation au total

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_id stable, posé par le producteur dans l'enveloppe JSON, jamais généré à la réception.
  • Une table processed_event (event_id UUID 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_transaction dé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