Что такое Transactional Outbox

Transactional Outbox — это архитектурный паттерн, который позволяет надёжно выполнить две связанные операции:

  1. изменить данные в локальной базе сервиса;
  2. отправить событие или сообщение в 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.

Прокрутить вверх