The problem
In our online shop, the order service exposes
POST /orders/{id}/place: the order moves from
the draft status to placed. It
then publishes order.order.placed.v1, so that
stock reserves the goods and
shipping prepares the shipment.
Two writes are involved: the order in PostgreSQL, the event in RabbitMQ. They share no transaction, and the sequence you pick decides how it fails.
Write, then publish. The commit succeeds, then the process stops before publishing: a redeployment, an out-of-memory error, an unreachable broker. The order is placed, but no service knows it: the stock will never be reserved.
Publish, then write. The event goes out,
then the transaction fails on a constraint or a conflict.
stock reserves the goods of an order still in
draft: a ghost event.
Two simultaneous requests. A double click,
or a client that retries, sends the same request twice. Both
read draft, apply the transition and publish
their own event. stock reserves twice.
The options
Publish from the request, before or after the commit
- What it solves: nothing to add, no table and no process.
-
What it costs: a lost event or a ghost
event, depending on the sequence chosen. In a handler
under
doctrine_transaction, theDispatchAfterCurrentBusStampsends after the commit: it avoids the ghost, not the loss.
A distributed transaction
- What it solves: in theory, the database and the broker commit together.
- What it costs: RabbitMQ does not take part. Its AMQP transactions only cover the publishes and acknowledgements of a channel, with no link to PostgreSQL.
An outbox read by a relay
- What it solves: the event is inserted into a table, in the same transaction as the order. A relay reads that table and publishes. Nothing is lost, nothing is published by mistake.
- What it costs: a table, a process, a latency tied to the polling interval, and possible duplicates.
Messenger's Doctrine transport can act as an outbox. But its retries break the ordering, and it deletes handled messages by default.
An outbox read by change data capture (CDC)
- What it solves: Debezium reads PostgreSQL's transaction log (WAL) through logical decoding. Its "Outbox Event Router" transformation turns each inserted row into a message.
- What it costs: infrastructure to run, Kafka Connect or Debezium Server and its RabbitMQ sink. And a replication slot retains the WAL until the connector has read it.
Our recommendation
Write the event to an outbox table, in the
same transaction as the order, then publish it with a
relay that reads that table.
The PostgreSQL transaction is the only atomicity we have: the order and its event go into it together. We move to CDC if polling weighs on the database, with measurements to back it up.
Write the outbox on the business fact
The event is born when the order moves to
placed. With Symfony Workflow, that is the
workflow.order.completed.place event,
dispatched once the transition is applied. A listener
inserts the outbox row there.
We avoid Doctrine postPersist or
postUpdate listeners: they see changed columns,
not an intent. Fixing an address triggers the same hook as
placing an order. And Doctrine does not write, in the
ongoing flush, what you persist there.
We do not write the outbox in the API Platform processors either: a transition run from the console would emit no event.
Lock the order before the transition
The processor reads the order with
LockMode::PESSIMISTIC_WRITE, that is a
SELECT ... FOR UPDATE. The second simultaneous
request waits for the first one to end, reads
placed again and fails cleanly. A row lock does
not block plain reads.
Prepare the event before the commit
Identifiers are UUID v7 generated in PHP: the event carries
the order's and its own before the flush().
Before the insert, we validate the envelope against its JSON
Schema. An invalid message rolls back the whole transaction:
the error shows up at the producer, not at a consumer.
Publish, mark, commit
The relay reads a batch of unpublished rows with
FOR UPDATE SKIP LOCKED. For each row, it
publishes, waits for the broker's confirmation (publisher
confirm), then sets published_at. It commits at
the end of the batch.
This sequence loses nothing. A stop between the publish and the commit rolls back the marking, and the row goes out again in the next batch. The consumer then gets a duplicate.
The confirmation means the broker has taken charge of the message, on disk if it is persistent and its queue durable. A message that no binding key matches is confirmed too, then dropped. An alternate exchange can collect it: see our article on the RabbitMQ topology.
Keep the published rows
Published rows form a dated log, useful to investigate or replay an event. A scheduled task purges what goes beyond the replay horizon you set.
The trade-offs we accept
A latency tied to polling. An event waits for the relay's next run. Reserving stock tolerates that delay; a real-time display, less so.
One more table and one more process. We
alert on the age of the oldest unpublished row and on the
attempts column.
Duplicates. The relay delivers "at least once". Each consumer must be idempotent, for instance with a table of processed events: that is the topic of a dedicated article.
A limited ordering. The ordering relies on
occurred_at, a single relay and stopping at the
first failure: a failing row blocks the following ones. For
several relays, split the rows by order: ordering only
matters within one order.
One confirmation at a time. Messenger waits for the confirmation of each message before the next one. For a persistent message, RabbitMQ states that it can take a few hundred milliseconds under load. The relay's throughput depends on it.
The signal that should trigger a review of this choice: a delay the consumers no longer tolerate, or a database loaded by the relay. CDC then becomes relevant. Regularly replaying long periods rather calls for a log: Kafka, or a RabbitMQ stream read by a client other than Messenger.
Implementation
The examples use Symfony 7.4, Doctrine ORM 3 (DBAL 4.3 or later), PostgreSQL and API Platform 4.
1. The outbox table
// src/Entity/OutboxMessage.php (order 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: 'outbox')]
#[ORM\Index(
name: 'outbox_unpublished_idx',
fields: ['occurredAt'],
options: ['where' => 'published_at IS NULL'],
)]
class OutboxMessage
{
#[ORM\Column(type: Types::DATETIMETZ_IMMUTABLE, nullable: true)]
private ?\DateTimeImmutable $publishedAt = null;
#[ORM\Column]
private int $attempts = 0;
public function __construct(
// Same value as the envelope's event_id.
#[ORM\Id]
#[ORM\Column(type: UuidType::NAME)]
private Uuid $id,
#[ORM\Column(length: 255)]
private string $eventName,
#[ORM\Column(length: 255)]
private string $routingKey,
// The full envelope, published as is.
#[ORM\Column(type: Types::JSONB)]
private array $envelope,
#[ORM\Column(type: Types::DATETIMETZ_IMMUTABLE)]
private \DateTimeImmutable $occurredAt,
) {
}
// getRoutingKey(), getEnvelope(), getAttempts(), setAttempts(), setPublishedAt().
}
The Types::JSONB type exists since DBAL 4.3,
which deprecates the jsonb option of the
json type.
2. The order and its workflow
The places are the constants of a
final class OrderStatus (DRAFT,
PLACED, CANCELLED), stored in the
status column.
# config/packages/workflow.yaml (order service)
framework:
workflows:
order:
type: state_machine
marking_store:
type: method
property: status
supports:
- App\Entity\Order
initial_marking: draft
places:
- draft
- placed
- cancelled
transitions:
place:
from: draft
to: placed
// src/Entity/Order.php (order service)
namespace App\Entity;
use Doctrine\ORM\Mapping as ORM;
use Symfony\Bridge\Doctrine\Types\UuidType;
use Symfony\Component\Uid\Uuid;
#[ORM\Entity]
#[ORM\Table(name: 'orders')]
class Order
{
#[ORM\Id]
#[ORM\Column(type: UuidType::NAME)]
private Uuid $id;
#[ORM\Column(length: 20)]
private string $status = OrderStatus::DRAFT;
// reference, customer, lines...
public function __construct(string $reference, Customer $customer)
{
// UUID v7 generated in PHP: the identifier exists before the flush().
$this->id = Uuid::v7();
// ...
}
// getId(), getStatus(), setStatus(), getCustomer(), getLines().
}
Order stays a Doctrine entity without any API
Platform attribute. The operation is added to the
OrderResource DTO resource, tied to the entity
by stateOptions.
read: false prevents API Platform from reading
the order before the processor, hence outside the
transaction.
// src/ApiResource/OrderResource.php (order service), excerpt
namespace App\ApiResource;
use ApiPlatform\Doctrine\Orm\State\Options;
use ApiPlatform\Metadata\ApiResource;
use ApiPlatform\Metadata\Post;
use App\Entity\Order;
use App\State\PlaceOrderProcessor;
#[ApiResource(
shortName: 'Order',
operations: [
// GetCollection, Get, Post and Patch, then:
new Post(
uriTemplate: '/orders/{id}/place',
status: 200,
input: false,
read: false,
processor: PlaceOrderProcessor::class,
),
],
stateOptions: new Options(entityClass: Order::class),
)]
final readonly class OrderResource
{
// ...
}
3. The processor: transaction, lock, transition
// src/State/PlaceOrderProcessor.php (order service)
namespace App\State;
use ApiPlatform\Metadata\Operation;
use ApiPlatform\State\ProcessorInterface;
use App\ApiResource\OrderResource;
use App\Entity\Order;
use App\Mapper\OrderMapper;
use Doctrine\DBAL\LockMode;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\HttpKernel\Exception\ConflictHttpException;
use Symfony\Component\HttpKernel\Exception\NotFoundHttpException;
use Symfony\Component\Workflow\WorkflowInterface;
/** @implements ProcessorInterface<null, OrderResource> */
final class PlaceOrderProcessor implements ProcessorInterface
{
public function __construct(
private readonly EntityManagerInterface $em,
private readonly WorkflowInterface $orderStateMachine,
private readonly OrderMapper $mapper,
) {
}
public function process(mixed $data, Operation $operation, array $uriVariables = [], array $context = []): OrderResource
{
// BEGIN, then flush() and COMMIT, or ROLLBACK on exception.
$order = $this->em->wrapInTransaction(function () use ($uriVariables): Order {
// SELECT ... FOR UPDATE: a concurrent request waits here.
$order = $this->em->find(Order::class, $uriVariables['id'], LockMode::PESSIMISTIC_WRITE)
?? throw new NotFoundHttpException('Order not found.');
if (!$this->orderStateMachine->can($order, 'place')) {
throw new ConflictHttpException('The order has already left the draft status.');
}
// The listener inserts the outbox row during apply().
$this->orderStateMachine->apply($order, 'place');
return $order;
});
return $this->mapper->toResource($order);
}
}
find() with a lock mode requires an active
transaction and reads the locked row again, even when the
entity is already loaded. lock(), on the other
hand, sets the lock without reloading the entity.
4. The listener: insert the outbox row
// src/EventListener/OrderPlacedOutboxListener.php (order service)
namespace App\EventListener;
use App\Entity\Order;
use App\Outbox\OutboxWriter;
use AppContracts\Order\OrderPlacedV1;
use AppContracts\Routing;
use Symfony\Component\Workflow\Attribute\AsCompletedListener;
use Symfony\Component\Workflow\Event\CompletedEvent;
final class OrderPlacedOutboxListener
{
public function __construct(
private readonly OutboxWriter $outbox,
) {
}
#[AsCompletedListener(workflow: 'order', transition: 'place')]
public function onPlace(CompletedEvent $event): void
{
/** @var Order $order */
$order = $event->getSubject();
$lines = [];
foreach ($order->getLines() as $line) {
$lines[] = [
'productId' => $line->getProduct()->getId()->toRfc4122(),
'quantity' => $line->getQuantity(), // for instance '2.000'
];
}
$this->outbox->add(
Routing::ORDER_PLACED_V1,
new OrderPlacedV1(
orderId: $order->getId()->toRfc4122(),
customerId: $order->getCustomer()->getId()->toRfc4122(),
lines: $lines,
),
correlationId: $order->getId()->toRfc4122(),
);
}
}
// src/Outbox/OutboxWriter.php (order service)
namespace App\Outbox;
use App\Entity\OutboxMessage;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\Serializer\Normalizer\AbstractObjectNormalizer;
use Symfony\Component\Serializer\Normalizer\NormalizerInterface;
use Symfony\Component\Uid\Uuid;
final class OutboxWriter
{
public function __construct(
private readonly EntityManagerInterface $em,
private readonly NormalizerInterface $normalizer,
) {
}
public function add(string $eventName, object $event, string $correlationId): void
{
$eventId = Uuid::v7();
$occurredAt = new \DateTimeImmutable();
$envelope = [
'event_id' => $eventId->toRfc4122(),
'event_name' => $eventName,
'occurred_at' => $occurredAt->format(\DateTimeInterface::RFC3339_EXTENDED),
'producer' => 'order',
'schema_version' => $event::SCHEMA_VERSION,
'correlation_id' => $correlationId,
// Absence rather than null: empty fields are not written.
'payload' => $this->normalizer->normalize($event, 'json', [
AbstractObjectNormalizer::SKIP_NULL_VALUES => true,
]),
];
// Validate json_encode($envelope) against its JSON Schema here. No flush(): the processor does it.
$this->em->persist(new OutboxMessage(
id: $eventId,
eventName: $eventName,
routingKey: $eventName,
envelope: $envelope,
occurredAt: $occurredAt,
));
}
}
5. The relay
// src/Outbox/OutboxRelay.php (order service)
namespace App\Outbox;
use App\Entity\OutboxMessage;
use Doctrine\ORM\EntityManagerInterface;
use Doctrine\ORM\Query\ResultSetMappingBuilder;
use Symfony\Component\DependencyInjection\Attribute\Autowire;
use Symfony\Component\Messenger\Bridge\Amqp\Transport\AmqpStamp;
use Symfony\Component\Messenger\Bridge\Amqp\Transport\AmqpTransport;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Exception\TransportException;
final class OutboxRelay
{
public const BATCH_SIZE = 100;
public function __construct(
private readonly EntityManagerInterface $em,
#[Autowire(service: 'messenger.transport.events')]
private readonly AmqpTransport $events,
) {
}
public function relay(): int
{
$published = $this->em->wrapInTransaction(function (): int {
$rsm = new ResultSetMappingBuilder($this->em);
$rsm->addRootEntityFromClassMetadata(OutboxMessage::class, 'o');
$messages = $this->em->createNativeQuery(
'SELECT '.$rsm->generateSelectClause(['o' => 'o']).' FROM outbox o
WHERE o.published_at IS NULL
ORDER BY o.occurred_at, o.id
LIMIT '.self::BATCH_SIZE.'
FOR UPDATE SKIP LOCKED',
$rsm,
)->getResult();
$count = 0;
foreach ($messages as $message) {
try {
// 1. Publish: send() returns once the broker has confirmed.
$this->events->send(new Envelope($message, [new AmqpStamp($message->getRoutingKey())]));
} catch (TransportException) {
// Nack or timeout: stop here to keep the ordering.
// New channel: a late ack will not confirm the next one.
$this->events->close();
$message->setAttempts($message->getAttempts() + 1);
break;
}
// 2. Mark.
$message->setPublishedAt(new \DateTimeImmutable());
++$count;
}
return $count;
}); // 3. Commit: wrapInTransaction() calls flush() then COMMIT.
$this->em->clear();
return $published;
}
}
The relay calls the events transport directly,
without going through the bus. It replaces the
EventPublisher and the
routing rules of
the topology article. Its serializer publishes the stored envelope as is.
// src/Outbox/OutboxSerializer.php (order service)
namespace App\Outbox;
use App\Entity\OutboxMessage;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Transport\Serialization\SerializerInterface;
final class OutboxSerializer implements SerializerInterface
{
public function encode(Envelope $envelope): array
{
/** @var OutboxMessage $message */
$message = $envelope->getMessage();
return [
'body' => json_encode($message->getEnvelope(), \JSON_THROW_ON_ERROR),
'headers' => ['Content-Type' => 'application/json'],
];
}
public function decode(array $encodedEnvelope): Envelope
{
throw new \LogicException('The order service does not consume this transport.');
}
}
The confirm_timeout option turns publisher
confirms on. A nack or a timeout throws a
TransportException. The transport already
publishes in persistent mode.
# config/packages/messenger.yaml (order service)
framework:
messenger:
transports:
events:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
serializer: App\Outbox\OutboxSerializer
options:
exchange:
name: !php/const AppContracts\Routing::EXCHANGE
type: topic
queues: []
# Wait for the broker's confirmation, 5 seconds at most.
confirm_timeout: 5
6. Run the relay
Symfony Scheduler triggers the relay every second. As long as the batches are full, the task chains them.
// src/Scheduler/RelayOutboxTask.php (order service)
namespace App\Scheduler;
use App\Outbox\OutboxRelay;
use Symfony\Component\Scheduler\Attribute\AsPeriodicTask;
#[AsPeriodicTask(frequency: 1)]
final class RelayOutboxTask
{
public function __construct(
private readonly OutboxRelay $relay,
) {
}
public function __invoke(): void
{
while (OutboxRelay::BATCH_SIZE === $this->relay->relay()) {
}
}
}
We run
php bin/console messenger:consume scheduler_default
as a single instance. If two relays overlap,
SKIP LOCKED hands them different rows, but the
ordering is no longer guaranteed.
Checklist
- The outbox written in the same transaction as the business change.
- A write on the business fact, not in a Doctrine hook.
- A lock on concurrent transitions.
- UUID v7 generated in PHP.
- An envelope validated against its JSON Schema.
- Publisher confirms, and the sequence: publish, mark, commit.
- A single active relay, stopped at the first failure.
- Idempotent consumers.
- An alert on the age of the oldest unpublished row.
- A defined retention and a scheduled purge.
Sources
- microservices.io, Transactional outbox.
-
Symfony, Messenger:
AMQP transport
(
confirm_timeout), serialization, Doctrine transport (delete_after_ack), transactional messages. - Symfony, AMQP transport code: waiting for the confirmation, persistent mode by default.
- Symfony, Workflow, events, Scheduler and Uid component.
-
Doctrine ORM,
Transactions and Concurrency,
Events
and
EntityManager code
(
find()under lock). -
Doctrine DBAL,
Types
and
4.3 upgrade notes:
jsonb. -
API Platform,
State Processors
and
main controller: the
readoption. - PostgreSQL, locking clause, row-level locks and logical decoding.
- RabbitMQ, Publisher Confirms and Semantics (AMQP transactions).
- Debezium, Outbox Event Router, PostgreSQL connector and Debezium Server.