Hello! Welcome back to your course on designing high-load distributed systems.
In this lesson, we will tackle a fundamental pattern for building reliable event-driven systems: the Transactional Outbox. Your experience with payment platforms and high-throughput services at Yandex has likely exposed you to the critical need for data consistency between a service's internal state and the events it communicates to the rest of the system. This lesson moves beyond the theory and dives into the practical implementation of ensuring that consistency.
Our learning outcome is to implement the Transactional Outbox pattern using a database table and a message relay. We will cover:
- The core problem the pattern solves (the "dual-write" problem).
- How to design and use a database outbox table within an atomic transaction.
- Different strategies for implementing the "message relay" component that publishes events.
- Key considerations for scalability and resilience in a high-load environment.
Let's begin.
1. The Challenge: Atomicity in Distributed Systems
In a microservices architecture, a common task is to update a database and then publish an event to a message broker (like Kafka or RabbitMQ) to notify other services of the change. For example, in a payment system, you might:
- Update the
invoicestable to mark an invoice asPAID. - Publish an
InvoicePaidevent to a Kafka topic.
The challenge lies in the "and". These two operations—a database write and a network call to a message broker—cannot be part of a single, atomic transaction. This creates a "dual-write" problem. Consider the failure modes:
- Database commit succeeds, but the message publish fails: The invoice is marked as paid, but no one is notified. The system is in an inconsistent state.
- Message publish succeeds, but the database commit fails (or happens after): Other services are notified that the invoice is paid, but the source of truth (the invoice database) doesn't reflect this. This can lead to incorrect downstream actions.
The Transactional Outbox pattern solves this by ensuring that the intent to publish an event is saved atomically with the business data change.
This diagram provides a clear visual of the pattern's components and flow.

2. The Solution: The Outbox Table
The core idea is to leverage the atomicity of a local database transaction. Instead of directly publishing an event, the service writes the event to a special outbox table within the same database and within the same transaction as the business data update.
If the transaction succeeds, both the business data and the event are durably stored. If it fails, both are rolled back. This guarantees that the intent to publish is never lost and is perfectly consistent with the service's state.
To understand the problem and the high-level solution, please start by reading the introductory sections of the following AWS guide.
Transactional outbox pattern - AWS Prescriptive Guidance
This AWS Prescriptive Guidance document clearly defines the dual-write problem and introduces the Transactional Outbox pattern as a solution. It sets the stage for why this pattern is critical for data consistency.
Please read the 'Intent', 'Motivation', and 'Applicability' sections. Focus on how it frames the dual-write problem and the scenarios where this pattern is most useful.
Designing the Outbox Table
A well-designed outbox table is crucial. It needs to store all the necessary information for the event to be created and published later.
Let's look at the best practices for the table schema and transactional integrity.
Outbox Pattern Best Practices for Reliable Messaging in ...
The article 'Outbox Pattern Best Practices' by pnrkarga provides an excellent, practical guide to implementation. We'll start with the sections on table design and transactional integrity.
Please read the sections 'What is the Outbox Pattern?', '1. Transactional Integrity', and '3. Optimize the Outbox Table Design'. Pay close attention to the recommended columns for the outbox table.
As you've read, a typical outbox table includes:
event_id: A unique identifier for the event (e.g., a UUID).aggregate_id: The ID of the business entity that was changed (e.g., theorder_id).event_type: A string identifying the event type (e.g., 'OrderCreated').payload: The body of the event, usually as a JSON or Avro byte array.status: To track the publishing state (e.g.,PENDING,PUBLISHED).created_at: Timestamp for ordering and diagnostics.
The application code would then look conceptually like this (in pseudocode):
// Assumes a framework that manages the transaction boundary
@Transactional
public void createOrder(OrderData data) {
// 1. Create and save the business entity
Order order = new Order(data);
orderRepository.save(order);
// 2. Create the event payload
OrderCreatedEvent event = new OrderCreatedEvent(order.getId(), order.getDetails());
// 3. Save the event to the outbox table
OutboxEvent outboxEvent = new OutboxEvent(
order.getId(),
"OrderCreated",
toJson(event)
);
outboxRepository.save(outboxEvent);
// The transaction commits here, making both writes atomic
}
3. The Message Relay: Publishing the Events
With events safely stored in the outbox table, we need a reliable mechanism to read them and publish them to the message broker. This component is the message relay. There are two primary approaches.
Approach 1: The Polling Publisher
This is the most straightforward implementation. A separate process or background thread periodically polls the outbox table for unpublished events, sends them to the message broker, and then marks them as published.

The simple flow for the relay is:
- Query the outbox table for events with
status = 'PENDING'. - For each event, publish it to the message broker.
- If publishing is successful, update the event's row in the outbox table to
status = 'PUBLISHED'or delete it.
This process must be resilient. What if the relay crashes after publishing but before updating the row? It will re-publish the same event on restart. This is why consumers of the events must be idempotent. The outbox pattern provides at-least-once delivery semantics.
Scaling the Polling Relay
In a high-load system, a single-threaded poller can become a bottleneck. To scale out, you can run multiple instances of the relay. However, this introduces a new challenge: how do you prevent multiple instances from picking up and publishing the same event?
The solution is to use a database-level locking mechanism to have each relay instance "claim" a batch of events. PostgreSQL's FOR UPDATE SKIP LOCKED is perfectly suited for this.
BEGIN;
SELECT * FROM outbox
WHERE status = 'PENDING'
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED;
-- The selected rows are now locked. Other concurrent transactions will skip them.
-- The relay process now publishes these 100 events.
-- Once published, update or delete the rows within the same transaction.
UPDATE outbox SET status = 'PUBLISHED' WHERE id IN (...);
COMMIT;
This strategy allows multiple relay instances to work on the outbox table in parallel without contention, dramatically increasing throughput.
Let's dive deeper into the implementation details of the relay and scaling strategies.
Outbox Pattern Best Practices for Reliable Messaging in ...
This next reading from the iyzico.engineering article covers the relay's responsibilities, including publishing, cleanup, and the critical scaling techniques we just discussed.
Please read sections '4. Reliable Event Publishing', '5. Event Acknowledgment and Cleanup', and '6. Scaling'. Focus on the retry logic, the importance of cleanup, and the use of 'SKIP LOCKED' for parallel processing.
Approach 2: Change Data Capture (CDC)
An alternative to polling is to use Change Data Capture (CDC). In this approach, the message relay tails the database's transaction log (e.g., PostgreSQL's WAL). When it sees a new row inserted into the outbox table, it reads the row and publishes the corresponding event.
Tools like Debezium are designed for this. They act as Kafka Connect sources that can stream table changes directly into Kafka topics.
Transactional outbox pattern - AWS Prescriptive Guidance
The AWS guide provides a good overview of the CDC approach as an alternative to polling. This will give you a broader perspective on implementation choices.
Read the subsections 'Using an outbox table with a relational database' and 'Using change data capture (CDC)'. Compare the two architectures shown in the diagrams.
Polling vs. CDC: Trade-offs
| Aspect | Polling Publisher | Change Data Capture (CDC) |
|---|---|---|
| Latency | Higher. Events are published with a delay based on the polling interval. | Near real-time. Changes are streamed from the transaction log. |
| DB Load | Higher. Involves frequent queries against the outbox table. | Lower. Taps into the database's native replication stream. |
| Complexity | Simpler to implement and manage within the application stack. | More complex. Requires deploying and managing a separate CDC platform (e.g., Debezium, Kafka Connect). |
| Coupling | Loosely coupled. The relay only needs SQL access. | Tightly coupled to the database version and its transaction log format. |
For many systems, the simplicity and robustness of a well-implemented polling publisher (especially with SKIP LOCKED) is sufficient. The CDC approach is powerful for use cases requiring very low latency but comes with significant operational overhead.
4. Code Examples and Further Study
To see how these concepts translate into code, both articles provide links to complete sample projects. Exploring these will solidify your understanding.
- The iyzico.engineering article (
3b6db) links to a Spring Boot, PostgreSQL, and RabbitMQ example, which is a very common stack. - The AWS guide (
73edd) links to a repository demonstrating both the polling and CDC approaches using AWS services (Lambda, RDS, DynamoDB).
I highly recommend cloning and running these examples to see the pattern in action.
Conclusion
In this lesson, we have dissected the Transactional Outbox pattern as a robust solution to the dual-write problem in distributed systems.
Key Takeaways:
- The pattern guarantees atomicity between a business state change and the publication of a corresponding event by using a single database transaction.
- It involves two main components: writing to an
outboxtable transactionally with the business data, and a separate message relay process. - The message relay is responsible for reading from the outbox and publishing to a message broker, providing at-least-once delivery guarantees.
- The relay can be implemented using a polling publisher or a Change Data Capture (CDC) pipeline, each with distinct trade-offs in latency, complexity, and database load.
- For high-throughput systems, a polling publisher can be scaled horizontally by using database locking mechanisms like
FOR UPDATE SKIP LOCKED.
Preview of the Next Lesson:
The Transactional Outbox pattern is a foundational element for building event-driven communication. In our next lesson, we will build upon this by exploring Sagas, a pattern for managing long-running, distributed transactions. We will compare choreography-based Sagas (which rely heavily on reliable eventing like the outbox pattern) with orchestration-based Sagas.