transactional outbox
Transactional Outbox
Understand the dual-write problem that silently corrupts most event-driven systems, how a single local transaction solves it, and why this pattern is the practical foundation of reliable messaging.
Prerequisites
Recommended Knowledge
- Database transactions and atomicity
- Message brokers and delivery semantics
- Two-phase commit and its costs
- Idempotency and duplicate handling
- Partial failure and timeout ambiguity
- Ordering guarantees and partition keys
- Change data capture concepts
- Sagas and event-driven workflows
The Dual-Write Problem
Almost every event-driven service needs to do two things when something happens: persist the change and tell other services about it. These are two separate systems, and there is no transaction spanning both.
// This code is wrong, and the reason is not obvious
async function placeOrder(order) {
await database.orders.insert(order); // Step one
await broker.publish("OrderPlaced", order); // Step two
}
| Failure Point | Database | Broker | Outcome |
|---|---|---|---|
| Before either | No record | No message | Consistent, order lost |
| After insert, before publish | Order exists | No message | Order never processed downstream |
| Publish times out ambiguously | Order exists | Unknown | May duplicate on retry |
| After both | Order exists | Message sent | Correct |
Simple Analogy
Recording a decision in a ledger, then walking to the post office to notify everyone affected. If you collapse on the way, the decision stands but nobody knows. Writing the letter into the ledger itself removes the walk.
Reversing the Order Does Not Help
// Also wrong, in the opposite direction
async function placeOrder(order) {
await broker.publish("OrderPlaced", order); // Now first
await database.orders.insert(order); // Now second
}
Why Retries Do Not Solve It
| Attempted Fix | Why It Fails |
|---|---|
| Retry the publish in a loop | Process can die mid-loop |
| Publish in a finally block | Process termination skips it entirely |
| Background retry queue in memory | Memory is lost on restart |
| Two-phase commit across both | Blocking, slow, often unsupported |
| Ignore it and reconcile later | Reconciliation needs the very events you lost |
The Outbox Solution
The insight is simple: if the two writes must be atomic, make them the same write. Record the message in the same database, in the same transaction, then deliver it separately.
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
aggregate_type VARCHAR(80) NOT NULL,
aggregate_id VARCHAR(80) NOT NULL,
event_type VARCHAR(80) NOT NULL,
payload JSONB NOT NULL,
headers JSONB NULL,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
published_at TIMESTAMP NULL,
attempts INT NOT NULL DEFAULT 0
);
CREATE INDEX idx_outbox_unpublished
ON outbox (id)
WHERE published_at IS NULL;
BEGIN;
INSERT INTO orders (order_id, customer_id, total, status, created_at)
VALUES (:order_id, :customer_id, :total, 'placed', CURRENT_TIMESTAMP);
INSERT INTO outbox (
aggregate_type, aggregate_id, event_type, payload
)
VALUES (
'order',
:order_id,
'OrderPlaced',
:event_payload
);
COMMIT;
async function placeOrder(database, order) {
return await database.transaction(async (tx) => {
const saved = await tx.orders.insert(order);
await tx.outbox.insert({
aggregateType: "order",
aggregateId: saved.orderId,
eventType: "OrderPlaced",
payload: {
orderId: saved.orderId,
customerId: saved.customerId,
items: saved.items,
total: saved.total,
placedAt: saved.createdAt
}
});
return saved;
});
}
The Message Relay
A separate process reads unpublished outbox rows and delivers them to the broker. Because the rows are durable, it can retry indefinitely without losing anything.
class OutboxRelay {
async pollOnce() {
const batch = await this.db.query(`
SELECT id, aggregate_type, aggregate_id,
event_type, payload, headers
FROM outbox
WHERE published_at IS NULL
ORDER BY id
LIMIT $1
FOR UPDATE SKIP LOCKED
`, [this.batchSize]);
if (batch.length === 0) return { published: 0 };
let published = 0;
for (const row of batch) {
try {
await this.broker.publish({
topic: this.topicFor(row.event_type),
key: row.aggregate_id,
value: row.payload,
headers: {
...row.headers,
messageId: String(row.id),
eventType: row.event_type
}
});
await this.markPublished(row.id);
published += 1;
} catch (error) {
await this.recordAttempt(row.id, error);
break;
}
}
return { published };
}
async markPublished(id) {
await this.db.query(`
UPDATE outbox
SET published_at = CURRENT_TIMESTAMP
WHERE id = $1
`, [id]);
}
}
Polling Versus Change Data Capture
| Aspect | Polling Relay | Change Data Capture |
|---|---|---|
| Mechanism | Query for unpublished rows | Tail the database write-ahead log |
| Latency | Bounded by poll interval | Near-immediate |
| Database load | Continuous queries | Minimal, reads the log |
| Operational complexity | Low, ordinary application code | Higher, connector infrastructure |
| Ordering | By primary key within a batch | Natural log order |
| Cleanup | Mark published, then delete | Delete immediately after insert |
| Database support | Any database | Requires log access |
The CDC Variant
-- With CDC, the row can be deleted in the same transaction
BEGIN;
INSERT INTO orders (order_id, customer_id, total, status)
VALUES (:order_id, :customer_id, :total, 'placed');
INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
VALUES ('order', :order_id, 'OrderPlaced', :payload);
DELETE FROM outbox WHERE aggregate_id = :order_id;
COMMIT;
-- The insert and delete both appear in the log; the connector
-- reads the insert and publishes, while the table stays empty
Ordering
The outbox preserves the order in which events were written, but only if the relay respects it and the broker partitions accordingly.
| Requirement | Mechanism | Failure If Omitted |
|---|---|---|
| Read in insertion order | Order by monotonic identifier | Events published out of sequence |
| Publish sequentially per key | Stop on failure, do not skip ahead | Later event arrives before earlier |
| Partition by aggregate | Use aggregate identifier as key | Broker reorders across partitions |
| Single relay per partition | Lock or shard the relay | Concurrent relays interleave |
async function publishInOrder(rows, broker) {
const byAggregate = new Map();
for (const row of rows) {
const key = row.aggregate_id;
if (!byAggregate.has(key)) byAggregate.set(key, []);
byAggregate.get(key).push(row);
}
// Different aggregates in parallel; within one, strictly sequential
await Promise.all(
Array.from(byAggregate.entries()).map(async ([key, group]) => {
for (const row of group) {
try {
await broker.publish({ key, value: row.payload });
await markPublished(row.id);
} catch (error) {
return { key, stoppedAt: row.id, error };
}
}
})
);
}
Running Multiple Relay Instances
A single relay is a throughput ceiling and a single point of failure. Running several requires preventing them from publishing the same rows.
| Approach | Mechanism | Ordering Preserved |
|---|---|---|
| Leader election | One active relay at a time | Fully |
| Row locking | Skip rows locked by others | Not across instances |
| Shard by aggregate hash | Each relay owns a hash range | Within each aggregate |
| Partitioned outbox table | One relay per partition | Within each partition |
-- Shard assignment keeps each aggregate on one relay
SELECT id, aggregate_id, event_type, payload
FROM outbox
WHERE published_at IS NULL
AND MOD(ABS(HASHTEXT(aggregate_id)), :total_shards) = :my_shard
ORDER BY id
LIMIT :batch_size
FOR UPDATE SKIP LOCKED;
Cleanup and Growth
Published rows must be removed or the table grows without bound, degrading the very query the relay depends on.
-- Delete in bounded batches to avoid long locks
DELETE FROM outbox
WHERE id IN (
SELECT id
FROM outbox
WHERE published_at IS NOT NULL
AND published_at < CURRENT_TIMESTAMP - INTERVAL '7 days'
ORDER BY id
LIMIT 10000
);
| Concern | Problem | Mitigation |
|---|---|---|
| Unbounded growth | Poll query slows over time | Scheduled deletion job |
| Index bloat | Partial index degrades | Periodic maintenance |
| Long delete transactions | Locks block the relay | Bounded batch deletes |
| Deleting too early | Lose replay capability | Retain for a defined window |
| Table scan on poll | Full scan as rows accumulate | Partial index on unpublished only |
The Inbox Pattern
The outbox guarantees at-least-once delivery. The receiving side needs a corresponding mechanism to make duplicate delivery harmless.
CREATE TABLE inbox (
message_id VARCHAR(120) PRIMARY KEY,
consumer VARCHAR(80) NOT NULL,
processed_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
async function consumeMessage(database, message, handler) {
return await database.transaction(async (tx) => {
const claimed = await tx.query(`
INSERT INTO inbox (message_id, consumer)
VALUES ($1, $2)
ON CONFLICT (message_id) DO NOTHING
RETURNING message_id
`, [message.id, handler.name]);
if (claimed.rowCount === 0) {
return { status: "duplicate-ignored" };
}
await handler.process(tx, message);
return { status: "processed" };
});
}
What to Put in the Payload
| Style | Contents | Trade-off |
|---|---|---|
| Notification only | Identifier and event type | Consumers must call back for detail |
| Event-carried state | Full relevant data | Larger messages, schema coupling |
| Delta | What changed | Consumer needs prior state |
| Full snapshot | Entire entity after the change | Large, but order-tolerant |
{
"messageId": "48120",
"eventType": "OrderPlaced",
"aggregateType": "order",
"aggregateId": "order-7734",
"occurredAt": "2026-09-24T06:41:09Z",
"schemaVersion": 2,
"payload": {
"orderId": "order-7734",
"customerId": "customer-4821",
"total": 249900,
"currency": "INR",
"items": [
{ "sku": "sku-501", "quantity": 2, "unitPrice": 124950 }
]
},
"traceId": "trace-9a2d5f"
}
Failure Modes
| Failure | Consequence | Mitigation |
|---|---|---|
| Relay stops running | Events accumulate undelivered | Alert on oldest unpublished age |
| Broker unavailable | Publishing halts | Retry, rows remain safe |
| Poison message | Partition blocked indefinitely | Attempt cap, move to dead letter |
| Crash after publish | Message republished | Consumer-side deduplication |
| Cleanup not running | Table and index growth | Monitor row count and query latency |
| Concurrent relays | Duplicate publishing, reordering | Leader election or sharding |
| Schema change breaks consumers | Downstream processing fails | Version the payload, evolve compatibly |
-- Poison message handling: bound the attempts
UPDATE outbox
SET attempts = attempts + 1,
last_error = :error_message
WHERE id = :id;
-- Route persistently failing rows aside
UPDATE outbox
SET published_at = CURRENT_TIMESTAMP,
dead_lettered = TRUE
WHERE id = :id
AND attempts >= :max_attempts;
Where the Outbox Fits
| Use Case | Role of the Outbox |
|---|---|
| Saga step transitions | Guarantees the next step is triggered |
| Search index updates | Ensures no document change is missed |
| Cache invalidation | Reliable invalidation on every write |
| Audit event emission | Audit record shares the write's fate |
| Cross-service notification | No lost domain events |
| Analytics streaming | Complete event stream |
| Webhook delivery | Durable queue for external callbacks |
Costs and Alternatives
What You Gain
- No lost or phantom events
- No distributed transaction
- Broker outage does not lose data
- Natural retry and replay
- Works with any database
What It Costs
- Extra write per transaction
- A relay process to operate
- Publish latency added
- Cleanup job required
- Duplicates pushed to consumers
| Alternative | Applicability | Limitation |
|---|---|---|
| Event sourcing | Events are the state | Significant architectural change |
| Direct CDC on business tables | Consumers accept row changes | Couples consumers to schema |
| Two-phase commit | Broker supports it | Blocking and slow |
| Periodic reconciliation | Divergence is tolerable | Detects late, cannot replay intent |
| Accept the risk | Events are non-critical | Silent data loss |
Monitoring
Signals Worth Tracking
- Age of the oldest unpublished row
- Unpublished row count
- Publish rate versus insert rate
- Relay poll duration and batch size
- Publish failures by error type
- Rows exceeding the attempt threshold
- Dead-lettered message count
- Total table size and cleanup lag
- Poll query latency over time
- Duplicate deliveries detected by the inbox
Verification
describe("transactional outbox", function () {
it("writes no outbox row when the business write fails", async function () {
await expect(
placeOrder(database, invalidOrder)
).rejects.toThrow();
const pending = await database.outbox.unpublished();
expect(pending).toHaveLength(0);
});
it("retains the event when the broker is unavailable", async function () {
await placeOrder(database, validOrder);
await faultInjector.stopBroker();
await relay.pollOnce();
const pending = await database.outbox.unpublished();
expect(pending).toHaveLength(1);
});
it("republishes after a crash between publish and mark", async function () {
await placeOrder(database, validOrder);
await faultInjector.crashAfterPublish();
await relay.restart();
await relay.pollOnce();
expect(broker.messagesFor(validOrder.orderId)).toHaveLength(2);
});
it("processes a duplicate exactly once at the consumer", async function () {
const message = buildMessage(validOrder);
await consumeMessage(database, message, handler);
await consumeMessage(database, message, handler);
expect(handler.invocationCount).toBe(1);
});
it("preserves order per aggregate", async function () {
await updateOrder(database, orderId, { status: "confirmed" });
await updateOrder(database, orderId, { status: "shipped" });
await relay.drain();
const sequence = broker.messagesFor(orderId).map(m => m.status);
expect(sequence).toEqual(["confirmed", "shipped"]);
});
});
Common Design Mistakes
Weak Design
- Publishing directly after committing
- Writing the outbox row outside the transaction
- Building the payload at publish time
- Skipping failed messages to continue
- Running concurrent relays without sharding
- No attempt cap on failing rows
- No cleanup job
- Assuming consumers get exactly-once
- Not alerting on unpublished age
Strong Design
- One transaction covering both writes
- Payload captured at write time
- Stops the partition on failure
- Shards or elects a single relay
- Caps attempts and dead-letters
- Deletes published rows in batches
- Partial index on unpublished only
- Pairs with an inbox for deduplication
- Alerts on oldest unpublished age
System Design Interview Discussion
| Question | What Your Answer Should Cover |
|---|---|
| What is the dual-write problem? | Two systems, no shared transaction |
| Why not just retry the publish? | The process can die before retrying |
| How does the outbox solve it? | Both writes in one local transaction |
| Polling or CDC? | Latency and load versus operational cost |
| How is ordering preserved? | Sequential per aggregate, partitioned keys |
| Can you scale the relay? | Shard by aggregate hash |
| Does this give exactly-once? | At-least-once plus consumer deduplication |
| How do you detect failure? | Oldest unpublished age as the key signal |
Design Checklist
Production Checklist
- Write the outbox row in the business transaction
- Capture the full payload at write time
- Include a version field in every payload
- Add a partial index on unpublished rows
- Order the relay query by a monotonic identifier
- Partition messages by aggregate identifier
- Stop the partition on publish failure
- Shard or elect a leader across relay instances
- Cap attempts and dead-letter poison messages
- Delete published rows in bounded batches
- Retain rows long enough for replay
- Implement an inbox on the consumer side
- Write the inbox record in the consumer transaction
- Alert on oldest unpublished row age
- Monitor publish rate against insert rate
- Test crash between publish and mark
Knowledge Check
What is the dual-write problem?
Writing to a database and publishing to a broker are separate operations with no shared transaction, so one can succeed while the other fails.
Why does reordering the two writes not help?
It swaps a lost event for a phantom event. Publishing first can announce a change that never persisted.
Why must the payload be captured at write time?
Building it at publish time reads current state, which may have changed since the transaction committed, producing an event that misrepresents what happened.
Does the outbox provide exactly-once delivery?
No. It guarantees at-least-once, because the relay may crash after publishing but before marking the row. Consumers must deduplicate.
Why is oldest unpublished age the key metric?
A rising value means events are being recorded but never delivered, which is a silent failure invisible to the application writing them.
Summary
The dual-write problem arises whenever a service must persist a change and notify others about it. These are two independent systems with no shared transaction, so a crash between them produces either a lost event or a phantom one.
No ordering of the two writes avoids this, and retries do not help because the process itself can die. The transactional outbox resolves it by making both writes the same write: the message is inserted into a table in the same local transaction as the business change.
A separate relay then delivers those rows, either by polling or by tailing the database log. Because the rows are durable, delivery can be retried indefinitely. Ordering is preserved by reading sequentially and partitioning by aggregate identifier, and multiple relays can run if sharded by that same key.
The guarantee is at-least-once, not exactly-once, because the relay may publish and crash before recording success. The inbox pattern completes the picture, recording processed message identifiers in the consumer's own transaction so duplicates become harmless.
Key Takeaway
If two writes must be atomic, make them one write. Insert the event into an outbox table inside the business transaction, capture the payload at write time, deliver it with a relay that stops rather than skips on failure, pair it with an inbox on the consumer side, and alert on the age of the oldest unpublished row.