Outbox/Inbox паттерн (PostgreSQL + RabbitMQ)
Представьте типичную ситуацию: пользователь оформляет заказ в интернет-магазине. Ваш сервис должен сделать две вещи — сохранить заказ в базе данных и опубликовать событие «заказ создан» в очередь сообщений (message queue), чтобы другие сервисы (доставка, уведомления, биллинг) могли на него отреагировать.
1. Введение
Казалось бы, что может пойти не так? Сохранили запись — отправили сообщение. Две строчки кода. Но именно здесь скрывается одна из самых коварных проблем распределённых систем — проблема двойной записи (dual write problem).
Дело в том, что база данных и брокер сообщений (message broker) — это два независимых ресурса, и у них нет общей транзакции. Рассмотрим варианты развития событий:
- Транзакция в БД зафиксирована, но приложение упало до отправки сообщения в RabbitMQ. Заказ есть, но никто в системе о нём не узнает — доставку никто не подготовит.
- Сообщение успешно отправлено в очередь, но транзакция в БД откатилась (например, сработал constraint). Другие сервисы уже начали обработку несуществующего заказа.
- Сообщение отправлено, но именно в момент подтверждения (ack) от брокера случился сетевой сбой. Ваш код думает, что публикация не удалась, и повторяет отправку — а на самом деле сообщение уже ушло. Получаем дублирование сообщений.
Ни один из этих сценариев не редкость — это статистика при достаточном количестве операций. И "просто добавить try-catch" здесь не решение: проблема не в обработке ошибок, а в том, что у вас нет атомарности между двумя разнородными хранилищами данных. Как быть в таком случае? Именно для решения этой проблемы существует паттерн Outbox (в паре с симметричным ему паттерном Inbox на стороне потребителя), который мы разберём в этой статье — от механики до практической реализации на PostgreSQL и RabbitMQ.
2. Механизм Outbox
Идея паттерна Outbox (буквально — «исходящий ящик») проста: раз мы не можем атомарно писать в два разных хранилища, давайте вообще не будем писать во второе хранилище напрямую. Вместо этого превратим запись данных в такую же строчку в базе данных, как и всё остальное, но атомарно (через транзакцию).
Как это выглядит на практике? Просто! Когда бизнес-логика сохраняет заказ, в той же самой транзакции (той же самой, что пишет в таблицу orders) она добавляет ещё одну запись — в таблицу outbox_messages. Это просто ещё один INSERT рядом с остальными. База данных гарантирует атомарность: либо зафиксируются обе записи (заказ и сообщение-в-очередь), либо не зафиксируется ни одна. Никакого "разрыва событий" между "сохранил" и "отправил" — потому что отправки в этот момент ещё не происходит в принципе.
Сама публикация в RabbitMQ выносится в отдельный процесс — так называемый релей (relay) или publisher. Его задача простая и монотонная: периодически проверять таблицу outbox_messages на предмет неотправленных записей, публиковать их в очередь и помечать как отправленные (или удалять). Есть события - отправляем, нет событий - ничего не делаем.
Из чего состоит механизм:
- Таблица outbox — хранит будущие сообщения: тип события, полезная нагрузка (payload) в виде JSON, время создания, статус (отправлено / не отправлено).
- Транзакционная запись — INSERT в outbox происходит строго вместе с бизнес-операцией, в одной транзакции.
- Релей — фоновый процесс, который вычитывает необработанные строки и публикует их в RabbitMQ.
- Подтверждение публикации — после успешного ack от брокера релей помечает запись как обработанную (или физически удаляет — это вопрос дизайна, к которому мы ещё вернёмся).
Ключевой момент здесь такой: мы не пытаемся сделать распределённую транзакцию между БД и брокером — мы её попросту избегаем, сводя задачу к одной локальной транзакции в PostgreSQL. Надёжность распределённой части (доставка в очередь) перекладывается на релей, который может безопасно повторять попытки — потому что источник истины (Source of Truth) один, это таблица в базе данных.
Теперь пара диаграмм, или если быть точнее sequence-диаграмм.

Диаграмма показывает две фазы: сначала "Приложение" пишет заказ и запись outbox одной транзакцией в PostgreSQL (шаги 1-3), затем независимый "Outbox relay" опрашивает таблицу, публикует сообщение в RabbitMQ и помечает запись обработанной (шаги 4-8) — эти две фазы разнесены во времени и никак не связаны общей транзакцией, в этом и суть паттерна.
3. Практическая настройка Outbox
Перейдём от теории к коду. Разберём три составляющие: схему таблицы, транзакционную запись через EF Core и варианты "Outbox reley".
Схема таблицы outbox
Минимальный набор полей выглядит так:
CREATE TABLE outbox_messages (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
type VARCHAR(255) NOT NULL, -- тип события, напр. "OrderCreated"
payload JSONB NOT NULL, -- сериализованное сообщение
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
processed_at TIMESTAMPTZ NULL -- NULL = ещё не отправлено
);
CREATE INDEX ix_outbox_unprocessed ON outbox_messages (created_at)
WHERE processed_at IS NULL;
Частичный индекс (partial index) по processed_at IS NULL — важная деталь: релей на каждом цикле опроса делает SELECT ... WHERE processed_at IS NULL ORDER BY created_at, и без такого индекса запрос со временем начнёт сканировать всю таблицу, включая давно обработанные записи.
Транзакционная запись через EF Core
Ключевой момент — вставка бизнес-сущности и сообщения outbox должны попасть в одну транзакцию. Если Вы работаете с одним DbContext, EF Core делает это автоматически в рамках SaveChanges():
public async Task CreateOrderAsync(Order order, CancellationToken ct)
{
var outboxMessage = new OutboxMessage
{
Type = nameof(OrderCreated),
Payload = JsonSerializer.Serialize(new OrderCreated(order.Id, order.CustomerId)),
CreatedAt = DateTimeOffset.UtcNow
};
_dbContext.Orders.Add(order);
_dbContext.OutboxMessages.Add(outboxMessage);
await _dbContext.SaveChangesAsync(ct); // одна транзакция на обе таблицы
}
Явного BeginTransaction() не требуется — EF Core сам оборачивает все изменения одного вызова SaveChangesAsync в транзакцию БД. Важно только не разбивать это на два вызова SaveChangesAsync — тогда гарантия атомарности пропадёт.
Варианты Publisher
Здесь есть развилка, и стоит понимать компромиссы каждого варианта:
- Polling publisher — фоновый
BackgroundService/IHostedService, который каждые N секунд (например, 1-2) выполняет запрос к таблице, публикует найденные сообщения в RabbitMQ и обновляетprocessed_at.
Плюсы: просто реализовать, не требует ничего кроме самой БД.
Минусы: задержка публикации равна интервалу опроса; при большой нагрузке — дополнительный постоянный поток запросов к БД. - CDC / логическая репликация PostgreSQL — вместо опроса используется механизм логического декодирования (logical decoding) через слот репликации (replication slot): изменения в таблице outbox доставляются релею почти в реальном времени, без поллинга. Обычно реализуется через Debezium или аналогичный коннектор.
Плюсы: минимальная задержка, нет нагрузки от периодических SELECT.
Минусы: заметно сложнее в эксплуатации — нужен отдельный компонент, слоты репликации требуют мониторинга, потому что незакрытый слот может раздуть Write-Ahead Log (WAL).
Для большинства проектов среднего размера polling publisher — разумная отправная точка: он проще в отладке и не добавляет инфраструктурной сложности и я уже писал ранее про реализацию на базе этого способа. К CDC имеет смысл переходить, когда задержка в секунду-две становится критичной для бизнес-сценария.
4. Механизм Inbox
Паттерн Outbox решает проблему на стороне отправителя. Но у получателя сообщения есть симметричная проблема — паттерн Inbox её решает.
В чём проблема на стороне потребителя? RabbitMQ, как и большинство брокеров, гарантирует доставку по модели at-least-once — "как минимум один раз". Это значит, что одно и то же сообщение может прийти к обработчику (consumer) дважды, а иногда и больше. Причины этого могут быть такими:
- Потребитель обработал сообщение, но упал до отправки подтверждения (ack) брокеру — брокер решит, что сообщение не обработано, и доставит его снова.
- Сработал механизм повторной доставки (redelivery) после разрыва соединения — брокер не знает, было сообщение обработано или нет, и на всякий случай отправляет его ещё раз.
- Сам потребитель реализует логику retry при временных сбоях (например, недоступна внешняя система) — и по ошибке в коде retry затрагивает уже успешно обработанное сообщение.
Если обработчик не идемпотентен (idempotent) — то есть повторная обработка одного и того же сообщения меняет состояние системы второй раз — начинаются проблемы: клиенту дважды спишут деньги, дважды отправят email, дважды создадут запись в другой системе. Поверьте, это очень неприятные проблемы, потому что их очень сложно отловить и не только по этому.
Механизм Inbox
Идея зеркальна Outbox: прежде чем применить бизнес-логику сообщения, обработчик проверяет — а не видел ли он уже это сообщение? Для этого:
- У каждого сообщения есть уникальный идентификатор (
message_id), который отправитель проставляет один раз при создании и не меняет при повторных доставках. - Обработчик, получив сообщение, в той же транзакции, где применяет бизнес-логику, добавляет запись в таблицу
inbox_messagesс этимmessage_id. - Если такая запись уже существует (нарушение уникального ограничения — unique constraint), обработчик понимает, что сообщение уже было обработано, и просто подтверждает получение (ack), не повторяя бизнес-логику.
Ключевой момент: проверка "видел ли я это сообщение" и применение бизнес-логики происходят в одной транзакции — точно так же, как в Outbox запись сообщения и бизнес-операция были в одной транзакции. Это тот же самый архитектурный приём, применённый в обратную сторону: локальная транзакция БД становится единственным источником истины (source of truth) вместо попытки координировать состояние между БД и логикой ретраев брокера.

Диаграмма показывает развилку после попытки вставки в inbox_messages: если запись новая — применяется бизнес-логика, если сработало нарушение уникального ограничения — логика пропускается. В обоих случаях транзакция коммитится и обработчик подтверждает сообщение брокеру одинаково — асимметрия только в том, выполняется ли содержательная часть обработки.
5. Практическая настройка Inbox
Схема таблицы inbox
CREATE TABLE inbox_messages (
message_id UUID PRIMARY KEY, -- приходит от отправителя, не генерируется здесь
message_type VARCHAR(255) NOT NULL,
received_at TIMESTAMPTZ NOT NULL DEFAULT now(),
processed_at TIMESTAMPTZ NULL
);
Обратите внимание на разницу со схемой outbox: там id генерировался при создании записи (gen_random_uuid()), здесь же message_id — это первичный ключ (primary key), значение которого приходит извне, от отправителя. Именно уникальность этого ключа и обеспечивает дедупликацию (deduplication) — вставить вторую строку с тем же message_id физически невозможно, база данных сама выбросит 23505 unique_violation.
Откуда берётся message_id? Очень просто! Идентификатор должен генерироваться один раз — на стороне отправителя, в момент создания сообщения outbox — и оставаться неизменным при всех последующих повторных доставках. Проще всего использовать тот же id, что и в таблице outbox_messages: он уже гарантированно уникален и создаётся один раз при вставке.
Транзакционная связка через EF Core
public async Task HandleOrderCreatedAsync(OrderCreated message, Guid messageId, CancellationToken ct)
{
try
{
_dbContext.InboxMessages.Add(new InboxMessage
{
MessageId = messageId,
MessageType = nameof(OrderCreated),
ReceivedAt = DateTimeOffset.UtcNow
});
// Бизнес-логика — те же изменения, что и запись в inbox,
// попадают в один и тот же SaveChangesAsync
await _shippingService.PrepareShipmentAsync(message.OrderId, ct);
await _dbContext.SaveChangesAsync(ct);
}
catch (DbUpdateException ex) when (IsUniqueViolation(ex))
{
// message_id уже существует — сообщение уже обработано, повторно логику не выполняем
return;
}
}
Обработчик исключения DbUpdateException с проверкой на код ошибки 23505 — это и есть точка, где происходит развилка с прошлой диаграммы: если SaveChangesAsync упал именно на нарушении уникальности, значит сообщение дублирующееся, и метод просто завершается — Rebus (или другой шинный фреймворк) после этого спокойно подтвердит (ack) сообщение брокеру, ведь исключение не "вылетело" наружу.
Важный нюанс: что если бизнес-логика вызывает внешний сервис
Если PrepareShipmentAsync не просто пишет в БД, а делает вызов к внешнему HTTP API — транзакционная гарантия по идемпотентности здесь уже не спасает: если вызов внешнего сервиса прошёл, а SaveChangesAsync после него упал по другой причине (не по дубликату), при повторной обработке внешний вызов произойдёт ещё раз. В таких случаях insert в inbox_messages стоит делать до вызова внешнего сервиса — отдельным SaveChangesAsync — чтобы проверка на дубликат срабатывала раньше, чем побочный эффект вовне. Это компромисс: вы жертвуете единой транзакцией ради более раннего обнаружения дубликата.
6. Гарантии доставки
Разберём, какие гарантии на самом деле даёт связка Outbox/Inbox, и какие пограничные ситуации (edge cases) стоит держать в голове.
At-least-once, а не exactly-once
Важно сразу снять иллюзию: Outbox/Inbox не превращает систему в exactly-once ("доставлено ровно один раз"). Правильнее говорить про effectively-once — эффект достигается один раз, но механика под капотом всё равно at-least-once с дедупликацией:
- Outbox гарантирует, что сообщение не потеряется — если бизнес-транзакция зафиксирована, запись в outbox тоже зафиксирована, и релей рано или поздно её отправит.
- Outbox не гарантирует, что сообщение не продублируется на пути к брокеру — если релей упал между публикацией в RabbitMQ и обновлением
processed_at, при перезапуске он опубликует то же сообщение повторно. - Inbox как раз закрывает эту дыру — гасит дубликаты на стороне потребителя.
Вместе они дают ту самую гарантию "эффект как будто ровно один раз", которую и ищут в большинстве бизнес-сценариев.
Что если релей упал в разных точках
- Упал до публикации — ничего не потеряно, при перезапуске найдёт необработанную запись и опубликует.
- Упал после публикации, но до обновления
processed_at— при перезапуске снова найдёт "необработанную" запись (она ведь не помечена) и опубликует повторно. Дубликат уходит в RabbitMQ — но именно поэтому Inbox на стороне потребителя обязателен, а не опционален. - Упал во время обновления
processed_at(например, сеть моргнула после физического UPDATE, но до ack клиенту) — по сути то же самое, что пункт 2: возможен дубликат, гасится Inbox'ом.
Отсюда практический вывод: Outbox без Inbox на другой стороне — это неполное решение. Если вы публикуете событие для внешнего сервиса, на который у Вас нет контроля над идемпотентностью обработчика — вы взяли на себя риск дублей.
Порядок доставки (ordering)
Polling publisher, вычитывающий записи ORDER BY created_at, обычно сохраняет порядок публикации в рамках одной таблицы outbox. Но гарантии это не даёт железно:
- При нескольких инстансах релея (несколько подов/процессов), работающих параллельно с одной таблицей, порядок публикации между инстансами не гарантирован — они могут выхватывать разные строки одновременно. Если порядок важен, нужен либо один активный релей (лидер), либо
SELECT ... FOR UPDATE SKIP LOCKEDс партиционированием по ключу агрегата. - RabbitMQ сам по себе не гарантирует порядок между несколькими consumer'ами одной очереди при параллельной обработке — если важен порядок, нужна одна очередь на один consumer или маршрутизация по ключу в один и тот же consumer.
Разрастание таблиц
Обе таблицы, outbox_messages и inbox_messages, растут бесконечно, если ничего не чистить. Нужен фоновый job, удаляющий обработанные записи старше некоторого срока (например, 7-30 дней) — этого времени обычно достаточно, чтобы пережить любые ретраи и разборы инцидентов, но не накапливать историю вечно. А можно сразу удалять записи, которые уже обработаны, тогда наличие записей будет служить поводом подумать "почему они тут так долго? может быть была ошибка при обработке?"
7. Смежный паттерн: Saga
Outbox/Inbox решают задачу надёжной доставки одного сообщения. Но что, если бизнес-операция состоит из нескольких шагов, разнесённых по разным сервисам — например, оформление заказа требует: резервирования товара на складе, списания оплаты и создания отгрузки, — и любой шаг может не выполниться?
Здесь на сцену выходит паттерн Saga — способ координировать распределённую последовательность локальных транзакций, каждая из которых публикует событие, запускающее следующий шаг, а в случае сбоя — компенсирующие действия (compensating actions), откатывающие уже выполненные шаги.
Почему Saga естественно опирается на Outbox/Inbox
Каждый шаг саги — это, по сути, тот же паттерн, что мы разобрали:
- Сервис-участник получает событие (Inbox гасит дубли, если событие пришло повторно).
- Выполняет свою локальную транзакцию (например, резервирует товар).
- Публикует следующее событие в цепочке (через свой Outbox, в той же транзакции, что и резервирование).
Получается, что сага — это не какой-то отдельный инфраструктурный механизм поверх шины сообщений, а скорее бизнес-хореография (choreography) или оркестрация (orchestration) поверх уже знакомой связки Outbox/Inbox на каждом шаге. Если в каждом сервисе-участнике эта связка уже надёжно реализована, сага "наследует" её гарантии — не теряет события и не выполняет шаги дважды.
Хореография против оркестрации — коротко
- Хореография — каждый сервис публикует событие, на которое реагирует следующий сервис в цепочке; никто не знает про сагу "целиком". Проще для небольшого числа шагов, но с ростом цепочки логику становится сложно отследить — она размазана по обработчикам разных сервисов.
- Оркестрация — выделенный координатор (оркестратор саги) явно знает все шаги и компенсации, и рассылает команды участникам, слушая их ответы. Легче отслеживать состояние всей саги в одном месте, но добавляет ещё один компонент, которому тоже нужны свои Outbox/Inbox.
Подробное сравнение этих двух подходов — тема отдельной статьи; здесь важно зафиксировать главное: без надёжной доставки на уровне каждого шага (то есть без Outbox/Inbox) сага в принципе не может дать содержательных гарантий — компенсирующее действие само может потеряться или продублироваться точно так же, как любое другое сообщение.
Так выглядит хореография — каждый сервис публикует своё событие через Outbox и реагирует на событие соседа через Inbox, без единой точки, знающей про сагу целиком.

А здесь — оркестрация: координатор явно знает все шаги и рассылает команды каждому участнику. Пунктирная стрелка от Payment service обратно к оркестратору иллюстрирует обратную связь: если участник сообщает об ошибке, именно оркестратор решает, какие компенсирующие действия запустить у уже выполненных шагов.
8. Мониторинг
Outbox/Inbox — это не паттерн, который можно один раз внедрить и забыть. Обе таблицы работают в фоне, и если не наблюдать за ними, деградация будет незаметной до тех пор, пока не превратится в инцидент. Разберём, что именно стоит выносить в метрики (Prometheus) и что — в логи (ELK).
Метрики в Prometheus
- Размер очереди outbox (backlog) — количество строк с
processed_at IS NULL. Простой gauge, снимаемый периодическим запросом:
SELECT count(*) FROM outbox_messages WHERE processed_at IS NULL;
Растущий backlog — первый сигнал, что релей не справляется или вовсе не работает.
- Задержка релея (publish lag) — разница между
now()иcreated_atсамой старой необработанной записи. Backlog может быть небольшим по количеству строк, но если среди них есть запись пятичасовой давности — это уже проблема, даже если очередь в целом не растёт. - Количество дублей в inbox — счётчик (counter), инкрементируемый в том самом
catch (DbUpdateException exception) when (IsUniqueViolation(exception))из раздела 5. Сам по себе рост этого счётчика — не авария (дубли — ожидаемая часть at-least-once доставки), но резкий скачок стоит отслеживать на предмет — не начал ли релей штормить повторными публикациями. - Частота ошибок релея — сколько раз попытка публикации в RabbitMQ завершилась ошибкой (недоступен брокер, таймаут). Здесь важны и абсолютные значения, и производная — резкий рост подряд идущих ошибок обычно означает проблему с самим брокером, а не с отдельным сообщением.
Что писать в логи (ELK)
- Каждую публикацию из релея — с
message_id, типом события и временем задержки отcreated_atдо момента публикации. Это даёт возможность восстановить путь конкретного сообщения при разборе инцидента. - Каждое срабатывание дедупликации в inbox — с
message_idи именем обработчика. Полезно при расследовании: если для одногоmessage_idтаких записей в логах много — вероятно, релей действительно шлёт этот message с проблемной регулярностью, и стоит смотреть в его сторону. - Ошибки релея при публикации — с полным контекстом (какая запись, какая попытка по счёту, если реализован retry с backoff).
Алерты, которые имеет смысл настроить
- Backlog outbox превышает пороговое значение N в течение M минут.
- Lag (задержка самой старой записи) превышает допустимое для бизнеса значение (например, 5 минут для сценариев, где пользователь ждёт подтверждения).
- Релей не публиковал ничего X минут при том, что backlog не равен нулю — верный признак, что процесс релея упал или завис.
Без этих трёх сигналов легко получить ситуацию, когда всё выглядит нормально в приложении (ошибок в бизнес-логике нет), а на деле сообщения неделями лежат неотправленными — потому что упал именно фоновый релей, а не что-то в основном потоке запросов.
9. Заключение
Мы разобрали Outbox/Inbox от проблемы двойной записи до конкретных схем таблиц, кода для EntityFrameworkCore, диаграмм взаимодействия и мониторинга. Но стоит закончить статью честным вопросом: а когда этот паттерн вообще нужен?
Когда паттерн оправдан
- У вас несколько сервисов, и потеря или дублирование события имеет реальную бизнес-цену (деньги, товар на складе, повторная отправка уведомления клиенту).
- Бизнес-операция публикует событие как часть более широкого процесса — саги, интеграции с внешней системой, аналитики, на которую полагаются другие команды.
- Вы уже сталкивались (или предвидите) с ситуацией, когда сообщение потерялось из-за падения процесса между записью в БД и отправкой в очередь.
Когда это чрезмерно
- Простой CRUD-сервис без интеграций — если событие никто не слушает, публиковать его через Outbox незачем.
- Прототип или MVP, где скорость разработки важнее надёжности доставки, а бизнес-цена потери сообщения близка к нулю.
- Сценарии, где допустима idempotent-обработка на уровне "просто перезапустить вручную при сбое" — то есть человек в цикле, а не автоматика.
Главный компромисс
Outbox/Inbox покупает надёжность ценой дополнительной инфраструктурной сложности: две новые таблицы, фоновый процесс-релей, мониторинг backlog и lag, задачи по очистке старых записей. Это разумная цена, если распределённая природа системы всё равно требует решать проблему двойной записи — но плохая идея, если вы вносите эту сложность "на всякий случай", без реального сценария, где она отрабатывает.
Как и с большинством архитектурных решений, здесь нет универсально правильного ответа — есть конкретный набор гарантий, которые нужны именно вашей системе, и цена, которую вы готовы за них заплатить.
Пишите правильный код и делайте правильную архитектуру.
Мои видео
Boosty.to | YouTube | Yandex.Дзен | RuTube | VK Video