⏱️ Reading time: 13 min
Saving an order to the database and publishing the corresponding event to Kafka seem like two trivial steps, but they’re never atomic: if the process dies between one and the other, the order exists without the rest of the system finding out. The transactional outbox pattern solves this problem without distributed transactions or locks between different systems.
📑 En este artículo
This problem is known as dual write: two writes that should happen together, in two systems that don’t share a transaction. It shows up between a database and a message queue, between two different databases, or between a database and a call to an external API.
TL;DR
- The transactional outbox pattern saves the event in the same transaction as the business data, in an outbox table.
- SELECT … FOR UPDATE SKIP LOCKED prevents two relays from processing the same outbox event twice.
- Debezium reads PostgreSQL’s WAL and publishes every new row in the outbox table to Kafka without touching it with SELECTs.
- Consumers must be idempotent because the outbox guarantees at-least-once delivery, never exactly-once.
- An outbox table that’s never purged grows without limit: you need a job that deletes or archives events that have already been published.
What Is the Transactional Outbox Pattern?
The transactional outbox pattern is a design pattern that guarantees a database write and the publication of an associated event happen atomically: it saves the event as just another row, within the same transaction, and delegates its actual delivery to a separate process called a relay.
It emerged as a response to a concrete problem in microservices architectures: when a service needs to persist a change and also notify other services about it, doing so in two separate steps leaves a window where the system can end up in a half-finished state. Chris Richardson documented the pattern in his microservices patterns catalog, and Gunnar Morling, from the Debezium team, popularized combining it with change data capture in a 2019 reference article.
Why It Matters
A service that updates its database and then publishes a message has two writes that should succeed or fail together, but they run in two systems that don’t share a transaction. If the process dies right after the database commit and before publishing, the event is lost forever: the order exists, but nobody finds out. If the order is reversed and the event is published before confirming the database write, an event can arrive about an order that later fails and never gets saved.
The naive solution is to wrap both operations in a two-phase distributed transaction (2PC). It works, but it requires the message broker to support XA, adds locks between systems, and doesn’t scale well. Kafka, for example, doesn’t offer native support for XA transactions with an external database. The key to the transactional outbox pattern is avoiding the problem entirely. Instead of writing to two systems, you write to just one.
flowchart TD
A["Service processes an order"] --> B["Saves the order to the database"]
A --> C["Publishes the event to Kafka"]
B --> D["Confirmed: the order is saved"]
C --> E["Network failure: the event is lost"]
How the Outbox Works Under the Hood
The outbox table lives in the same database as the business tables. When the service needs to persist a change and emit an event, it performs both writes in the same transaction: an INSERT into the business table (for example, orders) and an INSERT into the outbox table with the event payload. If the transaction fails, neither write gets recorded; if it commits, both are recorded together, always.
A separate process, called a relay or publisher, periodically reads new rows from the outbox table and sends them to the actual message queue (Kafka, RabbitMQ, SQS). Once it confirms delivery, it marks the row as published. This relay can be implemented in two ways: by polling, with a SELECT at set intervals, or by CDC, reading the database’s transaction log without running queries against the table.
sequenceDiagram
participant S as Service
participant DB as Database
participant R as Relay
participant K as Kafka
S->>DB: BEGIN transaction
S->>DB: INSERT into orders
S->>DB: INSERT into outbox
S->>DB: COMMIT
R->>DB: SELECT pending events
DB-->>R: outbox table rows
R->>K: publishes the event
R->>DB: UPDATE published_at
Practical Examples and How to Get Started
The following example implements an outbox with Node.js and PostgreSQL, without relying on Kafka or additional infrastructure: first a simple polling relay, then the improvement with SKIP LOCKED so multiple relay instances can run at the same time.
Requirements
Node.js 18 or later, a PostgreSQL instance running locally or in Docker, and the pg library. Kafka isn’t needed for this example: the relay can simulate publishing with a log before connecting it to a real broker.
npm init -y
npm install pg
The schema has two tables: the business table orders and the outbox table, which stores each event pending publication.
CREATE TABLE orders (
id SERIAL PRIMARY KEY,
customer TEXT NOT NULL,
total NUMERIC NOT NULL
);
CREATE TABLE outbox (
id SERIAL PRIMARY KEY,
aggregate_type TEXT NOT NULL,
aggregate_id INTEGER NOT NULL,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
published_at TIMESTAMPTZ
);
The service that creates an order inserts into both tables within the same transaction. If the COMMIT fails, neither row gets saved.
const { Pool } = require('pg');
const pool = new Pool();
async function crearPedido(customer, total) {
const client = await pool.connect();
try {
await client.query('BEGIN');
const { rows } = await client.query(
'INSERT INTO orders (customer, total) VALUES ($1, $2) RETURNING id',
[customer, total]
);
const orderId = rows[0].id;
await client.query(
`INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
VALUES ($1, $2, $3, $4)`,
['order', orderId, 'OrderCreated', JSON.stringify({ orderId, customer, total })]
);
await client.query('COMMIT');
return orderId;
} catch (err) {
await client.query('ROLLBACK');
throw err;
} finally {
client.release();
}
}
Running crearPedido('Ana', 49.90), the orders table gains one row and the outbox table gains another, atomically. A verification query confirms both exist:
SELECT o.id, o.customer, x.event_type, x.published_at
FROM orders o JOIN outbox x ON x.aggregate_id = o.id
WHERE o.id = 1;
id | customer | event_type | published_at
----+----------+---------------+--------------
1 | Ana | OrderCreated | NULL
published_at being NULL means the event hasn’t gone out yet. That’s where the relay comes in: a process that polls the outbox table, publishes what it finds, and marks the row.
async function relayPolling() {
const { rows } = await pool.query(
`SELECT * FROM outbox WHERE published_at IS NULL ORDER BY id LIMIT 50
FOR UPDATE SKIP LOCKED`
);
for (const evento of rows) {
await publicarEnKafka(evento.event_type, evento.payload);
await pool.query('UPDATE outbox SET published_at = now() WHERE id = $1', [evento.id]);
}
}
setInterval(relayPolling, 2000);
FOR UPDATE SKIP LOCKED is the key piece if you run more than one relay instance: each instance only picks up rows that no one else is processing at that moment, instead of blocking while waiting for the other. Without that clause, two relays competing for the same row end up either publishing the same event twice or locking each other out.
💡 Tip: running multiple relay instances tolerates failures, it doesn’t multiply throughput. With SKIP LOCKED, each event is processed by only one instance, so adding relays increases availability, not publishing speed.
On the other side, the consumer needs its own deduplication table, sometimes called an inbox, to ignore an event it already processed before.
CREATE TABLE processed_events (
event_id INTEGER PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
async function manejarEventoOrderCreated(evento) {
const client = await pool.connect();
try {
await client.query('BEGIN');
const { rowCount } = await client.query(
'INSERT INTO processed_events (event_id) VALUES ($1) ON CONFLICT DO NOTHING',
[evento.id]
);
if (rowCount === 0) {
await client.query('ROLLBACK');
return; // already processed before, ignore it
}
await client.query(
'UPDATE inventory SET reserved = reserved + 1 WHERE product_id = $1',
[evento.payload.productId]
);
await client.query('COMMIT');
} catch (err) {
await client.query('ROLLBACK');
throw err;
} finally {
client.release();
}
}
Real-World Use Cases
The most common case is syncing a state change with a notification to other services: an order that switches to paid and triggers sending an email, updating inventory, and calculating commissions, all in different services that don’t share a database.
Another frequent use is feeding a search system or a derived cache: every write to the main table generates an outbox event that a consumer uses to reindex in Elasticsearch or invalidate an entry in Redis, without coupling the main transaction to the availability of those systems.
It also shows up as a foundation for partial event sourcing: instead of rebuilding the entire state from events, a service keeps its normal state in relational tables and uses the outbox only to expose changes as an event stream to the rest of the organization.
Common Mistakes and Best Practices
- Publishing before the transaction commits. If the relay were to read an event whose transaction hasn’t committed yet, a consumer could process an order that’s later canceled with a rollback. You should always read only rows that have already been committed.
- An unpurged outbox table. If published rows are never deleted or archived, the table grows without limit and ends up degrading the indexes across the whole database. A periodic job that deletes rows with
published_atolder than a few days is mandatory in production. - Non-idempotent consumers. The outbox provides at-least-once guarantees: the same event can arrive duplicated if the relay fails right after publishing and before marking the row. The consumer needs a deduplication key, and the outbox row’s own id works fine.
- Event ordering across different aggregates. Publishing without a partition key doesn’t preserve the order within the same order. You need to use
aggregate_idas the key when sending to Kafka:
await kafkaProducer.send({
topic: 'orders',
messages: [{ key: String(evento.aggregate_id), value: JSON.stringify(evento.payload) }]
});
Comparison with Alternatives
| Option | When to Use It | Advantage | Limitation |
|---|---|---|---|
| Outbox with polling | Low to medium volume, no CDC infrastructure available | Easy to implement with what you already have | Latency equal to the polling interval and extra SELECT load |
| Outbox with CDC (Debezium) | High volume, low latency, you already run Kafka | Minimal latency, reads the database log without touching it | Requires deploying and maintaining Kafka Connect |
| Distributed transactions (2PC/XA) | Legacy systems with XA support on both ends | True atomicity between the two resources | Blocking, doesn’t scale, and most modern brokers don’t support it |
Going Deeper
The CDC implementation replaces polling with direct reads of the transaction log, the WAL in PostgreSQL. Debezium runs as a Kafka Connect connector, subscribes to PostgreSQL’s logical replication, and emits a Kafka event for every new row in the outbox table, without running a single SELECT against it. The advantage isn’t just about latency. A polling relay competes for locks with the application’s normal transactions, while reading the WAL doesn’t interfere with anything.
Neither variant delivers true exactly-once. Both polling and Debezium can publish the same event more than once if the process fails between publishing and confirming the offset. The guarantee they do hold is at-least-once with order preserved per partition, which pushes the responsibility for deduplication onto the consumer.
An outbox event’s lifecycle moves through well-defined states, and modeling it this way helps decide what to retry:
stateDiagram-v2
[*] --> Pending
Pending --> Published: the relay sends it to Kafka
Published --> Confirmed: the consumer processes it
Published --> Pending: network failure, retried
Confirmed --> [*]
Monitoring Outbox Lag
The simplest health signal is measuring how long pending events have gone unpublished. If that number keeps growing steadily, the relay has fallen behind.
SELECT extract(epoch FROM now() - min(created_at)) AS lag_seconds
FROM outbox
WHERE published_at IS NULL;
lag_seconds
-------------
3.42
⚠️ Watch out: if several services write to the same outbox table with different server clocks, ordering only by created_at can process near-simultaneous events in a different order than they actually happened. The table’s auto-incrementing id is a more reliable ordering reference than the timestamp.
A detail that’s rarely documented: the outbox doesn’t need a new table for every aggregate type. A generic schema with aggregate_type, aggregate_id, event_type, and a JSONB payload works for the whole application, and it’s the same schema the Debezium outbox connector expects without any further transformation.
Your next step: set up the two tables from this article in a local PostgreSQL database, run crearPedido three times in a row, and confirm with the verification SELECT that each order has its outbox row before you write the relay.
Frequently Asked Questions
What problem does the transactional outbox pattern solve?
It solves the dual write problem: it prevents a database write and the publication of an event to a message queue from ending up inconsistent when the process fails between the two.
Does the transactional outbox guarantee exactly-once delivery?
No. It guarantees at-least-once delivery: an event can arrive duplicated at the consumer, which needs to be idempotent to process it only once despite retries.
Do you need Kafka to use an outbox?
Not necessarily. The outbox works with any message queue, such as RabbitMQ or SQS, or even with a relay that calls an external API directly. Kafka is just the most common choice in microservices architectures.
What’s the difference between polling-based outbox and Debezium-based outbox?
Polling runs a periodic SELECT against the outbox table. Debezium reads the transaction log and doesn’t touch the table with queries, which reduces latency and load on the database.
How do you prevent the outbox table from growing without limit?
With a periodic job that deletes or archives rows where published_at is not NULL after a reasonable amount of time, splitting off the history if it needs to be kept for auditing.
References
- microservices.io: Chris Richardson’s patterns catalog with the original definition of the transactional outbox.
- Debezium Blog: Gunnar Morling’s article on how to combine the outbox with CDC.
- PostgreSQL Docs: official reference for the
SELECT ... FOR UPDATE SKIP LOCKEDclause. - Wikipedia: general definition of change data capture (CDC).
📱 Enjoy this content? Follow @programacion on Telegram for daily tech content in Spanish: quick summaries, fresh content every day. @programacion
Featured image: Foto de Markus Spiske en Unsplash
Did it work for you? Got a different error? Say so below: questions get answered and help the next reader.
Leave a comment
0 Comments