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
billingarrive 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
stockcontenait 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
rdkafkane 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ésactiverenable.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.getsur un stream, or le transport AMQP de Messenger lit parbasic.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.
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_idd'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_eventalignée sur l'horizon de rejeu de l'outbox. -
Un essai avec
--dry-runavant chaque rejeu. - Les signaux qui justifient Kafka, suivis.
Sources
- Apache Kafka, téléchargements : version 4.3.1 du 25 juin 2026.
- Apache Kafka, annonces des versions 4.0 (fin de ZooKeeper, mode KRaft) et 4.2 (share groups, DLQ de Kafka Streams).
- Apache Kafka, introduction : rétention par topic, ordre par partition.
- Apache Kafka, conception : groupes de consommateurs, offsets, compaction, share groups.
-
Apache Kafka,
Kafka Connect
(
errors.deadletterqueue.topic.name) et opérations (kafka-consumer-groups.sh --reset-offsets). - Apache Kafka, KIP-1191 : DLQ pour les share groups, proposition acceptée.
- librdkafka, journal des versions : consommateur des share groups en préversion depuis la v2.15.0.
- php-rdkafka, version 6.0.5 : pas d'API pour les share groups.
-
RabbitMQ,
Streams and Super Streams
et
code des streams
:
basic.getrefusé. - RabbitMQ, Default Exchange, messages non routables et leur confirmation.
-
Symfony,
Messenger
: transports, renvoi vers Enqueue,
exchange[name],confirm_timeout. -
Symfony,
code du transport AMQP
: lecture par
get(), exchange par défaut pour les retries. -
Symfony,
arguments et options des commandes
:
#[Argument]et#[Option]. -
Doctrine ORM,
Batch Processing
:
toIterable()etclear(). - PostgreSQL, partitionnement de tables.
- koco/messenger-kafka : transport Kafka communautaire.