Le problème

Dans notre boutique en ligne, le service order expose POST /orders/{id}/place : la commande passe du statut draft à placed. Il publie alors order.order.placed.v1, pour que stock réserve la marchandise et que shipping prépare l'expédition.

Deux écritures sont en jeu : la commande dans PostgreSQL, l'événement dans RabbitMQ. Elles ne partagent aucune transaction, et l'ordre choisi décide de la panne.

Écrire, puis publier. Le commit réussit, puis le processus s'arrête avant de publier : redéploiement, manque de mémoire, broker injoignable. La commande est passée, mais aucun service ne le sait : le stock ne sera jamais réservé.

Publier, puis écrire. L'événement part, puis la transaction échoue sur une contrainte ou un conflit. stock réserve la marchandise d'une commande restée en draft : c'est un événement fantôme.

Deux requêtes simultanées. Un double clic, ou un client qui réessaie, envoie deux fois la même requête. Les deux lisent draft, appliquent la transition et publient chacune leur événement. stock réserve deux fois.

Les options

Publier depuis la requête, avant ou après le commit

  • Ce qu'elle règle : rien à ajouter, ni table ni processus.
  • Ce qu'elle coûte : une perte ou un événement fantôme, selon l'ordre choisi. Dans un handler sous doctrine_transaction, le DispatchAfterCurrentBusStamp envoie après le commit : il évite le fantôme, pas la perte.

Une transaction distribuée

  • Ce qu'elle règle : en théorie, la base et le broker valident ensemble.
  • Ce qu'elle coûte : RabbitMQ n'y participe pas. Ses transactions AMQP ne couvrent que les publications et les acquittements d'un canal, sans lien avec PostgreSQL.

Une outbox lue par un relais

  • Ce qu'elle règle : l'événement est inséré dans une table, dans la même transaction que la commande. Un relais lit cette table et publie. Rien ne se perd, rien n'est publié à tort.
  • Ce qu'elle coûte : une table, un processus, une latence liée à l'intervalle de lecture et des doublons possibles.

Le transport Doctrine de Messenger peut servir d'outbox. Mais ses retries rompent l'ordre, et il supprime par défaut les messages traités.

Une outbox lue par capture des changements (CDC)

  • Ce qu'elle règle : Debezium lit le journal de transactions (WAL) de PostgreSQL par décodage logique. Sa transformation « Outbox Event Router » fait de chaque ligne insérée un message.
  • Ce qu'elle coûte : une infrastructure à exploiter, Kafka Connect ou Debezium Server et son sink RabbitMQ. Et un slot de réplication retient le WAL tant que le connecteur ne l'a pas lu.

Notre recommandation

Écrire l'événement dans une table outbox, dans la même transaction que la commande, puis le publier par un relais qui lit cette table.

La transaction PostgreSQL est notre seule atomicité : la commande et son événement y entrent ensemble. Nous passons au CDC si la lecture périodique pèse sur la base, mesures à l'appui.

Écrit avec la commande, publié ensuite : au pire un doublon, jamais une perte Client order API, processor PostgreSQL orders, outbox Relais Scheduler RabbitMQ app.events Transaction de la requête POST /orders/{id}/place SELECT ... FOR UPDATE transition place listener : ligne d'outbox UPDATE + INSERT COMMIT 200 Transaction du relais FOR UPDATE SKIP LOCKED publier confirmation published_at COMMIT Arrêt entre la confirmation et le COMMIT : la ligne n'est pas marquée. Elle repart au lot suivant : le consommateur reçoit un doublon, jamais une perte.

Écrire l'outbox sur le fait métier

L'événement naît quand la commande passe à placed. Avec Symfony Workflow, c'est l'événement workflow.order.completed.place, émis une fois la transition appliquée. Un listener y insère la ligne d'outbox.

Nous évitons les listeners Doctrine postPersist ou postUpdate : ils voient des colonnes modifiées, pas une intention. Corriger une adresse y déclenche le même hook que passer une commande. Et Doctrine n'écrit pas, dans le flush en cours, ce qu'on y persiste.

Nous n'écrivons pas non plus l'outbox dans les processors API Platform : une transition lancée en console n'émettrait aucun événement.

Verrouiller la commande avant la transition

Le processor lit la commande avec LockMode::PESSIMISTIC_WRITE, soit un SELECT ... FOR UPDATE. La seconde requête simultanée attend la fin de la première, relit placed et échoue proprement. Un verrou de ligne ne bloque pas les lectures simples.

Préparer l'événement avant le commit

Les identifiants sont des UUID v7 générés en PHP : l'événement porte celui de la commande et le sien avant le flush(). Avant l'insertion, nous validons l'enveloppe contre son JSON Schema. Un message invalide annule toute la transaction : l'erreur apparaît chez le producteur, pas chez un consommateur.

Publier, marquer, commiter

Le relais lit un lot de lignes non publiées avec FOR UPDATE SKIP LOCKED. Pour chaque ligne, il publie, attend la confirmation du broker (publisher confirm), puis renseigne published_at. Il commite à la fin du lot.

Cet ordre ne perd rien. Un arrêt entre la publication et le commit annule le marquage, et la ligne repart au lot suivant. Le consommateur reçoit alors un doublon.

La confirmation signifie que le broker a pris le message en charge, sur disque s'il est persistant et sa queue durable. Un message qu'aucune binding key ne retient est confirmé aussi, puis écarté. Une alternate exchange peut le recueillir : voir notre article sur la topologie RabbitMQ.

Garder les lignes publiées

Les lignes publiées forment un journal daté, utile pour enquêter ou rejouer un événement. Une tâche planifiée purge ce qui dépasse l'horizon de rejeu que vous fixez.

Les compromis assumés

Une latence liée à la lecture périodique. Un événement attend le passage suivant du relais. Réserver du stock tolère ce délai ; un affichage en temps réel, moins.

Une table et un processus en plus. Nous alertons sur l'âge de la plus ancienne ligne non publiée et sur la colonne attempts.

Des doublons. Le relais livre « au moins une fois ». Chaque consommateur doit être idempotent, par exemple avec une table des événements traités : c'est le sujet d'un article dédié.

Un ordre limité. L'ordre repose sur occurred_at, un relais unique et l'arrêt au premier échec : une ligne en échec bloque les suivantes. Pour plusieurs relais, répartissez les lignes par commande : l'ordre ne compte qu'au sein d'une commande.

Une confirmation à la fois. Messenger attend la confirmation de chaque message avant le suivant. Pour un message persistant, RabbitMQ indique qu'elle peut prendre quelques centaines de millisecondes sous charge. Le débit du relais en dépend.

Le signal qui doit faire réexaminer ce choix : un retard que les consommateurs ne tolèrent plus, ou une base chargée par le relais. Le CDC devient alors pertinent. Rejouer régulièrement de longues périodes 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, Doctrine ORM 3 (DBAL 4.3 ou plus), PostgreSQL et API Platform 4.

1. La table outbox

// src/Entity/OutboxMessage.php (service order)
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: 'outbox')]
#[ORM\Index(
    name: 'outbox_unpublished_idx',
    fields: ['occurredAt'],
    options: ['where' => 'published_at IS NULL'],
)]
class OutboxMessage
{
    #[ORM\Column(type: Types::DATETIMETZ_IMMUTABLE, nullable: true)]
    private ?\DateTimeImmutable $publishedAt = null;

    #[ORM\Column]
    private int $attempts = 0;

    public function __construct(
        // Identique à l'event_id de l'enveloppe.
        #[ORM\Id]
        #[ORM\Column(type: UuidType::NAME)]
        private Uuid $id,
        #[ORM\Column(length: 255)]
        private string $eventName,
        #[ORM\Column(length: 255)]
        private string $routingKey,
        // L'enveloppe complète, publiée telle quelle.
        #[ORM\Column(type: Types::JSONB)]
        private array $envelope,
        #[ORM\Column(type: Types::DATETIMETZ_IMMUTABLE)]
        private \DateTimeImmutable $occurredAt,
    ) {
    }

    // getRoutingKey(), getEnvelope(), getAttempts(), setAttempts(), setPublishedAt().
}

Le type Types::JSONB existe depuis DBAL 4.3, qui déprécie l'option jsonb du type json.

2. La commande et son workflow

Les places sont les constantes d'une final class OrderStatus (DRAFT, PLACED, CANCELLED), stockées dans la colonne status.

# config/packages/workflow.yaml (service order)
framework:
  workflows:
    order:
      type: state_machine
      marking_store:
        type: method
        property: status
      supports:
        - App\Entity\Order
      initial_marking: draft
      places:
        - draft
        - placed
        - cancelled
      transitions:
        place:
          from: draft
          to: placed
// src/Entity/Order.php (service order)
namespace App\Entity;

use Doctrine\ORM\Mapping as ORM;
use Symfony\Bridge\Doctrine\Types\UuidType;
use Symfony\Component\Uid\Uuid;

#[ORM\Entity]
#[ORM\Table(name: 'orders')]
class Order
{
    #[ORM\Id]
    #[ORM\Column(type: UuidType::NAME)]
    private Uuid $id;

    #[ORM\Column(length: 20)]
    private string $status = OrderStatus::DRAFT;

    // reference, customer, lines...

    public function __construct(string $reference, Customer $customer)
    {
        // UUID v7 généré en PHP : l'identifiant existe avant le flush().
        $this->id = Uuid::v7();
        // ...
    }

    // getId(), getStatus(), setStatus(), getCustomer(), getLines().
}

Order reste une entité Doctrine sans attribut API Platform. L'opération s'ajoute à la ressource DTO OrderResource, reliée à l'entité par stateOptions. read: false empêche API Platform de lire la commande avant le processor, donc hors de la transaction.

// src/ApiResource/OrderResource.php (service order), extrait
namespace App\ApiResource;

use ApiPlatform\Doctrine\Orm\State\Options;
use ApiPlatform\Metadata\ApiResource;
use ApiPlatform\Metadata\Post;
use App\Entity\Order;
use App\State\PlaceOrderProcessor;

#[ApiResource(
    shortName: 'Order',
    operations: [
        // GetCollection, Get, Post et Patch, puis :
        new Post(
            uriTemplate: '/orders/{id}/place',
            status: 200,
            input: false,
            read: false,
            processor: PlaceOrderProcessor::class,
        ),
    ],
    stateOptions: new Options(entityClass: Order::class),
)]
final readonly class OrderResource
{
    // ...
}

3. Le processor : transaction, verrou, transition

// src/State/PlaceOrderProcessor.php (service order)
namespace App\State;

use ApiPlatform\Metadata\Operation;
use ApiPlatform\State\ProcessorInterface;
use App\ApiResource\OrderResource;
use App\Entity\Order;
use App\Mapper\OrderMapper;
use Doctrine\DBAL\LockMode;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\HttpKernel\Exception\ConflictHttpException;
use Symfony\Component\HttpKernel\Exception\NotFoundHttpException;
use Symfony\Component\Workflow\WorkflowInterface;

/** @implements ProcessorInterface<null, OrderResource> */
final class PlaceOrderProcessor implements ProcessorInterface
{
    public function __construct(
        private readonly EntityManagerInterface $em,
        private readonly WorkflowInterface $orderStateMachine,
        private readonly OrderMapper $mapper,
    ) {
    }

    public function process(mixed $data, Operation $operation, array $uriVariables = [], array $context = []): OrderResource
    {
        // BEGIN, puis flush() et COMMIT, ou ROLLBACK sur exception.
        $order = $this->em->wrapInTransaction(function () use ($uriVariables): Order {
            // SELECT ... FOR UPDATE : une requête concurrente attend ici.
            $order = $this->em->find(Order::class, $uriVariables['id'], LockMode::PESSIMISTIC_WRITE)
                ?? throw new NotFoundHttpException('Commande introuvable.');

            if (!$this->orderStateMachine->can($order, 'place')) {
                throw new ConflictHttpException('La commande a déjà quitté le statut draft.');
            }

            // Le listener insère la ligne d'outbox pendant apply().
            $this->orderStateMachine->apply($order, 'place');

            return $order;
        });

        return $this->mapper->toResource($order);
    }
}

find() avec un mode de verrou exige une transaction active et relit la ligne verrouillée, même si l'entité est déjà chargée. lock(), lui, pose le verrou sans relire l'entité.

4. Le listener : insérer la ligne d'outbox

// src/EventListener/OrderPlacedOutboxListener.php (service order)
namespace App\EventListener;

use App\Entity\Order;
use App\Outbox\OutboxWriter;
use AppContracts\Order\OrderPlacedV1;
use AppContracts\Routing;
use Symfony\Component\Workflow\Attribute\AsCompletedListener;
use Symfony\Component\Workflow\Event\CompletedEvent;

final class OrderPlacedOutboxListener
{
    public function __construct(
        private readonly OutboxWriter $outbox,
    ) {
    }

    #[AsCompletedListener(workflow: 'order', transition: 'place')]
    public function onPlace(CompletedEvent $event): void
    {
        /** @var Order $order */
        $order = $event->getSubject();

        $lines = [];
        foreach ($order->getLines() as $line) {
            $lines[] = [
                'productId' => $line->getProduct()->getId()->toRfc4122(),
                'quantity' => $line->getQuantity(), // par exemple '2.000'
            ];
        }

        $this->outbox->add(
            Routing::ORDER_PLACED_V1,
            new OrderPlacedV1(
                orderId: $order->getId()->toRfc4122(),
                customerId: $order->getCustomer()->getId()->toRfc4122(),
                lines: $lines,
            ),
            correlationId: $order->getId()->toRfc4122(),
        );
    }
}
// src/Outbox/OutboxWriter.php (service order)
namespace App\Outbox;

use App\Entity\OutboxMessage;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\Serializer\Normalizer\AbstractObjectNormalizer;
use Symfony\Component\Serializer\Normalizer\NormalizerInterface;
use Symfony\Component\Uid\Uuid;

final class OutboxWriter
{
    public function __construct(
        private readonly EntityManagerInterface $em,
        private readonly NormalizerInterface $normalizer,
    ) {
    }

    public function add(string $eventName, object $event, string $correlationId): void
    {
        $eventId = Uuid::v7();
        $occurredAt = new \DateTimeImmutable();

        $envelope = [
            'event_id' => $eventId->toRfc4122(),
            'event_name' => $eventName,
            'occurred_at' => $occurredAt->format(\DateTimeInterface::RFC3339_EXTENDED),
            'producer' => 'order',
            'schema_version' => $event::SCHEMA_VERSION,
            'correlation_id' => $correlationId,
            // Absence plutôt que null : les champs vides ne sont pas écrits.
            'payload' => $this->normalizer->normalize($event, 'json', [
                AbstractObjectNormalizer::SKIP_NULL_VALUES => true,
            ]),
        ];

        // Valider ici json_encode($envelope) contre son JSON Schema. Pas de flush() : le processor s'en charge.
        $this->em->persist(new OutboxMessage(
            id: $eventId,
            eventName: $eventName,
            routingKey: $eventName,
            envelope: $envelope,
            occurredAt: $occurredAt,
        ));
    }
}

5. Le relais

// src/Outbox/OutboxRelay.php (service order)
namespace App\Outbox;

use App\Entity\OutboxMessage;
use Doctrine\ORM\EntityManagerInterface;
use Doctrine\ORM\Query\ResultSetMappingBuilder;
use Symfony\Component\DependencyInjection\Attribute\Autowire;
use Symfony\Component\Messenger\Bridge\Amqp\Transport\AmqpStamp;
use Symfony\Component\Messenger\Bridge\Amqp\Transport\AmqpTransport;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Exception\TransportException;

final class OutboxRelay
{
    public const BATCH_SIZE = 100;

    public function __construct(
        private readonly EntityManagerInterface $em,
        #[Autowire(service: 'messenger.transport.events')]
        private readonly AmqpTransport $events,
    ) {
    }

    public function relay(): int
    {
        $published = $this->em->wrapInTransaction(function (): int {
            $rsm = new ResultSetMappingBuilder($this->em);
            $rsm->addRootEntityFromClassMetadata(OutboxMessage::class, 'o');

            $messages = $this->em->createNativeQuery(
                'SELECT '.$rsm->generateSelectClause(['o' => 'o']).' FROM outbox o
                WHERE o.published_at IS NULL
                ORDER BY o.occurred_at, o.id
                LIMIT '.self::BATCH_SIZE.'
                FOR UPDATE SKIP LOCKED',
                $rsm,
            )->getResult();

            $count = 0;
            foreach ($messages as $message) {
                try {
                    // 1. Publier : send() rend la main après la confirmation du broker.
                    $this->events->send(new Envelope($message, [new AmqpStamp($message->getRoutingKey())]));
                } catch (TransportException) {
                    // Nack ou délai dépassé : on s'arrête pour garder l'ordre.
                    // Nouveau canal : un ack tardif ne confirmera pas le suivant.
                    $this->events->close();
                    $message->setAttempts($message->getAttempts() + 1);
                    break;
                }

                // 2. Marquer.
                $message->setPublishedAt(new \DateTimeImmutable());
                ++$count;
            }

            return $count;
        }); // 3. Commiter : wrapInTransaction() appelle flush() puis COMMIT.

        $this->em->clear();

        return $published;
    }
}

Le relais appelle directement le transport events, sans passer par le bus. Il remplace l'EventPublisher et les règles routing de l'article sur la topologie. Son sérialiseur publie l'enveloppe stockée telle quelle.

// src/Outbox/OutboxSerializer.php (service order)
namespace App\Outbox;

use App\Entity\OutboxMessage;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Transport\Serialization\SerializerInterface;

final class OutboxSerializer implements SerializerInterface
{
    public function encode(Envelope $envelope): array
    {
        /** @var OutboxMessage $message */
        $message = $envelope->getMessage();

        return [
            'body' => json_encode($message->getEnvelope(), \JSON_THROW_ON_ERROR),
            'headers' => ['Content-Type' => 'application/json'],
        ];
    }

    public function decode(array $encodedEnvelope): Envelope
    {
        throw new \LogicException('Le service order ne consomme pas ce transport.');
    }
}

L'option confirm_timeout active les publisher confirms. Un nack ou un délai dépassé lève une TransportException. Le transport publie déjà en mode persistant.

# config/packages/messenger.yaml (service order)
framework:
  messenger:
    transports:
      events:
        dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
        serializer: App\Outbox\OutboxSerializer
        options:
          exchange:
            name: !php/const AppContracts\Routing::EXCHANGE
            type: topic
          queues: []
          # Attendre la confirmation du broker, 5 secondes au plus.
          confirm_timeout: 5

6. Lancer le relais

Symfony Scheduler déclenche le relais chaque seconde. Tant que les lots sont pleins, la tâche enchaîne.

// src/Scheduler/RelayOutboxTask.php (service order)
namespace App\Scheduler;

use App\Outbox\OutboxRelay;
use Symfony\Component\Scheduler\Attribute\AsPeriodicTask;

#[AsPeriodicTask(frequency: 1)]
final class RelayOutboxTask
{
    public function __construct(
        private readonly OutboxRelay $relay,
    ) {
    }

    public function __invoke(): void
    {
        while (OutboxRelay::BATCH_SIZE === $this->relay->relay()) {
        }
    }
}

Nous lançons php bin/console messenger:consume scheduler_default en un seul exemplaire. Si deux relais se chevauchent, SKIP LOCKED leur donne des lignes différentes, mais l'ordre n'est plus garanti.

Checklist

  • L'outbox écrite dans la même transaction que le changement métier.
  • Une écriture sur le fait métier, pas dans un hook Doctrine.
  • Un verrou sur les transitions concurrentes.
  • Des UUID v7 générés en PHP.
  • Une enveloppe validée contre son JSON Schema.
  • Des publisher confirms, et l'ordre : publier, marquer, commiter.
  • Un seul relais actif, arrêté au premier échec.
  • Des consommateurs idempotents.
  • Une alerte sur l'âge de la plus ancienne ligne non publiée.
  • Une rétention définie et une purge planifiée.

Sources