Волобуев Петр | Технологический радарВолобуев Петр | Технологический радар

Outbox

pattern
Использую

Паттерн доехал до прода. Мы выбрали RabbitMQ вместо Kafka, а значит надёжное хранение событий и replay из коробки нам не достались — при этом есть зоны, где доставка сообщений должна быть гарантированной, а обработка идемпотентной. Outbox закрыл ровно эту дыру.

Реализацию не писали сами: взяли транзакционный outbox из MassTransit поверх RabbitMQ — событие пишется в ту же транзакцию, что и доменные изменения, фоновый relay доносит его до брокера, а inbox на стороне потребителя даёт дедупликацию для идемпотентности.

Из хорошего чтива по теме — статья на Хабре.

Использую

Outbox pattern — это способ надежно публиковать события/сообщения из сервиса так, чтобы изменение данных в БД и отправка сообщения происходили атомарно (в рамках одной транзакции), без двухфазных коммитов и без “потерянных” событий.

Зачем он нужен

Классическая проблема: ты в обработчике запроса

  1. пишешь изменения в БД
  2. публикуешь событие в брокер (Kafka/RabbitMQ/…)

Если упасть между (1) и (2) — данные в БД уже есть, а событие не ушло. Если наоборот (событие ушло, а транзакция откатилась) — получатель увидит событие про то, чего “не случилось”. Outbox устраняет оба сценария.

Идея паттерна

Вместо “сразу в брокер” сервис пишет событие в таблицу outbox в той же транзакции, что и доменные изменения:

  • Транзакция:

    • UPDATE/INSERT доменных таблиц
    • INSERT INTO outbox (...) VALUES (...)
  • Отдельный фоновой процесс (publisher/relay) читает outbox и отправляет в брокер.

  • После успешной отправки помечает запись как отправленную (или удаляет).

В итоге: если в БД есть изменение — в outbox гарантированно есть событие, и оно рано или поздно будет опубликовано.

Типовая схема данных Outbox

Минимально полезные поля:

  • Id (UUID/ULID)
  • OccurredAt (когда событие “произошло”)
  • Type (имя события/контракта)
  • Payload (JSON/Protobuf/base64)
  • Headers/Metadata (traceId, tenantId, correlationId, schemaVersion)
  • Status (Pending/Sent/Failed)
  • RetryCount, NextAttemptAt, LastError
  • PartitionKey/AggregateId (для порядка по сущности)

Как “забирать” сообщения из outbox

Два популярных подхода:

1) Polling (опрос)

Publisher периодически делает запрос:

  • выбрать N “pending” строк (часто ORDER BY occurred_at)
  • заблокировать их (FOR UPDATE SKIP LOCKED в Postgres) — чтобы несколько воркеров не взяли одно и то же
  • отправить в брокер
  • отметить Sent

Плюсы: проще, portable. Минусы: задержка (latency), нагрузка на БД.

2) CDC (Change Data Capture)

Сервис всё равно пишет в outbox, но публикация делается через чтение WAL/binlog (Debezium и т.п.), и дальше отправка в Kafka/… Плюсы: низкая задержка и меньше polling-нагрузки. Минусы: сложнее эксплуатация.

Гарантии и важные свойства

  • Outbox даёт at-least-once delivery: событие может быть отправлено повторно.

  • Поэтому потребители обязаны быть идемпотентными:

    • хранить “processed message ids”
    • использовать upsert по ключу события/агрегата
    • дедупликация на стороне consumer

Exactly-once в распределенных системах почти всегда достигается только “через дизайн” (идемпотентность + дедуп), а не за счет магии транспорта.

Порядок сообщений

Если нужен порядок по агрегату (например, OrderId):

  • кладёшь AggregateId/PartitionKey
  • отправляешь в Kafka в партицию по ключу (тогда порядок в рамках ключа сохраняется)
  • при polling - выбираешь/публикуешь последовательно по ключу (или обеспечиваешь сериализацию на publisher’е)

Глобальный порядок обычно не нужен и дорог.

Когда Outbox особенно нужен

  • микросервисы + брокер сообщений
  • бизнес-события должны “не потеряться”
  • нельзя (или не хочется) тащить 2PC/Distributed Transactions
  • хочется прогнозируемой надежности и трассируемости