Transactional Outbox — это архитектурный паттерн, который позволяет надёжно выполнить две связанные операции:
- изменить данные в локальной базе сервиса;
- отправить событие или сообщение в Kafka, RabbitMQ или другой брокер.
Главная проблема в том, что база данных и Kafka являются двумя независимыми системами, и обычная локальная транзакция не может атомарно охватить их обе.
Transactional Outbox устраняет этот двойной разрыв: сервис не отправляет событие непосредственно в Kafka, а записывает его в специальную таблицу outbox в той же транзакции, в которой изменяет бизнес-данные. Затем отдельный процесс переносит записанное событие из outbox в брокер. (Debezium)
Проблема Dual Write
Представим сервис заказов:
Order Service
├── PostgreSQL
└── Kafka
При создании заказа нужно:
1. INSERT INTO orders ...
2. kafka.send("OrderCreated")
Это две записи в разные системы:
PostgreSQL ← первая запись
Kafka ← вторая запись
Такую ситуацию называют Dual Write, то есть двойной записью.
Вариант 1: сначала база, потом Kafka
orderRepository.save(order);
kafkaTemplate.send("orders", event);
Возможен сценарий:
1. Заказ записан в базу.
2. Приложение упало.
3. Событие в Kafka не отправлено.
Получается:
orders:
заказ существует
Kafka:
OrderCreated отсутствует
Другие сервисы никогда не узнают о заказе.
Вариант 2: сначала Kafka, потом база
kafkaTemplate.send("orders", event);
orderRepository.save(order);
Теперь возможен обратный сценарий:
1. OrderCreated отправлен в Kafka.
2. Запись в базу не удалась.
Получается:
Kafka:
OrderCreated существует
orders:
заказа не существует
Другие сервисы начинают обрабатывать событие о сущности, которая фактически не была создана.
Как работает Transactional Outbox
В базе сервиса создаётся дополнительная таблица:
CREATE TABLE outbox (
id UUID PRIMARY KEY,
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id UUID NOT NULL,
event_type VARCHAR(100) NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMP NOT NULL,
published_at TIMESTAMP NULL
);
При создании заказа сервис выполняет одну локальную транзакцию:
BEGIN;
INSERT INTO orders (
id,
customer_id,
amount,
status
)
VALUES (
'order-123',
'customer-456',
1000,
'CREATED'
);
INSERT INTO outbox (
id,
aggregate_type,
aggregate_id,
event_type,
payload,
created_at
)
VALUES (
'event-789',
'Order',
'order-123',
'OrderCreated',
'{
"orderId": "order-123",
"customerId": "customer-456",
"amount": 1000
}',
NOW()
);
COMMIT;
База данных гарантирует:
либо записаны и заказ, и событие;
либо не записано ничего.
Вот здесь и находится ядро паттерна.
Не происходит попытки сделать одну транзакцию между PostgreSQL и Kafka. Вместо этого обе критические записи сначала сводятся к одной транзакционной системе, то есть к локальной базе данных.
Полный поток
Client
│
│ POST /orders
▼
Order Service
│
│ одна DB-транзакция
▼
PostgreSQL
├── orders
│ order-123
│
└── outbox
OrderCreated(order-123)
│
│ Relay / CDC
▼
Kafka
│
▼
Payment Service
Notification Service
Analytics Service
Последовательность:
1. Клиент создаёт заказ.
2. Order Service открывает транзакцию.
3. Order Service записывает заказ.
4. Order Service записывает событие в outbox.
5. Транзакция фиксируется.
6. Отдельный publisher обнаруживает запись outbox.
7. Publisher отправляет событие в Kafka.
8. Получатели обрабатывают событие.
Таким образом, событие становится частью состояния самого сервиса, а не кратковременной попыткой вызвать внешнюю систему.
Кто отправляет записи из Outbox
Есть два основных варианта.
1. Polling Publisher
Отдельный процесс периодически читает таблицу:
SELECT *
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED;
После этого он:
1. отправляет события в Kafka;
2. помечает их как опубликованные;
3. позднее удаляет или архивирует.
Пример:
@Scheduled(fixedDelay = 1000)
public void publishEvents() {
List<OutboxEvent> events =
outboxRepository.findUnpublishedEvents();
for (OutboxEvent event : events) {
kafkaTemplate.send(
event.getAggregateId(),
event.getPayload()
);
outboxRepository.markPublished(event.getId());
}
}
Схема:
Outbox table
│
│ SELECT every second
▼
Outbox Publisher
│
▼
Kafka
Это относительно простой вариант, но сервис постоянно опрашивает таблицу.
2. Change Data Capture
Вместо опроса можно читать журнал изменений базы данных.
Например:
PostgreSQL WAL
│
▼
Debezium
│
▼
Kafka Connect
│
▼
Kafka
Приложение только записывает данные:
orders + outbox
Debezium отслеживает изменения таблицы outbox через журнал транзакций базы и преобразует их в сообщения Kafka. Debezium прямо предоставляет Outbox Event Router для такой реализации паттерна. (Debezium)
Здесь важно различить:
Transactional Outbox
это сам архитектурный паттерн;
Debezium
это один из инструментов его реализации.
Можно реализовать Outbox без Debezium, а Debezium можно использовать не только для Outbox.
Что именно гарантирует Transactional Outbox
Паттерн гарантирует согласованность между:
бизнес-изменением
и
намерением опубликовать событие
То есть если заказ успешно создан, запись о необходимости отправить OrderCreated также существует.
Но паттерн сам по себе обычно не гарантирует:
ровно одну физическую доставку сообщения
Чаще всего получается семантика:
at-least-once delivery
Сообщение будет доставлено как минимум один раз, но иногда может быть отправлено повторно.
Откуда берутся дубликаты
Допустим, publisher действует так:
1. прочитал outbox event;
2. отправил его в Kafka;
3. должен отметить published_at;
Но между шагами 2 и 3 процесс упал:
Kafka:
сообщение уже существует
Outbox:
published_at всё ещё NULL
После перезапуска publisher снова увидит событие и снова отправит его.
OrderCreated(event-789)
OrderCreated(event-789)
Поэтому получатель должен быть идемпотентным.
Идемпотентный consumer
Каждое событие должно иметь уникальный eventId:
{
"eventId": "event-789",
"eventType": "OrderCreated",
"aggregateId": "order-123",
"payload": {
"amount": 1000
}
}
Получатель хранит обработанные идентификаторы:
CREATE TABLE processed_messages (
consumer_name VARCHAR(100) NOT NULL,
event_id UUID NOT NULL,
processed_at TIMESTAMP NOT NULL,
PRIMARY KEY (consumer_name, event_id)
);
При получении сообщения:
BEGIN;
INSERT INTO processed_messages (
consumer_name,
event_id,
processed_at
)
VALUES (
'payment-service',
'event-789',
NOW()
)
ON CONFLICT DO NOTHING;
Если запись уже существовала, событие ранее обрабатывалось:
event-789 уже обработан → ничего не делать
Если записи не было:
event-789 новый → выполнить бизнес-операцию
Важно, чтобы запись в processed_messages и бизнес-изменение consumer тоже выполнялись в одной локальной транзакции:
BEGIN;
INSERT INTO processed_messages ...;
INSERT INTO payments ...;
COMMIT;
Таким образом:
Outbox на стороне producer
+
Inbox / processed_messages на стороне consumer
создают надёжный асинхронный канал поверх доставки с возможными повторами.
Пример на Spring
Упрощённо операция выглядит так:
@Service
@RequiredArgsConstructor
public class OrderService {
private final OrderRepository orderRepository;
private final OutboxRepository outboxRepository;
private final ObjectMapper objectMapper;
@Transactional
public UUID createOrder(CreateOrderCommand command) {
Order order = Order.create(
command.customerId(),
command.amount()
);
orderRepository.save(order);
OrderCreated event = new OrderCreated(
UUID.randomUUID(),
order.getId(),
order.getCustomerId(),
order.getAmount(),
Instant.now()
);
OutboxMessage outboxMessage = new OutboxMessage(
event.eventId(),
"Order",
order.getId(),
"OrderCreated",
serialize(event),
Instant.now()
);
outboxRepository.save(outboxMessage);
return order.getId();
}
private String serialize(Object event) {
try {
return objectMapper.writeValueAsString(event);
} catch (JsonProcessingException exception) {
throw new IllegalStateException(
"Cannot serialize outbox event",
exception
);
}
}
}
Аннотация:
@Transactional
здесь охватывает обе SQL-операции:
INSERT INTO orders
INSERT INTO outbox
Она не охватывает Kafka и не должна её охватывать.
Почему нельзя просто отправить сообщение после commit
Иногда делают так:
@Transactional
public void createOrder() {
orderRepository.save(order);
}
public void handleRequest() {
createOrder();
kafkaTemplate.send(event);
}
На первый взгляд всё нормально: сначала гарантированно фиксируем заказ, затем отправляем событие.
Но между этими действиями всё равно остаётся щель:
commit базы
↓
щель падения
↓
send в Kafka
Transactional Outbox материализует намерение отправить событие до выхода из транзакции:
commit:
order
outbox event
После этого событие уже не зависит от продолжительности жизни текущего процесса.
Transactional Outbox и Saga
Эти паттерны решают разные задачи.
Transactional Outbox
Решает вопрос:
Как надёжно опубликовать событие, связанное с локальной транзакцией?
локальное изменение
+
будущее сообщение
Saga
Решает вопрос:
Как провести бизнес-операцию через несколько сервисов, если общей транзакции между ними нет?
Например:
Create Order
↓
Reserve Money
↓
Reserve Inventory
↓
Arrange Delivery
Transactional Outbox часто используется внутри каждого шага Saga:
Order Service:
создать Order
записать OrderCreated в outbox
Payment Service:
зарезервировать деньги
записать PaymentReserved в outbox
Inventory Service:
зарезервировать товар
записать InventoryReserved в outbox
То есть:
Saga управляет последовательностью распределённых изменений.
Outbox обеспечивает надёжную передачу событий между шагами.
Transactional Outbox и двухфазный commit
Альтернативой теоретически может быть распределённая транзакция:
Database
+
Message Broker
+
2PC coordinator
Но тогда обе системы должны поддерживать совместимый транзакционный протокол, появляется координатор и возрастает операционная связанность.
Outbox не пытается сделать базу и брокер одной транзакционной системой. Он заменяет распределённую атомарность следующим устройством:
1. атомарная запись внутри одной базы;
2. гарантированная асинхронная доставка наружу;
3. идемпотентная обработка повторов.
Одно из ключевых преимуществ паттерна именно в том, что он позволяет обойтись без 2PC. (microservices.io)
Что хранить в Outbox
Типичная запись содержит:
id
aggregate_type
aggregate_id
event_type
payload
created_at
Дополнительно:
correlation_id
causation_id
trace_id
tenant_id
schema_version
partition_key
published_at
retry_count
Пример:
{
"id": "event-789",
"aggregateType": "Order",
"aggregateId": "order-123",
"eventType": "OrderCreated",
"schemaVersion": 2,
"correlationId": "request-456",
"occurredAt": "2026-07-14T01:30:00Z",
"payload": {
"customerId": "customer-456",
"amount": 1000,
"currency": "EUR"
}
}
Важная тонкость: порядок событий
Допустим, один заказ последовательно изменился:
OrderCreated
OrderPaid
OrderCancelled
Желательно, чтобы события одного aggregate попадали в одну Kafka partition:
Kafka key = orderId
Тогда:
OrderCreated(order-123)
OrderPaid(order-123)
OrderCancelled(order-123)
будут упорядочены внутри одной партиции.
Но Outbox сам по себе не создаёт глобальный порядок событий всей системы. Он может сохранить причинный порядок внутри конкретного агрегата, если relay и брокер правильно используют aggregateId как ключ.
Очистка таблицы
Таблица outbox постоянно растёт, поэтому нужна стратегия очистки:
неопубликованные записи → хранить;
недавно опубликованные → временно хранить;
старые опубликованные → удалить или архивировать.
Например:
DELETE FROM outbox
WHERE published_at < NOW() - INTERVAL '7 days';
При CDC-подходе иногда применяется отдельная политика удаления после того, как WAL-изменение гарантированно было прочитано коннектором.
Нельзя бездумно удалять запись сразу после попытки отправки. Нужно понимать, в какой точке система считает публикацию завершённой и как восстанавливается после падения.
Главная формула
Transactional Outbox можно свести к такой конструкции:
Вместо:
DB write
+
Broker write
делаем:
одна транзакция:
DB business write
DB outbox write
затем:
outbox → broker
Или ещё короче:
Событие сначала становится частью локального состояния сервиса, а уже затем доставляется во внешний мир.
Паттерн не создаёт магическое exactly once. Он устраняет потерю события между локальной транзакцией и брокером, а возможные повторы закрываются идентификаторами событий и идемпотентными consumer.