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.
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_eventtable (event_idUUID 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_transactiondeclared 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
- Symfony, Messenger: Writing Idempotent Handlers, Middleware for Doctrine, Middleware, Message Deduplication (Symfony 7.3), Message Serializer For Custom Data Formats.
- Symfony, source code of DoctrineTransactionMiddleware and RejectRedeliveredMessageMiddleware, 7.4 branch.
-
PostgreSQL:
INSERT, ON CONFLICT clause,
Index Uniqueness Checks,
Transactions
and
error codes
(
25P02). - Doctrine: DBAL, executeStatement() and ORM, no INSERT in DQL.
-
RabbitMQ:
Automatic Requeueing,
Reliability Guide
(at-least-once delivery,
redeliveredflag) and stream deduplication. - rabbitmq-message-deduplication, community plugin.