The problem
Back to our online shop. The order service
writes each event to its outbox table, in the
transaction of the order. A relay then publishes these rows
to the app.events exchange, and the
stock and shipping services read
them from their queues.
One day, events that are already past must be processed again. Three situations come up often:
-
a new consumer: the
billingservice arrives and must know the orders placed since the start of the month. Its queue only receives the messages published after it was created; - lost messages: a queue deleted by mistake, or messages removed from the failure transport without being processed;
-
processing to redo: a
stockhandler had a bug, fixed since, and the orders of a given period must go through again.
RabbitMQ does not help here: a queue deletes a message as soon as it is acknowledged. The question soon follows: should we move to Kafka, whose log keeps the events after they are read?
The options
We checked these facts on 2 October 2026, against Kafka 4.3.1, its latest version.
Moving to Kafka
What it solves. A topic keeps its events after they are read, for a retention period set per topic. Each consumer group knows its position (the offset) and can move it back to a date to read everything again.
Order is guaranteed within a partition, and the message key chooses the partition. With the order ID as the key, the events of a given order stay in order. Compaction keeps at least the last value of each key: a compacted topic serves as current state.
What it costs.
- No per-message TTL or DLQ in the broker. A consumer group keeps a single offset per partition. A failed message is retried in place, blocking its partition, or the client code republishes it elsewhere.
- DLQs outside the broker only, in Kafka Connect and, since version 4.2, in Kafka Streams.
- Capped parallelism. Within a group, each partition is read by a single consumer at a time. The number of partitions therefore caps the number of useful consumers.
-
Share groups that are still young.
Production-ready since Kafka 4.2, they acknowledge record
by record and count delivery attempts, but their DLQ is
only an accepted proposal. On the PHP side, librdkafka
offers them as a preview only, and the
rdkafkaextension does not expose them yet. - Heavier operations. Since Kafka 4.0, KRaft mode replaces ZooKeeper. A cluster still has to be sized, monitored and upgraded.
-
No official transport for Messenger. The
Symfony documentation points to Enqueue, and community
transports generally rely on the
rdkafkaextension.koco/messenger-kafka, at version 0.x, recommends disablingenable.auto.offset.store: otherwise, every message is acknowledged even when its handler fails.
Using RabbitMQ streams
- What they solve: a stream is a log inside RabbitMQ. Reading deletes nothing, and a consumer attaches at an offset or at a date. Super streams partition it to scale out.
-
What they cost: no dead letter exchange
and no per-message TTL. Above all, RabbitMQ refuses
basic.geton a stream, and Messenger's AMQP transport reads withbasic.get. The consumers would need another client.
Keeping RabbitMQ and turning the outbox into a log
-
What it solves: the outbox already holds
every published event, with its
event_id, its routing key and its full JSON envelope. Keeping the published rows for a long time gives us a log. Replaying means republishing a selection of rows to the queue of a single consumer. - What it costs: a table that grows, and a console command to write and maintain. The log stays specific to each producer.
Our recommendation
While replays stay occasional, we keep RabbitMQ: a long-lived outbox becomes a log that we republish to the queue of a single consumer.
The building blocks already exist. The outbox keeps every
event with its ID. Idempotent consumers, thanks to their
processed_event table, skip what they have
already processed: republishing one event too many has no
effect.
We change neither the broker nor the PHP client. Messenger's AMQP transport, its retries and its failure transport stay as they are. For the three situations above, we get the core of what Kafka's log would bring.
What a replay processes again, and what it skips
The consumer has the last word: an event missing from its
processed_event table is processed, an event
already there is ignored. The three situations therefore
differ:
- for a new consumer, every replayed event is processed;
- for lost messages, only those that had not succeeded are processed;
-
for a bug that produced a wrong effect without an error,
the event is already recorded. You must fix the data, then
remove the matching rows from
processed_event, as an explicit step on the consumer's side.
This last case is not specific to the outbox: with Kafka, an idempotent consumer would also ignore an event read again.
The signals that justify Kafka
We revisit this choice when one of these signals appears:
- volumes that the polling relay no longer keeps up with, seen in the publication delay;
- many consumers that need to replay, to the point where replaying becomes routine;
- stream processing: time windows, joins between streams, aggregates computed continuously;
- long retention of the history, required by the business or by regulation;
- a team able to operate a Kafka cluster and its clients.
What changes in the design with Kafka
Naming applies to topics, and each service's queue becomes a consumer group, with one offset per partition.
Order requires grouping: the events of an order share a
topic and a key, the order ID. One topic per event type
would lose the order between placed and
cancelled.
The trade-offs we accept
One log per producer. Each producing
service keeps its outbox and runs its own replays. Replaying
events from order and from
stock takes two replays, with no guaranteed
order between them.
A table that grows. Retention is decided
like a data retention period: long enough for the planned
replays, no longer. On a large table, monthly partitioning
replaces the DELETE with dropping a partition.
Replay bypasses the bindings. The default exchange delivers to the named queue, whatever its binding keys. Only replay keys that the target queue is bound to: otherwise, the consumer receives an event it cannot handle.
Order is not guaranteed. A replay covers a
single routing key: the order between keys is lost, and its
events arrive after the normal flow. A copy of state must
check a version before writing: the
processed_event table does not tell whether an
event arrives too late.
Two horizons to align. Do not replay to an
existing consumer events older than the purge of its
processed_event table. It would no longer
recognise them and would process them a second time.
The signal that should make us revisit this choice: replaying becomes a routine operation rather than an exceptional one. Or the relay's measured delay exceeds what the consumers tolerate.
Implementation
The examples use Symfony 7.4 (same options in 8.1),
Messenger's AMQP transport, Doctrine ORM 3 and PostgreSQL.
The exchange and the queues are those of the article
RabbitMQ: one topic exchange and one queue per consuming
service. The outbox table and the relay's serializer
are those of
our article on the Outbox pattern.
1. Keep the published rows
The relay marks each published row
(published_at) and leaves it in place. An index
on the routing key and the date serves the replay's
selections.
// src/Entity/OutboxMessage.php (order service, excerpt)
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'],
)]
// Added for the replay: selection by routing key and period.
#[ORM\Index(name: 'outbox_replay_idx', fields: ['routingKey', 'occurredAt'])]
class OutboxMessage
{
// id (= event_id), routing_key, envelope (the full JSON envelope, as JSONB),
// occurred_at, published_at...: see the article on the Outbox pattern.
}
The purge becomes a scheduled task, set to the chosen replay horizon:
-- scheduled task of the order service
DELETE FROM outbox
WHERE published_at IS NOT NULL
AND occurred_at < :horizon;
2. A transport that targets a single queue
RabbitMQ's default exchange is a direct exchange with an
empty name. Every queue is bound to it automatically, with
its own name as the routing key. Publishing to this exchange
with the stock.events key therefore delivers
the message to the stock.events queue only.
Messenger already relies on it for its retries: the delay queue sends the message back to the failing queue only. For the replay, we declare a transport whose exchange has an empty name, which the AMQP transport supports.
# config/packages/messenger.yaml (order service, excerpt)
framework:
messenger:
transports:
replay:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
# The relay's serializer: it publishes the stored JSON envelope as is.
serializer: App\Outbox\OutboxSerializer
options:
# Empty name: the default exchange, the routing key is a queue name.
exchange:
name: ''
# The replay declares and binds no queue.
queues: []
# Waits for the broker's confirmation of each message (in seconds).
confirm_timeout: 5
This transport does not appear under routing:
only the replay command uses it, directly.
3. Publish the JSON envelope as it is stored
The envelope column holds the full JSON
envelope, event_id included. The
replay transport reuses the relay's serializer,
App\Outbox\OutboxSerializer, which publishes
this envelope without rebuilding it. The consumer therefore
receives the same event as on the first publication.
The consumer must read the event type from the JSON envelope
(event_name), not from the AMQP routing key.
Here, the routing key of a replayed message is the queue
name.
4. Write the replay console command
The console command selects the published rows by routing
key and period, in their original order. It sends them one
by one on the replay transport, like the relay,
with the target queue as the routing key.
// src/Command/ReplayOutboxCommand.php (order service)
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: 'Republishes events to the queue of one consumer')]
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: 'Target queue')] string $queue,
#[Argument(description: 'Routing key to replay')] string $routingKey,
#[Argument(description: 'Start, inclusive (RFC 3339)')] string $from,
#[Argument(description: 'End, exclusive (RFC 3339)')] string $to,
#[Option(description: 'Count without publishing')] 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)
// Without an explicit type, Doctrine would write the date without its UTC offset.
->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) {
// Default exchange: the routing key is the name of the target queue.
$this->replay->send(new Envelope($message, [new AmqpStamp($queue)]));
}
if (0 === ++$count % 100) {
$this->entityManager->clear();
}
}
$output->writeln(sprintf('%d event(s) for %s%s.', $count, $queue, $dryRun ? ' (dry run)' : ''));
return Command::SUCCESS;
}
}
The console command does not modify the outbox. Running an
interrupted replay again republishes events already
delivered: the consumers skip them thanks to
processed_event.
5. Run a replay
First check that the target queue exists and that it is bound to the replayed key. A message sent to a missing queue is discarded by RabbitMQ, and the broker's confirmation still arrives.
# Bindings of the target queue
rabbitmqctl list_bindings source_name destination_name routing_key | grep stock.events
# Count, then republish
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
For a new consumer, the order of the steps matters. Deploy
billing and run
messenger:setup-transports: its queue exists
and already receives the normal flow. Then replay the
history you need, the overlap with the normal flow being
absorbed by processed_event.
Checklist
- The published outbox rows kept, with a written retention and a scheduled purge.
-
An index on
(routing_key, occurred_at)for the replay's selections. - A replay transport on the default exchange, with broker confirmation.
-
The JSON envelope republished as is, with its original
event_id. - Idempotent consumers, which read the event type from the JSON envelope.
- A target queue checked, bound to the replayed keys.
-
A
processed_eventpurge aligned with the outbox replay horizon. -
A
--dry-runtrial before each replay. - The signals that justify Kafka, monitored.
Sources
- Apache Kafka, downloads: version 4.3.1 of 25 June 2026.
- Apache Kafka, release announcements of 4.0 (end of ZooKeeper, KRaft mode) and 4.2 (share groups, Kafka Streams DLQ).
- Apache Kafka, introduction: per-topic retention, order per partition.
- Apache Kafka, design: consumer groups, offsets, compaction, share groups.
-
Apache Kafka,
Kafka Connect
(
errors.deadletterqueue.topic.name) and operations (kafka-consumer-groups.sh --reset-offsets). - Apache Kafka, KIP-1191: DLQ for share groups, accepted proposal.
- librdkafka, changelog: share consumer in preview since v2.15.0.
- php-rdkafka, version 6.0.5: no share group API.
-
RabbitMQ,
Streams and Super Streams
and
stream code:
basic.getrefused. - RabbitMQ, Default Exchange, unroutable messages and their confirmation.
-
Symfony,
Messenger: transports, pointer to Enqueue,
exchange[name],confirm_timeout. -
Symfony,
AMQP transport code: reading with
get(), default exchange for retries. -
Symfony,
command arguments and options:
#[Argument]and#[Option]. -
Doctrine ORM,
Batch Processing:
toIterable()andclear(). - PostgreSQL, table partitioning.
- koco/messenger-kafka: community Kafka transport.