Transactional Outbox/Inbox: события без распределённой транзакции

Как надёжно публиковать события при изменении данных без распределённой транзакции: outbox, relay/poller, inbox и дедупликация, idempotency key, ordering, poison messages и почему exactly-once — это иллюзия

Почти любая event-driven система упирается в один и тот же вопрос: как одновременно изменить данные в своей базе и опубликовать событие об этом изменении — так, чтобы либо случилось и то и другое, либо ничего? Без решения этой задачи система рано или поздно окажется в состоянии, когда заказ создан, а событие OrderCreated потеряно (или наоборот — событие ушло, а транзакция откатилась).

Это проблема dual-write: два независимых хранилища (БД и брокер), и нет общей транзакции, которая их охватывает. Transactional Outbox решает её красиво и без распределённых транзакций. Парный к нему Inbox защищает потребителя от неизбежных дубликатов. Вместе они — рабочий скелет production event-driven систем, без которого остальные паттерны серии стоят на песке.

В статье

Проблема dual-write

Наивный код выглядит так:

tx.begin()
saveOrder(order)        // запись в БД
tx.commit()
broker.publish(event)   // публикация в брокер

Между commit() и publish() процесс может упасть. Тогда заказ есть, а события нет — потребители никогда не узнают. Поменять порядок (сначала publish, потом commit) не помогает: тогда возможно событие без заказа.

Соблазн — обернуть оба действия в распределённую транзакцию (2PC между БД и брокером). Но это дорого, блокирующе и поддерживается не всеми брокерами. Outbox обходит проблему, сведя всё к одной локальной транзакции БД.

Outbox: событие как часть транзакции

Идея проста: событие пишется в ту же базу, в той же транзакции, что и бизнес-данные — в специальную таблицу outbox.

tx.begin()
saveOrder(order)                 // бизнес-данные
saveToOutbox(event)              // событие — в той же транзакции
tx.commit()                      // атомарно: либо оба, либо ничего

Теперь нет dual-write: или закоммитились и заказ, и запись в outbox, или ничего. Событие гарантированно сохранено вместе с данными. Осталось доставить его в брокер — этим занимается отдельный процесс (relay).

Минимальная структура записи outbox: id, aggregate_id, event_type, payload, created_at, published_at (NULL пока не отправлено), attempts.

Relay: как события попадают в брокер

Relay (он же message relay, publisher, dispatcher) — процесс, который читает неопубликованные записи из outbox и шлёт их в брокер. Два подхода:

  • Polling publisher — relay периодически опрашивает таблицу (SELECT ... WHERE published_at IS NULL ORDER BY id LIMIT N), публикует, помечает published_at. Просто, работает на любой БД, но создаёт нагрузку опросом и добавляет задержку.
  • CDC (Change Data Capture) — relay читает не таблицу, а лог транзакций БД (WAL в PostgreSQL, binlog в MySQL) через инструмент вроде Debezium. Меньше задержка, нет нагрузки опросом, но это дополнительный инфраструктурный компонент.

Для старта polling почти всегда достаточно. CDC оправдан, когда важна низкая задержка или объём событий велик.

Важно: relay гарантирует доставку at least once — если процесс упал между публикацией и проставлением published_at, после рестарта он отправит событие повторно. С этим разбирается потребитель.

Inbox и дедупликация на стороне потребителя

Раз доставка at-least-once, потребитель обязан быть готов к дубликатам. Inbox — зеркальный паттерн: таблица обработанных событий на стороне потребителя.

tx.begin()
inserted = inbox.tryInsert(event.id)   // INSERT ... ON CONFLICT DO NOTHING
if not inserted:                        // строки не появилось — уже обрабатывали
    tx.rollback()
    ack()
    return
applyBusinessLogic(event)
tx.commit()                             // обработка и отметка — атомарно
ack()

Критичный момент — атомарность проверки и вставки. Соблазнительный вариант «сначала contains, потом save» двумя шагами содержит гонку: два конкурентных обработчика одного события успеют проверить «нет такого» до того, как любой из них вставит запись, и оба применят бизнес-логику. Правильно — одна операция с уникальным ограничением на event.id: INSERT ... ON CONFLICT DO NOTHING. Если строки не добавилось — событие уже обработано. Запись в inbox и бизнес-логика идут в одной транзакции, поэтому «обработано» и «зафиксировано как обработанное» атомарны. Inbox можно периодически чистить от старых записей (по TTL).

Idempotency key, correlation и causation id

Три идентификатора, без которых распределённую систему невозможно отлаживать и делать надёжной:

  • Idempotency key — уникальный ключ операции (часто = event.id). По нему inbox и любые внешние вызовы распознают повтор. Главный инструмент борьбы с дубликатами.
  • Correlation id — сквозной идентификатор всей цепочки, порождённой одним исходным запросом. Один и тот же на всех событиях и вызовах процесса. Позволяет собрать всю историю «что произошло из-за вот этого клика».
  • Causation id — id события/команды, которое непосредственно породило текущее. Восстанавливает дерево причинно-следственных связей: что чем вызвано.

Correlation отвечает на «к какому процессу это относится», causation — на «что именно это вызвало». Оба стоит класть в metadata каждого события с самого начала: задним числом добавить их в работающую систему мучительно.

Ordering: порядок, который легко потерять

At-least-once и параллельная обработка легко ломают порядок: событие B может быть обработано раньше A, даже если опубликовано позже. Где это критично (например, AccountCredited до AccountOpened) — варианты:

  • Партиционирование по ключу агрегата — все события одного агрегата идут в одну партицию/очередь и обрабатываются последовательно (модель Kafka).
  • Версия/порядковый номер в событии — потребитель отбрасывает или буферизует события «из будущего».
  • Проектировать обработчики устойчивыми к переупорядочиванию — там, где это возможно (коммутативные операции).

Глобальный порядок через всю систему — дорогая иллюзия. Обычно достаточно порядка в пределах одного агрегата.

Retry и poison messages

Обработка события может упасть. Политика повторов:

  • Retry с backoff — временные сбои (сеть, недоступность БД) лечатся повтором с нарастающей задержкой.
  • Poison message — событие, которое падает всегда (битый payload, баг в обработчике). Бесконечный retry заблокирует очередь. Такое событие после N попыток уходит в dead-letter queue (DLQ) для ручного разбора.
  • Различать «упало временно» и «упадёт всегда» — иначе либо теряете данные, либо вечно крутите один и тот же сбой.

DLQ — не свалка, а очередь, которую нужно мониторить и разгребать. Сообщения в ней — это сигнал о проблеме.

Exactly-once как системная иллюзия

«Exactly-once delivery» в общем распределённом случае на уровне транспорта недостижим: между «доставил» и «подтвердил доставку» всегда есть окно, где возможен сбой и повтор. Поэтому в типичном межсервисном сценарии — где на другом конце внешняя БД и сторонние API — обработку приходится проектировать как at-least-once. (Оговорка: у Kafka есть транзакции и exactly-once semantics, но в ограниченной модели «Kafka → обработка → Kafka»; как только в цепочке появляется внешняя система, гарантия снова сводится к идемпотентности на вашей стороне.)

Но exactly-once processing — достижимый и правильный целевой результат. Он получается не от магии брокера, а из связки: at-least-once доставка + идемпотентные обработчики (inbox + idempotency key). Событие может прийти дважды, но эффект применится один раз.

Запомнить: не ищите брокер с «exactly-once», стройте идемпотентную обработку. Это и есть инженерный ответ.

Лунная база: событие о расходе ресурса

Модуль жизнеобеспечения списывает кислород и должен сообщить об этом подсистеме мониторинга и прогноза:

  1. В одной транзакции: уменьшить остаток кислорода + записать OxygenConsumed в outbox.
  2. Relay публикует событие в брокер (at-least-once).
  3. Подсистема прогноза получает событие, проверяет inbox по event.id, применяет к своей модели, фиксирует id.
  4. Если событие пришло повторно (relay перезапустился) — inbox отбрасывает дубль, прогноз не искажается.

Без outbox возможен сценарий: кислород списан, а монитор не узнал — и аварийный порог сработает с опозданием. Цена ошибки на лунной базе наглядна, но в обычном биллинге она ровно та же.

Мониторинг: за чем следить

Outbox — не «настроил и забыл». Минимум метрик:

  • лаг outbox — возраст самой старой неопубликованной записи (now - min(created_at) WHERE published_at IS NULL). Растёт → relay не справляется или упал;
  • размер outbox — число неотправленных записей;
  • размер DLQ — сколько poison messages ждут разбора;
  • частота повторов — всплеск retry сигналит о деградации downstream.

Лаг outbox — важнейшая метрика: именно он показывает, насколько read-сторона системы отстаёт от реальности.

Checklist: готов ли ваш outbox

  1. Пишется ли событие в outbox в той же транзакции, что и бизнес-данные?
  2. Есть ли relay и понятно ли, polling это или CDC?
  3. Идемпотентны ли потребители (inbox или эквивалент)?
  4. Проставляются ли idempotency / correlation / causation id с самого начала?
  5. Определён ли порядок там, где он критичен (партиционирование по агрегату)?
  6. Куда уходят poison messages и кто разбирает DLQ?
  7. Мониторится ли лаг outbox и размер DLQ?

Если на пункт 1 ответ «нет» — у вас не outbox, а dual-write с лишней таблицей. Если «не знаем» на 3 — система генерирует дубли эффектов, просто вы их ещё не заметили.

Документация и первоисточники

Обсуждение в Telegram

Присоединиться →

Комментарии