The problem

Back to our online shop. The order service records an order and publishes order.order.placed.v1 through an outbox relay. The stock service reads its stock.events queue and reserves the goods of each line.

With acknowledgements, RabbitMQ guarantees "at least once" delivery. It does not guarantee "exactly once". The same event can therefore reach stock several times, for ordinary reasons:

  • the relay publishes again: it published the outbox row, then stopped before marking it as published. On restart, it publishes it again;
  • the broker redelivers: the worker committed its transaction, then lost its connection before acknowledging. RabbitMQ then puts the message back in the queue;
  • someone replays: a period of events is republished after an incident. Or a message is retried by hand while another copy has already been processed.

Without precautions, every delivery produces its effect. The stock is reserved twice, the customer receives two shipping notices. Retries cannot help: they run the processing again, they do not know it already succeeded.

RabbitMQ's redelivered flag is not enough to spot duplicates. An event republished by the relay arrives as a new message, without that flag.

On the Symfony side, a message redelivered by RabbitMQ is not handled directly. A default Messenger middleware sends it into the retry loop: the duplicate therefore comes back later.

The options

Natural idempotence

  • What it solves: an effect that sets a target state ("the order is shipped") can be replayed without harm. An upsert on a key is the typical example.
  • What it costs: it does not cover cumulative effects. Reserving two units, decrementing a counter or sending an email are not target states.

Deduplication by business key

  • What it solves: a uniqueness constraint applies to the effect itself, for example a single reservation per order and per product.
  • What it costs: each effect needs a key, and one does not always exist. The deduplication rule is then rewritten in every handler.

A table of processed events

  • What it solves: each consuming service records the identifier of the events it has processed, in its own database. A duplicate is recognised before any effect, whatever the handler.
  • What it costs: one more write per message, one table per consuming service and a purge to organise.

Deduplication by the broker or by Messenger

  • What it solves: discarding duplicates before they reach the consumer.
  • What it costs: RabbitMQ offers none for classic or quorum queues. Streams deduplicate on publication, but our topology does not use them. A community plugin, outside the official distribution, only remembers identifiers in a bounded cache (exchange) or while the original waits in the queue.

Symfony 7.3 added a DeduplicateStamp to Messenger. It takes a lock on dispatch and releases it when processing ends: a copy that arrives afterwards goes through. In every case, the only reliable record of the processing is the consumer's database, committed together with the effect.

Our recommendation

We recommend a processed_event table in the consuming service's database, filled before any business effect with INSERT ... ON CONFLICT DO NOTHING, in the same transaction as that effect.

Two deliveries of the same event: a single stock reservation 1st delivery Message received from stock.events Transaction ( doctrine_transaction ) INSERT processed_event 1 row Handler: reservation created COMMIT: everything is saved Crash before the acknowledgement Unacknowledged: back as a retry 2nd delivery: same event_id Message received from stock.events Transaction ( doctrine_transaction ) INSERT processed_event 0 rows Handler not called COMMIT: nothing new Acknowledgement: message removed A single reservation in total

The identifier comes from the producer

The event_id is a UUID v7 set by the producer when it writes its outbox row. It travels in the message's JSON envelope, next to event_name, occurred_at and payload. The consumer never generates it: an identifier created on receipt would be new on every delivery.

The same transaction as the effect

Messenger's doctrine_transaction middleware opens a transaction before the handlers. If they succeed, it calls flush() then commit(). On an exception, it rolls everything back: the processed_event row disappears with the reservations, and the retry processes the event again.

If the commit goes through but the acknowledgement is lost, the next delivery finds the row and stops.

Avoid the error rather than catch it

A plain INSERT would raise a unique violation on a duplicate. Catching it is not enough: after an error, PostgreSQL refuses any other statement of the transaction (code 25P02) until the rollback. ON CONFLICT DO NOTHING raises no error, and Doctrine DBAL's executeStatement() returns the number of inserted rows: zero for a duplicate.

If two workers receive the same event at the same moment, the primary key settles it. The second INSERT waits for the first one's transaction to end. If it commits, the second inserts nothing; if it rolls back, the second inserts and processes.

In a middleware, not in each handler

Messenger only passes the message to the handler: the event_id stays on the Messenger envelope, in a stamp. A middleware sees that envelope and runs once per message, before all of its handlers. Inserted in each handler, the row would block a second handler of the same event: it would find the row and do nothing.

The trade-offs we accept

One more write per message. Each processed event adds a row and an index entry to the consumer's database. On a high-traffic queue, this cost needs watching.

One table per consuming service, and a purge. stock, shipping and billing each have their own processed_event table. Without a purge, it grows forever; with a purge that is too short, duplicates come back (step 8).

The key is the event_id at the scale of the service. If a service binds the same event to two of its transports, each with its own handlers (fromTransport), the second copy is discarded.

Only effects in the database are protected. The transaction covers what the handler writes to PostgreSQL. A sent email or an HTTP call is not undone by a rollback: if the commit then fails, the retry will do it again.

These effects need their own idempotency key, or a separate message sent after the commit.

Ordering is not handled. The table tells whether an event was already processed, not whether it arrives too late. For a copy of state, the version check of step 6 covers this case.

The signal that should trigger a review of this choice: the insert into processed_event shows up in the measurements as the bottleneck. Or most effects of a service live outside its database, and the table no longer protects them.

Implementation

The examples use Symfony 7.4, Doctrine ORM 3 with DBAL 4, PostgreSQL and Messenger's AMQP transport. The classes mentioned also exist in Symfony 8. The events transport and the stock.events queue are those of the article RabbitMQ: one topic exchange and one queue per consuming service.

1. Keep the event_id in a stamp

The transport's serializer reads the JSON envelope. It builds the DTO from payload and stores event_id in a stamp.

// src/Messenger/EventIdStamp.php (stock service)
namespace App\Messenger;

use Symfony\Component\Messenger\Stamp\StampInterface;

final class EventIdStamp implements StampInterface
{
    public function __construct(
        public readonly string $eventId,
    ) {
    }
}
// src/Messenger/EventSerializer.php (stock service, excerpt)
use Symfony\Component\Messenger\Exception\MessageDecodingFailedException;

// ...

public function decode(array $encodedEnvelope): Envelope
{
    try {
        $body = json_decode($encodedEnvelope['body'], true, flags: \JSON_THROW_ON_ERROR);
        // toEvent() maps event_name to the DTO class (OrderPlacedV1...).
        $event = $this->toEvent($body['event_name'], $body['payload']);
    } catch (\Throwable $e) {
        // The AMQP receiver then rejects the message, without putting it back in the queue.
        throw new MessageDecodingFailedException($e->getMessage(), 0, $e);
    }

    // decodeStamps() rebuilds the RedeliveryStamp from the header written by encode().
    $stamps = $this->decodeStamps($encodedEnvelope['headers'] ?? []);

    // The other fields of the JSON envelope go into a second stamp, left out here.
    return new Envelope($event, [new EventIdStamp($body['event_id']), ...$stamps]);
}

On a retry, Messenger encodes the message again with this serializer. encode() must therefore write event_id back from the stamp, so that the copy keeps its identifier.

encode() and decode() must also carry the retry count in a header, and decode() must turn it back into a RedeliveryStamp. Without it, the retry counter starts again from zero and the message never reaches the failure transport. The failure transport then keeps the whole Messenger envelope, EventIdStamp included, with its default serializer.

2. Declare the table

The entity is there for the schema: it declares the table for the migrations. Inserts go through DBAL, because DQL has no INSERT, and a duplicate flush() would fail.

// src/Entity/ProcessedEvent.php (stock service)
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: 'processed_event')]
#[ORM\Index(columns: ['processed_at'])]
class ProcessedEvent
{
    #[ORM\Id]
    #[ORM\Column(name: 'event_id', type: UuidType::NAME)]
    private Uuid $eventId;

    #[ORM\Column(name: 'processed_at', type: Types::DATETIMETZ_IMMUTABLE)]
    private \DateTimeImmutable $processedAt;
}

php bin/console make:migration generates the table's migration: event_id with the native UUID type as primary key, and an index on processed_at for the purge.

3. Write the middleware

// src/Messenger/ProcessedEventMiddleware.php (stock service)
namespace App\Messenger;

use Doctrine\DBAL\Connection;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Middleware\MiddlewareInterface;
use Symfony\Component\Messenger\Middleware\StackInterface;
use Symfony\Component\Messenger\Stamp\ReceivedStamp;

final class ProcessedEventMiddleware implements MiddlewareInterface
{
    public function __construct(
        private readonly Connection $connection,
    ) {
    }

    public function handle(Envelope $envelope, StackInterface $stack): Envelope
    {
        $eventId = $envelope->last(EventIdStamp::class)?->eventId;

        // Message sent by this service, or without event_id: nothing to deduplicate.
        if (null === $envelope->last(ReceivedStamp::class) || null === $eventId) {
            return $stack->next()->handle($envelope, $stack);
        }

        if (!$this->connection->isTransactionActive()) {
            throw new \LogicException('doctrine_transaction must come before ProcessedEventMiddleware.');
        }

        $inserted = $this->connection->executeStatement(
            'INSERT INTO processed_event (event_id, processed_at) VALUES (:eventId, now())
             ON CONFLICT (event_id) DO NOTHING',
            ['eventId' => $eventId],
        );

        if (0 === (int) $inserted) {
            // Already processed: no handler is called, the worker acknowledges the message.
            return $envelope;
        }

        return $stack->next()->handle($envelope, $stack);
    }
}

The middleware also runs when messages are sent. The check on ReceivedStamp limits it to messages received from a transport.

4. Configure the bus

# config/packages/messenger.yaml (stock service, excerpt)
framework:
  messenger:
    buses:
      messenger.bus.default:
        middleware:
          # Opens the transaction; flush() and commit if all goes well, rollback otherwise.
          - doctrine_transaction
          # Discards the events already processed, in this transaction.
          - App\Messenger\ProcessedEventMiddleware
    transports:
      events:
        dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
        serializer: App\Messenger\EventSerializer
        # exchange and queues options: see the article on the topology

The order of the list matters: Messenger runs these middleware in the declared order, before calling the handlers. The worker hands a message received without a bus name to the default bus. That is the case of a JSON envelope coming from another service.

5. Write the handler

The handler only contains the business effect. It does not call flush(): doctrine_transaction does it before the commit. Reservation is a classic Doctrine entity, with a UUID v7 id generated in its constructor.

// src/MessageHandler/ReserveStockOnOrderPlaced.php (stock service)
namespace App\MessageHandler;

use App\Entity\Reservation;
use AppContracts\Order\OrderPlacedV1;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;
use Symfony\Component\Uid\Uuid;

#[AsMessageHandler]
final class ReserveStockOnOrderPlaced
{
    public function __construct(
        private readonly EntityManagerInterface $entityManager,
    ) {
    }

    public function __invoke(OrderPlacedV1 $event): void
    {
        foreach ($event->lines as $line) {
            $reservation = new Reservation();
            $reservation->setOrderId(Uuid::fromString($event->orderId));
            $reservation->setProductId(Uuid::fromString($line['productId']));
            $reservation->setQuantity($line['quantity']); // decimal string, '2.000'

            $this->entityManager->persist($reservation);
        }
    }
}

6. A copy of state: add a version check

The shipping service keeps a copy of each order's status, to know whether a parcel can leave. If the producer publishes a version number of the order, incremented on every change, the upsert only applies newer versions.

-- shipping service: local copy of the order statuses
INSERT INTO order_status (order_id, status, version)
VALUES (:orderId, :status, :version)
ON CONFLICT (order_id) DO UPDATE
    SET status = EXCLUDED.status, version = EXCLUDED.version
    WHERE order_status.version < EXCLUDED.version;

The middleware already discards duplicates. The version check also discards an event that arrives after a newer one.

7. Test the double delivery

The test goes through the bus, not through the handler alone: it checks the transaction, the middleware and the handler together. With the ReceivedStamp, the message is handled on the spot, as if it came out of the queue.

// tests/Messenger/DoubleDeliveryTest.php (stock service)
namespace App\Tests\Messenger;

use App\Entity\Reservation;
use App\Messenger\EventIdStamp;
use AppContracts\Order\OrderPlacedV1;
use Doctrine\DBAL\Connection;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Bundle\FrameworkBundle\Test\KernelTestCase;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Messenger\Stamp\ReceivedStamp;
use Symfony\Component\Uid\Uuid;

final class DoubleDeliveryTest extends KernelTestCase
{
    public function testTheSameEventDeliveredTwiceReservesStockOnce(): void
    {
        $container = static::getContainer();
        $orderId = Uuid::v7();
        $eventId = (string) Uuid::v7();

        $event = new OrderPlacedV1(
            orderId: (string) $orderId,
            customerId: (string) Uuid::v7(),
            lines: [['productId' => (string) Uuid::v7(), 'quantity' => '2.000']],
        );

        $bus = $container->get(MessageBusInterface::class);
        $envelope = new Envelope($event, [new EventIdStamp($eventId), new ReceivedStamp('events')]);

        // Two deliveries of the same event, as after a crash before the acknowledgement.
        $bus->dispatch($envelope);
        $bus->dispatch($envelope);

        $reservations = $container->get(EntityManagerInterface::class)
            ->getRepository(Reservation::class)
            ->count(['orderId' => $orderId]);
        self::assertSame(1, $reservations);

        $processed = $container->get(Connection::class)->fetchOne(
            'SELECT count(*) FROM processed_event WHERE event_id = :eventId',
            ['eventId' => $eventId],
        );
        self::assertSame(1, (int) $processed);
    }
}

A second test is worth writing: make the handler fail, then check that processed_event does not contain the event_id. That is what lets the retry process the event again.

8. Purge according to the replay horizon

-- scheduled task of the stock service
DELETE FROM processed_event WHERE processed_at < :horizon;

A row can go once no copy of its event can arrive anymore. This horizon depends on how far back the producers can replay their outbox, and on how long messages can stay in the failure transport. Purging below it lets a replay process events that were already applied.

Checklist

  • A stable event_id, set by the producer in the JSON envelope, never generated on receipt.
  • A processed_event table (event_id UUID as primary key, processed_at) in the database of each consuming service.
  • The insert before any business effect, in the same transaction as that effect.
  • doctrine_transaction declared before the deduplication middleware.
  • INSERT ... ON CONFLICT DO NOTHING, and no handler called when no row is inserted.
  • Effects outside the database (email, HTTP call) protected by their own key.
  • A version check for copies of state that depend on ordering.
  • A test that delivers the same message twice and checks for a single effect.
  • A purge aligned with the replay horizon.

Sources