The Transactional Outbox Pattern
Why saving to the database and publishing an event can never be atomic, how an outbox table plus a relay fixes it, and what ordering and duplicates still cost.
A service that saves an order and then publishes OrderPlaced to a broker is doing two writes to two systems, and no crash-proof ordering of them exists: commit first and a crash loses the event, publish first and a rollback leaves an event about an order that never existed. The transactional outbox turns the second write into a row in an outbox table, committed in the same database transaction as the order, and a separate relay reads that table and publishes. The event can no longer be lost; it can be published twice, so consumers must deduplicate.
Context
The problem appeared as soon as systems kept state in one place and notified others through another. The classic answer was a distributed transaction across both: XA (the X/Open eXtended Architecture standard, 1991) with 2PC (two-phase commit), run by application servers in the Java EE era. Most modern brokers, including Kafka and RabbitMQ, do not take part in XA, and Pat Helland's 2007 paper Life beyond Distributed Transactions argued that scalable systems should not rely on it anyway. What replaced it is local transactions plus asynchronous messaging, and the outbox is the bridge between the two. It is catalogued as "Transactional outbox" on microservices.io and is built into tooling such as Debezium's outbox event router, NServiceBus, MassTransit and Axon.
You have met the problem if you have ever written this, which is how it usually starts:
await db.transaction(async tx => {
await tx.insert(orders).values(order)
}) // committed
await kafka.send({topic: 'orders', // what if we crash here,
messages: [{key: order.id, // or Kafka is down?
value: JSON.stringify({type: 'OrderPlaced', ...order})}]})- Dual write
- Writing the same fact to two systems (a database and a broker, a database and a search index) from application code, with no transaction spanning both.
- Outbox table
- A table in the service’s own database holding events waiting to be published, inserted in the same transaction as the business change.
- Relay
- The process that reads the outbox and publishes each row to the broker: either a poller querying the table or a CDC (change data capture) reader tailing the database log.
- Aggregate
- The entity an event is about, such as one order. Events for the same aggregate must be published in the order they were written.
- Inbox table
- The consumer-side mirror: a table of processed message IDs, written in the same transaction as the consumer’s own change, so a redelivered message is ignored.
Why it matters
In an event-driven system, other services act on events: the warehouse reserves stock, billing charges the card, search reindexes the product. A lost event is a silent inconsistency, an order that is never shipped or a price that never updates, with no error anywhere and usually discovered by a customer. A ghost event is worse: an email for an order that was rolled back, or stock reserved for nothing. Dual writes fail rarely enough to pass every test and often enough to cause weekly incidents at scale, because deploys, out-of-memory kills and broker failovers all happen between the two writes eventually.
Why the two writes cannot be ordered safely
The database commit and the broker publish are two independent operations with no shared coordinator. Whichever comes first, a crash, a timeout or a failed second write between them leaves the two systems disagreeing, and the application cannot tell from the outside which one happened.
| Order of writes | Failure between them | Result |
|---|---|---|
| Commit, then publish | Crash, deploy or broker down after commit | Lost event: other services never learn about the order |
| Publish, then commit | Commit fails (constraint, deadlock, crash) | Ghost event: consumers act on an order that does not exist |
| Publish inside the tx | Commit fails after the send returned | Ghost event, and the transaction held locks during a network call |
Why not a distributed transaction
2PC across the database and the broker would make both writes atomic, but it needs both to support XA (Kafka and RabbitMQ do not; Kafka's own transactions, since 0.11 in 2017, only make writes atomic within Kafka), it holds locks across a network round-trip in every request, and a coordinator failure can leave participants blocked holding those locks. The outbox gets the guarantee that matters, nothing lost and nothing invented, from a single local transaction.
The outbox: one transaction, then a relay
The fix moves the publish out of the request. The service writes the order and an outbox row describing the event in the same local transaction, so either both are committed or neither is. Publishing becomes someone else's job: a relay reads committed outbox rows, sends them to the broker, and records that they were sent. If the relay crashes after sending but before recording, it sends again on restart.
CREATE TABLE outbox (
id uuid PRIMARY KEY, -- becomes the message ID
aggregate_id text NOT NULL, -- e.g. the order ID
type text NOT NULL, -- 'OrderPlaced'
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
published_at timestamptz -- NULL = not sent yet
);
CREATE INDEX outbox_unpublished ON outbox (created_at)
WHERE published_at IS NULL;
-- in the request: both rows or neither
BEGIN;
INSERT INTO orders (id, customer_id, total) VALUES ($1, $2, $3);
INSERT INTO outbox (id, aggregate_id, type, payload)
VALUES (gen_random_uuid(), $1, 'OrderPlaced', $4);
COMMIT;Two ways to build the relay
A polling publisher queries the outbox for unpublished rows every few hundred milliseconds. With FOR UPDATE SKIP LOCKED (PostgreSQL since 9.5, MySQL since 8.0) several relay instances can run without picking the same rows. A log-tailing relay uses CDC instead: it reads the database's replication log (PostgreSQL logical decoding, the MySQL binlog) and turns every inserted outbox row into a message. Debezium does this and ships an outbox event router that maps the row's columns to the Kafka topic, key and payload.
-- one relay iteration, inside a transaction
BEGIN;
SELECT id, aggregate_id, type, payload
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED; -- other relays skip these rows
-- publish each row to the broker, wait for the broker's ack, then:
UPDATE outbox SET published_at = now() WHERE id = ANY($1);
COMMIT;
-- crash before COMMIT: rows unlock and are sent again (duplicates)| Polling publisher | Log tailing (CDC) | |
|---|---|---|
| Latency | Polling interval, typically 100 ms - 1 s | Near real time, as the log is written |
| DB load | A query per tick, plus UPDATEs | Reads the replication stream, no queries |
| Moving parts | A loop in your service | Debezium or similar, plus Kafka Connect |
| Ordering | Needs care with concurrent relays | Commit order from the log |
| Cleanup | Mark published, delete later | Rows can be deleted right after insert |
Ordering, duplicates and cleanup
Ordering per aggregate
Consumers usually need events about the same order in sequence (OrderPlaced before OrderCancelled), not a global order. Publish with the aggregate ID as the message key so a partitioned broker like Kafka keeps them on one partition. On the relay side, two pollers with SKIP LOCKED can send consecutive events for one order in parallel and swap them. Either run a single active relay, or have each relay instance own a subset of aggregates, for example by hashing the aggregate ID. CDC preserves commit order naturally.
Duplicates and the inbox
The outbox guarantees at-least-once publishing, so every consumer must tolerate the same message twice. The robust version is an inbox table: the consumer records the message ID in the same transaction as its own change, and a unique constraint turns a redelivery into a no-op.
async function handle(msg: {id: string; type: string; payload: Order}) {
await db.transaction(async tx => {
const inserted = await tx.execute(sql`
INSERT INTO inbox (message_id) VALUES (${msg.id})
ON CONFLICT (message_id) DO NOTHING
RETURNING message_id`)
if (inserted.rows.length === 0) return // seen before: skip
await reserveStock(tx, msg.payload) // same transaction
})
// ack the broker only after the commit
}Cleanup
An outbox that is never cleaned grows by one row per event forever, and the polling query slows down with it. Delete published rows on a schedule (or partition the table by day and drop old partitions), keeping a few days for debugging and replay. With CDC, the row only has to exist long enough to reach the log, so many setups insert and delete it in the same transaction.
Walking through the failures
The pattern earns its keep in the failure cases. Each one either cannot happen any more or turns into a duplicate the consumer already handles:
- 1Crash before COMMIT. Neither the order nor the outbox row exists. The client gets an error and retries; nothing was announced.
- 2Crash right after COMMIT. Both rows exist. The relay finds the unpublished outbox row on its next tick and publishes it. Nothing is lost.
- 3Broker down for ten minutes. Requests keep succeeding; outbox rows pile up. The relay retries and drains the backlog when the broker returns. The outbox is a buffer, so alert on its size.
- 4Relay crashes after publishing, before marking. The rows are still unpublished in the table and are sent again. Consumers see a duplicate with the same message ID and the inbox drops it.
- 5Consumer crashes after its change, before acking. Its transaction already recorded the message ID, so the redelivered message is skipped. Without the inbox, stock would be reserved twice.
Pitfalls
- Publishing from the request "just this once"
A second code path that sends to the broker directly after commit brings the dual write back for that event type. All publishing has to go through the outbox, or the guarantee holds only for some events and nobody knows which.
- Consumers that are not idempotent
The relay will publish some messages twice, because it cannot mark a row and publish it atomically either. A consumer that sends an email or charges a card on every delivery turns each relay restart into duplicate side effects. Dedup by message ID in the same transaction as the change.
- Paging the outbox by an auto-increment ID
Concurrent transactions commit in a different order than their IDs were assigned, so a cursor of "last ID sent" skips rows that commit late. The loss is rare, silent and hard to reproduce. Use a published flag, or the database log.
- Parallel relays reordering one aggregate
Scaling the poller with
SKIP LOCKEDspreads rows of one order across instances, which then publish them in either order. Partition relay work by aggregate, keep one active relay, or include a per-aggregate version so consumers can detect and drop out-of-order events. - Never cleaning the table
Millions of published rows bloat the table and the index the poller scans, and the relay's query gets slower exactly when a backlog needs it to be fast. Delete or drop published rows on a schedule and monitor the unpublished count and the age of the oldest unpublished row.
Interview questions
Q1What is the dual write problem?
It is updating two systems, typically a database and a message broker, from application code with no transaction covering both. Any failure between the two writes leaves them inconsistent: commit then crash loses the event, publish then a failed commit announces something that never happened. No ordering of the two calls avoids it, because the failure can always land between them.
Q2Walk me through implementing a transactional outbox in a PostgreSQL service.
I add an outbox table with an ID, aggregate ID, type, JSON payload and a published_at column. Every command inserts its business rows and the event row in one transaction. A relay loop selects unpublished rows ordered by creation with FOR UPDATE SKIP LOCKED, publishes them keyed by aggregate ID, waits for the broker's acknowledgement, then marks them published and commits. A scheduled job deletes old published rows, and I alert on the oldest unpublished row's age.
Q3What happens when the relay crashes after publishing but before marking the rows?
The rows are still unpublished, so after restart the relay sends them again and consumers receive duplicates. That is the expected at-least-once behaviour, not a bug. Consumers handle it by deduplicating on the message ID, ideally with an inbox table written in the same transaction as their own change.
Q4Why not use a distributed transaction instead?
Because most brokers, Kafka and RabbitMQ included, do not participate in XA, and 2PC holds locks across a network round-trip on every request and can block if the coordinator fails. The outbox gives the needed guarantee, nothing lost and nothing invented, from one local transaction, and moves the cross-system part into an asynchronous retry loop.
Q5Polling the outbox or CDC with something like Debezium: how do you choose?
Polling is the simplest: a loop in the service, no extra infrastructure, latency of one polling interval and some query load. CDC reads the database log, so it is near real time, keeps commit order and puts no query load on the table, at the cost of running Debezium and Kafka Connect and managing replication slots. I start with polling and move to CDC when latency, volume or ordering demands it.
Q6How do you keep events for one order in order?
Publish with the order ID as the message key, so a partitioned broker keeps that order's events on one partition in sequence. On the relay side, make sure two instances cannot publish rows for the same aggregate concurrently: a single active relay, work partitioned by aggregate hash, or CDC, which follows commit order. A version number per aggregate lets consumers detect anything that still arrives out of order.
Q7Does the outbox give you exactly-once delivery?
No, it gives at-least-once publishing with no lost or phantom events. Exactly-once effects come from combining it with idempotent consumers: the consumer stores the message ID and its own change atomically, so a duplicate delivery changes nothing.
- Saving to a database and publishing to a broker are two writes with no shared transaction; any ordering of them can lose an event or invent one.
- The outbox writes the event as a row in the same local transaction as the business change, so both commit or neither does.
- A relay publishes outbox rows, by polling with FOR UPDATE SKIP LOCKED or by tailing the database log with CDC, and retries until the broker accepts them.
- Publishing is at-least-once, so consumers deduplicate by message ID, ideally with an inbox table in the same transaction as their change.
- Keep per-aggregate order with the aggregate ID as the message key and no concurrent relays per aggregate; never page the outbox by an increasing ID.
- Clean up published rows and alert on the age of the oldest unpublished one: the outbox is also a buffer.