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 billing service 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 stock handler 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 rdkafka extension 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 rdkafka extension. koco/messenger-kafka, at version 0.x, recommends disabling enable.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.get on a stream, and Messenger's AMQP transport reads with basic.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.

Replaying from the outbox: only the targeted queue receives the republished events RabbitMQ and a long-retention outbox outbox order service order.order.cancelled.v1 the replay filters by key and period relay app.events topic exchange shipping.events stock.events order.order.placed.v1 order.order.placed.v1 replay default exchange key = queue name only stock.events gets the replay Kafka: the log stays after reading, each group keeps its position a topic, partition 0 0 1 2 3 4 5 6 7 8 shipping group: offset 8 stock group: reads again from offset 3 (position moved back to a date)

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_event purge aligned with the outbox replay horizon.
  • A --dry-run trial before each replay.
  • The signals that justify Kafka, monitored.

Sources