Le problème

Reprenons notre boutique en ligne. Le service order écrit chaque événement dans sa table outbox, dans la transaction de la commande. Un relais publie ensuite ces lignes sur l'exchange app.events, et les services stock et shipping les lisent dans leurs queues.

Un jour, il faut traiter de nouveau des événements déjà passés. Trois situations reviennent souvent :

  • un nouveau consommateur : le service billing arrive et doit connaître les commandes passées depuis le début du mois. Sa queue ne reçoit que les messages publiés après sa création ;
  • des messages perdus : une queue supprimée par erreur, ou des messages retirés du failure transport sans avoir été traités ;
  • un traitement à refaire : un handler de stock contenait un bug, corrigé depuis, et les commandes d'une période doivent repasser.

RabbitMQ n'aide pas ici : une queue efface un message dès qu'il est acquitté. La question arrive alors vite : faut-il passer à Kafka, dont le log garde les événements après leur lecture ?

Les options

Nous avons vérifié ces faits le 2 octobre 2026, sur Kafka 4.3.1, sa dernière version.

Passer à Kafka

Ce qu'il règle. Un topic garde ses événements après leur lecture, pendant une rétention fixée par topic. Chaque groupe de consommateurs connaît sa position (l'offset) et peut la ramener à une date pour tout relire.

L'ordre est garanti dans une partition, et la clé du message choisit la partition. Avec l'identifiant de commande comme clé, les événements d'une même commande restent dans l'ordre. La compaction garde au moins la dernière valeur de chaque clé : un topic compacté sert d'état courant.

Ce qu'il coûte.

  • Pas de TTL ni de DLQ par message dans le broker. Un groupe de consommateurs garde un seul offset par partition. Un message en échec se retente sur place, en bloquant sa partition, ou le code client le republie ailleurs.
  • Des DLQ hors du broker seulement, dans Kafka Connect et, depuis la version 4.2, dans Kafka Streams.
  • Un parallélisme plafonné. Dans un groupe, chaque partition est lue par un seul consommateur à la fois. Le nombre de partitions plafonne donc le nombre de consommateurs utiles.
  • Des share groups encore jeunes. Prêts pour la production depuis Kafka 4.2, ils acquittent message par message et comptent les tentatives, mais leur DLQ n'est qu'une proposition acceptée. Côté PHP, librdkafka ne les propose qu'en préversion, et l'extension rdkafka ne les expose pas encore.
  • Une exploitation plus lourde. Depuis Kafka 4.0, le mode KRaft remplace ZooKeeper. Il reste un cluster à dimensionner, à surveiller et à mettre à jour.
  • Pas de transport officiel pour Messenger. La documentation de Symfony renvoie vers Enqueue, et les transports communautaires s'appuient en général sur l'extension rdkafka. koco/messenger-kafka, en version 0.x, recommande de désactiver enable.auto.offset.store : sinon, chaque message est acquitté même si son handler échoue.

Utiliser les streams de RabbitMQ

  • Ce qu'ils règlent : un stream est un journal dans RabbitMQ. La lecture n'efface rien, et un consommateur s'attache à un offset ou à une date. Les super streams le partitionnent pour monter en charge.
  • Ce qu'ils coûtent : pas de dead letter exchange ni de TTL par message. Surtout, RabbitMQ refuse basic.get sur un stream, or le transport AMQP de Messenger lit par basic.get. Les consommateurs auraient besoin d'un autre client.

Garder RabbitMQ et faire de l'outbox un journal

  • Ce qu'elle règle : l'outbox contient déjà chaque événement publié, avec son event_id, sa clé de routage et son enveloppe JSON complète. En gardant longtemps les lignes publiées, nous obtenons un journal. Rejouer, c'est republier une sélection de lignes vers la queue d'un seul consommateur.
  • Ce qu'elle coûte : une table qui grossit, et une commande console à écrire et à maintenir. Le journal reste propre à chaque producteur.

Notre recommandation

Tant que le rejeu reste ponctuel, nous gardons RabbitMQ : l'outbox, conservée longtemps, devient un journal que nous republions vers la queue d'un seul consommateur.

Les briques existent déjà. L'outbox garde chaque événement avec son identifiant. Les consommateurs idempotents écartent, grâce à leur table processed_event, ce qu'ils ont déjà traité : republier un événement de trop n'a pas d'effet.

Nous ne changeons ni de broker ni de client PHP. Le transport AMQP de Messenger, ses retries et son failure transport restent en place. Pour les trois situations du début, nous obtenons l'essentiel de ce que le log de Kafka apporterait.

Rejouer depuis l'outbox : seule la queue visée reçoit les événements republiés RabbitMQ et outbox à rétention longue outbox service order order.order.cancelled.v1 le rejeu filtre par clé et par période relais app.events topic exchange shipping.events stock.events order.order.placed.v1 order.order.placed.v1 rejeu exchange par défaut clé = nom de la queue seule stock.events reçoit le rejeu Kafka : le log reste après lecture, chaque groupe garde sa position un topic, partition 0 0 1 2 3 4 5 6 7 8 groupe shipping : offset 8 groupe stock : relit depuis l'offset 3 (position ramenée à une date)

Ce que le rejeu retraite, et ce qu'il écarte

Le consommateur garde le dernier mot : un événement absent de sa table processed_event est traité, un événement présent est ignoré. Les trois situations diffèrent donc :

  • pour un nouveau consommateur, chaque événement rejoué est traité ;
  • pour des messages perdus, seuls ceux qui n'avaient pas abouti sont traités ;
  • pour un bug qui a produit un mauvais effet sans erreur, l'événement est déjà enregistré. Il faut corriger les données, puis retirer les lignes concernées de processed_event, par un geste explicite chez le consommateur.

Ce dernier cas n'est pas propre à l'outbox : avec Kafka, un consommateur idempotent ignorerait aussi un événement relu.

Les signaux qui justifient Kafka

Nous réexaminons ce choix quand l'un de ces signaux apparaît :

  • des volumes que le relais par polling ne suit plus, constatés sur le retard de publication ;
  • de nombreux consommateurs qui ont besoin de rejouer, au point que le rejeu devient une routine ;
  • du traitement de flux : fenêtres de temps, jointures entre flux, agrégats calculés en continu ;
  • une rétention longue de l'historique imposée par le métier ou par la réglementation ;
  • une équipe capable d'exploiter un cluster Kafka et ses clients.

Ce qui change dans la conception avec Kafka

Le nommage porte sur des topics, et la queue de chaque service devient un groupe de consommateurs, avec un offset par partition.

L'ordre impose de regrouper : les événements d'une commande partagent un topic et une clé, l'identifiant de commande. Un topic par type d'événement perdrait l'ordre entre placed et cancelled.

Les compromis assumés

Un journal par producteur. Chaque service producteur garde son outbox et lance ses propres rejeux. Rejouer des événements de order et de stock demande deux rejeux, sans ordre garanti entre eux.

Une table qui grossit. La rétention se décide comme une durée de conservation : assez longue pour les rejeux prévus, pas plus. Sur une table volumineuse, un partitionnement par mois remplace le DELETE par la suppression d'une partition.

Le rejeu contourne les bindings. L'exchange par défaut livre à la queue nommée, quelles que soient ses binding keys. Ne rejouez vers une queue que des clés auxquelles elle est liée : sinon, le consommateur reçoit un événement qu'il ne sait pas traiter.

L'ordre n'est pas garanti. Un rejeu traite une seule clé de routage : l'ordre entre clés se perd, et ses événements arrivent après le flux normal. Une copie d'état doit vérifier une version avant d'écrire : la table processed_event ne dit pas si un événement arrive trop tard.

Deux horizons à aligner. Ne rejouez pas vers un consommateur existant des événements plus anciens que la purge de sa table processed_event. Il ne les reconnaîtrait plus et les traiterait une seconde fois.

Le signal qui doit faire réexaminer ce choix : le rejeu devient une opération courante plutôt qu'un geste d'exception. Ou le retard du relais, mesuré, dépasse ce que les consommateurs tolèrent.

Mise en œuvre

Les exemples utilisent Symfony 7.4 (options identiques en 8.1), le transport AMQP de Messenger, Doctrine ORM 3 et PostgreSQL. L'exchange et les queues sont ceux de l'article RabbitMQ : une seule topic exchange et une queue par service consommateur. La table outbox et le sérialiseur du relais sont ceux de notre article sur le pattern Outbox.

1. Garder les lignes publiées

Le relais marque chaque ligne publiée (published_at) et la laisse en place. Un index sur la clé de routage et la date sert les sélections du rejeu.

// src/Entity/OutboxMessage.php (service order, extrait)
namespace App\Entity;

use Doctrine\ORM\Mapping as ORM;

#[ORM\Entity]
#[ORM\Table(name: 'outbox')]
#[ORM\Index(
    name: 'outbox_unpublished_idx',
    fields: ['occurredAt'],
    options: ['where' => 'published_at IS NULL'],
)]
// Ajouté pour le rejeu : sélection par clé de routage et par période.
#[ORM\Index(name: 'outbox_replay_idx', fields: ['routingKey', 'occurredAt'])]
class OutboxMessage
{
    // id (= event_id), routing_key, envelope (l'enveloppe JSON complète, en JSONB),
    // occurred_at, published_at... : voir l'article sur le pattern Outbox.
}

La purge devient une tâche planifiée, réglée sur l'horizon de rejeu choisi :

-- tâche planifiée du service order
DELETE FROM outbox
WHERE published_at IS NOT NULL
  AND occurred_at < :horizon;

2. Un transport qui vise une seule queue

L'exchange par défaut de RabbitMQ est une exchange directe au nom vide. Chaque queue y est liée automatiquement, avec son nom comme clé de routage. Publier sur cette exchange avec la clé stock.events remet donc le message à la seule queue stock.events.

Messenger s'en sert déjà pour ses retries : la queue de délai renvoie le message à la seule queue en échec. Pour le rejeu, nous déclarons un transport dont l'exchange porte un nom vide, ce que le transport AMQP prévoit.

# config/packages/messenger.yaml (service order, extrait)
framework:
  messenger:
    transports:
      replay:
        dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
        # Le sérialiseur du relais : il publie l'enveloppe JSON stockée telle quelle.
        serializer: App\Outbox\OutboxSerializer
        options:
          # Nom vide : l'exchange par défaut, la clé de routage est un nom de queue.
          exchange:
            name: ''
          # Le rejeu ne déclare ni ne lie aucune queue.
          queues: []
          # Attend la confirmation du broker pour chaque message (en secondes).
          confirm_timeout: 5

Ce transport n'apparaît pas dans routing : seule la commande de rejeu l'utilise, directement.

3. Publier l'enveloppe JSON telle qu'elle est stockée

La colonne envelope contient l'enveloppe JSON complète, event_id compris. Le transport replay reprend le sérialiseur du relais, App\Outbox\OutboxSerializer, qui publie cette enveloppe sans la reconstruire. Le consommateur reçoit donc le même événement qu'à la première publication.

Le consommateur doit lire le type de l'événement dans l'enveloppe JSON (event_name), pas dans la clé de routage AMQP. Ici, la clé de routage d'un message rejoué est le nom de la queue.

4. Écrire la commande console de rejeu

La commande console sélectionne les lignes publiées par clé de routage et par période, dans l'ordre d'origine. Elle les envoie une à une sur le transport replay, comme le relais, avec la queue cible comme clé de routage.

// src/Command/ReplayOutboxCommand.php (service order)
namespace App\Command;

use Doctrine\DBAL\Types\Types;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\Console\Attribute\Argument;
use Symfony\Component\Console\Attribute\AsCommand;
use Symfony\Component\Console\Attribute\Option;
use Symfony\Component\Console\Command\Command;
use Symfony\Component\Console\Output\OutputInterface;
use Symfony\Component\DependencyInjection\Attribute\Autowire;
use Symfony\Component\Messenger\Bridge\Amqp\Transport\AmqpStamp;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Transport\Sender\SenderInterface;

#[AsCommand(name: 'app:outbox:replay', description: 'Republie des événements vers la queue d\'un consommateur')]
final class ReplayOutboxCommand
{
    public function __construct(
        private readonly EntityManagerInterface $entityManager,
        #[Autowire(service: 'messenger.transport.replay')]
        private readonly SenderInterface $replay,
    ) {
    }

    public function __invoke(
        OutputInterface $output,
        #[Argument(description: 'Queue cible')] string $queue,
        #[Argument(description: 'Clé de routage à rejouer')] string $routingKey,
        #[Argument(description: 'Début, inclus (RFC 3339)')] string $from,
        #[Argument(description: 'Fin, exclue (RFC 3339)')] string $to,
        #[Option(description: 'Compter sans publier')] bool $dryRun = false,
    ): int {
        $query = $this->entityManager->createQuery(
            'SELECT m FROM App\Entity\OutboxMessage m
             WHERE m.routingKey = :routingKey
               AND m.occurredAt >= :from AND m.occurredAt < :to
               AND m.publishedAt IS NOT NULL
             ORDER BY m.occurredAt ASC, m.id ASC'
        )
            ->setParameter('routingKey', $routingKey)
            // Sans type explicite, Doctrine écrirait la date sans son décalage horaire.
            ->setParameter('from', new \DateTimeImmutable($from), Types::DATETIMETZ_IMMUTABLE)
            ->setParameter('to', new \DateTimeImmutable($to), Types::DATETIMETZ_IMMUTABLE);

        $count = 0;
        foreach ($query->toIterable() as $message) {
            if (!$dryRun) {
                // Exchange par défaut : la clé de routage est le nom de la queue cible.
                $this->replay->send(new Envelope($message, [new AmqpStamp($queue)]));
            }

            if (0 === ++$count % 100) {
                $this->entityManager->clear();
            }
        }

        $output->writeln(sprintf('%d événement(s) pour %s%s.', $count, $queue, $dryRun ? ' (essai)' : ''));

        return Command::SUCCESS;
    }
}

La commande console ne modifie pas l'outbox. Relancer un rejeu interrompu republie des événements déjà livrés : les consommateurs les écartent grâce à processed_event.

5. Lancer un rejeu

Vérifiez d'abord que la queue cible existe et qu'elle est liée à la clé rejouée. Un message envoyé vers une queue absente est écarté par RabbitMQ, et la confirmation du broker arrive quand même.

# Bindings de la queue cible
rabbitmqctl list_bindings source_name destination_name routing_key | grep stock.events

# Compter, puis republier
php bin/console app:outbox:replay stock.events order.order.placed.v1 \
  2026-09-14T08:00:00Z 2026-09-14T12:00:00Z --dry-run
php bin/console app:outbox:replay stock.events order.order.placed.v1 \
  2026-09-14T08:00:00Z 2026-09-14T12:00:00Z

Pour un nouveau consommateur, l'ordre des étapes compte. Déployez billing et lancez messenger:setup-transports : sa queue existe et reçoit déjà le flux normal. Rejouez ensuite l'historique voulu, le recouvrement avec le flux normal étant absorbé par processed_event.

Checklist

  • Les lignes publiées de l'outbox gardées, avec une rétention écrite et une purge planifiée.
  • Un index sur (routing_key, occurred_at) pour les sélections du rejeu.
  • Un transport de rejeu sur l'exchange par défaut, avec confirmation du broker.
  • L'enveloppe JSON republiée telle quelle, avec son event_id d'origine.
  • Des consommateurs idempotents, qui lisent le type de l'événement dans l'enveloppe JSON.
  • Une queue cible vérifiée, liée aux clés rejouées.
  • Une purge de processed_event alignée sur l'horizon de rejeu de l'outbox.
  • Un essai avec --dry-run avant chaque rejeu.
  • Les signaux qui justifient Kafka, suivis.

Sources