Outbox
patternПаттерн доехал до прода. Мы выбрали RabbitMQ вместо Kafka, а значит надёжное хранение событий и replay из коробки нам не достались — при этом есть зоны, где доставка сообщений должна быть гарантированной, а обработка идемпотентной. Outbox закрыл ровно эту дыру.
Реализацию не писали сами: взяли транзакционный outbox из MassTransit поверх RabbitMQ — событие пишется в ту же транзакцию, что и доменные изменения, фоновый relay доносит его до брокера, а inbox на стороне потребителя даёт дедупликацию для идемпотентности.
Из хорошего чтива по теме — статья на Хабре.
Outbox pattern — это способ надежно публиковать события/сообщения из сервиса так, чтобы изменение данных в БД и отправка сообщения происходили атомарно (в рамках одной транзакции), без двухфазных коммитов и без “потерянных” событий.
Зачем он нужен
Классическая проблема: ты в обработчике запроса
- пишешь изменения в БД
- публикуешь событие в брокер (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,LastErrorPartitionKey/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
- хочется прогнозируемой надежности и трассируемости