Outbox + Inbox: событие и данные в одной транзакции
Outbox + Inbox: событие и данные в одной транзакции
Представьте: use-case завершает урок, записывает прогресс в PostgreSQL, а затем публикует событие lesson.completed в Kafka. Между этими двумя операциями приложение падает. Прогресс сохранён, но событие не отправлено. Achievement-сервис ничего не узнал. Analytics не записала метрику. Пользователь не получил бейдж.
Это называется dual-write problem - запись в два независимых хранилища без общей транзакции. Outbox закрывает именно эту щель, и важно понимать, что он закрывает только её.
Outbox атомарно фиксирует бизнес-изменение и намерение опубликовать событие - одной локальной транзакцией одной базы (три смысла слова «атомарно»). Доставка из outbox остаётся at-least-once: она требует повторов на стороне relay и идемпотентности на стороне подписчика.
Разница видна на схеме. Атомарна только верхняя часть:
транзакция БД ← здесь атомарность есть
├── UPDATE orders
└── INSERT outbox
↓
relay (poller) ← здесь её нет и быть не может
↓
брокер
Relay - обычный процесс, и с ним может случиться что угодно: упасть
до публикации, опубликовать и упасть до отметки published_at, отправить
событие второй раз, не отправлять его час, потому что лежал. Ни один
из этих случаев outbox не отменяет - он лишь гарантирует, что событие
есть что отправить.
Отсюда практический вывод: outbox не заменяет идемпотентность подписчика, а требует её. Пара «outbox + inbox» работает вместе именно потому, что каждая половина закрывает свою щель: первая - между данными и намерением, вторая - между доставкой и эффектом.
Dual-write problem
Две операции - запись в БД и публикация в брокер - не могут быть атомарными, потому что PostgreSQL и Kafka - разные системы. Нельзя обернуть их в одну транзакцию. Какой бы порядок вы ни выбрали, есть окно для сбоя:
- Сначала БД, потом брокер. Если упали после коммита в БД - данные изменились, но событие не ушло. Подписчики не узнали.
- Сначала брокер, потом БД. Если упали после публикации - подписчики получили событие о данных, которых нет в БД. Ещё хуже.
Dual-write нельзя решить retry-ами. Нужен другой подход.
Outbox: событие как часть транзакции
Идея: не отправляем событие в брокер напрямую. Вместо этого пишем его в специальную таблицу outbox в той же транзакции, что и бизнес-данные. Отдельный процесс (poller) читает outbox и отправляет события в брокер.
Если транзакция коммитится - и данные, и событие сохранены. Если откатывается - ничего не сохранено. Атомарность гарантирована средствами БД.
DDL для outbox-таблицы
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
event_id VARCHAR(64) NOT NULL UNIQUE,
event_type VARCHAR(128) NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
published_at TIMESTAMPTZ,
attempts INT NOT NULL DEFAULT 0
);
-- Poller выбирает неотправленные события по этому индексу
CREATE INDEX idx_outbox_unpublished
ON outbox (created_at)
WHERE published_at IS NULL;
Поля:
event_id- уникальный идентификатор события (UUID/ULID)event_type- тип события (lesson.completed,user.registered)payload- JSON с данными событияpublished_at- NULL, пока событие не отправлено в брокерattempts- счётчик попыток отправки (для мониторинга застрявших)
Запись в outbox внутри транзакции
Use-case сохраняет бизнес-данные и событие в одной транзакции:
// CompleteLessonUseCase завершает урок и записывает событие в outbox.
type CompleteLessonUseCase struct {
db *sql.DB
}
// Execute выполняет завершение урока атомарно с outbox-записью.
func (uc *CompleteLessonUseCase) Execute(ctx context.Context, userID, lessonID int64) error {
tx, err := uc.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
// 1. Бизнес-логика: обновляем прогресс
_, err = tx.ExecContext(ctx,
`UPDATE lesson_progress
SET completed = true, completed_at = now()
WHERE user_id = $1 AND lesson_id = $2`,
userID, lessonID,
)
if err != nil {
return fmt.Errorf("update progress: %w", err)
}
// 2. Записываем событие в outbox (в той же транзакции!)
eventID := ulid.Make().String()
payload, _ := json.Marshal(LessonCompletedPayload{
UserID: userID,
LessonID: lessonID,
})
_, err = tx.ExecContext(ctx,
`INSERT INTO outbox (event_id, event_type, payload)
VALUES ($1, $2, $3)`,
eventID, "lesson.completed", payload,
)
if err != nil {
return fmt.Errorf("insert outbox: %w", err)
}
// Коммит: и прогресс, и событие сохранены атомарно
return tx.Commit()
}
<?php
// src/Application/UseCase/CompleteLessonUseCase.php
declare(strict_types=1);
namespace App\Application\UseCase;
use Doctrine\DBAL\Connection;
use Symfony\Component\Uid\Uuid;
final class CompleteLessonUseCase
{
public function __construct(
private readonly Connection $connection,
) {}
public function execute(int $userId, int $lessonId): void
{
$this->connection->transactional(function (Connection $tx) use ($userId, $lessonId): void {
// 1. Бизнес-логика: обновляем прогресс
$tx->executeStatement(
'UPDATE lesson_progress SET completed = true, completed_at = now()
WHERE user_id = :user_id AND lesson_id = :lesson_id',
['user_id' => $userId, 'lesson_id' => $lessonId],
);
// 2. Записываем событие в outbox (в той же транзакции!)
$tx->insert('outbox', [
'event_id' => Uuid::v4()->toRfc4122(),
'event_type' => 'lesson.completed',
'payload' => json_encode([
'user_id' => $userId,
'lesson_id' => $lessonId,
], JSON_THROW_ON_ERROR),
]);
});
// Коммит автоматический: и прогресс, и событие сохранены атомарно
}
}
Обратите внимание: ни одной строки кода, связанной с Kafka, NATS или любым брокером. Use-case знает только о БД. Отправка - забота другого компонента.
Outbox Poller: горутина-отправитель
Отдельный процесс периодически читает неотправленные события из outbox и публикует их в брокер:
// Poller досылает события из outbox в брокер.
type Poller struct {
db *sql.DB
publisher Publisher
batchSize int
}
func NewPoller(db *sql.DB, publisher Publisher, batchSize int) *Poller {
return &Poller{db: db, publisher: publisher, batchSize: batchSize}
}
// outboxEvent - строка, вычитанная из outbox.
type outboxEvent struct {
id int64
eventID string
eventType string
payload []byte
}
// PollBatch публикует одну пачку событий.
//
// Всё - в одной транзакции, и это не небрежность. `FOR UPDATE SKIP LOCKED`
// держит строки занятыми **до конца транзакции**: закоммить её сразу после
// выборки, и второй экземпляр поллера тут же увидит те же события. Тогда
// блокировка защищала бы ровно на время чтения, то есть ни от чего.
//
// Порядок шагов внутри:
//
// 1. пачка целиком считывается и `rows` закрывается, и только потом идут
// публикации. Запрос к базе внутри цикла `rows.Next()` занял бы второе
// соединение (первое держит открытый набор строк), а на пуле из одного
// это тупик: процесс просто встанет, без ошибки;
// 2. публикация в брокер - внешний эффект, транзакция его не откатит.
// Поэтому упасть между `Publish` и отметкой значит отправить событие
// второй раз на следующем круге. Это at-least-once, и именно поэтому
// подписчик обязан быть идемпотентным - см. inbox.
func (p *Poller) PollBatch(ctx context.Context) error {
tx, err := p.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
events, err := claimBatch(ctx, tx, p.batchSize)
if err != nil {
return err
}
if len(events) == 0 {
return tx.Commit()
}
for _, e := range events {
if err := p.publisher.Publish(ctx, e.eventType, e.payload); err != nil {
// Публикация не удалась: считаем попытку и идём дальше. Событие
// останется неопубликованным и попадёт в следующую пачку.
if _, uErr := tx.ExecContext(ctx,
`UPDATE outbox SET attempts = attempts + 1 WHERE id = $1`, e.id,
); uErr != nil {
return fmt.Errorf("bump attempts for %s: %w", e.eventID, uErr)
}
continue
}
if _, err := tx.ExecContext(ctx,
`UPDATE outbox SET published_at = now() WHERE id = $1`, e.id,
); err != nil {
return fmt.Errorf("mark published %s: %w", e.eventID, err)
}
}
return tx.Commit()
}
// claimBatch забирает пачку неопубликованных событий, блокируя строки
// от других экземпляров поллера. Транзакцию не коммитит - её жизнью
// управляет вызывающий, и от этого зависит смысл блокировки.
func claimBatch(ctx context.Context, tx *sql.Tx, batchSize int) ([]outboxEvent, error) {
rows, err := tx.QueryContext(ctx,
`SELECT id, event_id, event_type, payload
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
FOR UPDATE SKIP LOCKED
LIMIT $1`,
batchSize,
)
if err != nil {
return nil, fmt.Errorf("query outbox: %w", err)
}
defer rows.Close()
var events []outboxEvent
for rows.Next() {
var e outboxEvent
if err := rows.Scan(&e.id, &e.eventID, &e.eventType, &e.payload); err != nil {
return nil, fmt.Errorf("scan row: %w", err)
}
events = append(events, e)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate rows: %w", err)
}
return events, nil
}
<?php
// src/Application/Cron/OutboxPollCommand.php
declare(strict_types=1);
namespace App\Application\Cron;
use App\Application\Port\EventPublisherPort;
use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;
use Symfony\Component\Console\Attribute\AsCommand;
use Symfony\Component\Console\Command\Command;
use Symfony\Component\Console\Input\InputInterface;
use Symfony\Component\Console\Output\OutputInterface;
use Throwable;
// В PHP-FPM долгоживущий цикл не подходит - используй один из двух путей:
// 1) Symfony Messenger worker (`bin/console messenger:consume outbox`) -
// отдельный процесс с встроенным циклом и graceful shutdown.
// 2) Однопроходная команда из cron каждые N секунд (вариант ниже).
#[AsCommand(name: 'app:outbox-poll')]
final class OutboxPollCommand extends Command
{
private const BATCH_SIZE = 100;
public function __construct(
private readonly Connection $connection,
private readonly EventPublisherPort $publisher,
private readonly LoggerInterface $logger,
) {
parent::__construct();
}
protected function execute(InputInterface $input, OutputInterface $output): int
{
// FOR UPDATE SKIP LOCKED - несколько параллельных воркеров не пересекаются
$rows = $this->connection->fetchAllAssociative(
'SELECT id, event_id, event_type, payload
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
LIMIT :lim
FOR UPDATE SKIP LOCKED',
['lim' => self::BATCH_SIZE],
);
foreach ($rows as $row) {
try {
$this->publisher->publishRaw($row['event_type'], (string) $row['payload']);
$this->connection->executeStatement(
'UPDATE outbox SET published_at = now() WHERE id = :id',
['id' => $row['id']],
);
} catch (Throwable $e) {
$this->connection->executeStatement(
'UPDATE outbox SET attempts = attempts + 1 WHERE id = :id',
['id' => $row['id']],
);
$this->logger->error('publish failed', [
'event_id' => $row['event_id'],
'err' => $e->getMessage(),
]);
}
}
return Command::SUCCESS;
}
}
Типичный интервал - 100-500ms. Можно ускорить через LISTEN/NOTIFY в PostgreSQL: use-case после коммита шлёт NOTIFY outbox_new, а poller слушает канал и просыпается мгновенно. Поллер целиком, вместе с use-case, тремя подписчиками и тестами на FakeEventBus, собран в мини-проекте трека.
1. Запрос к базе внутри цикла rows.Next(). Пока набор строк открыт,
он держит соединение. UPDATE через db попросит второе - и на пуле
из одного соединения процесс встанет навсегда. Ни ошибки, ни паники: поллер
просто перестаёт двигаться. На машине разработчика с пулом по умолчанию всё
работает, поэтому такое доезжает до прода. Поэтому пачка сначала считывается
целиком, и только потом публикуется.
2. Выборка без блокировки строк. Раздел типичных ошибок ниже требует
SELECT ... FOR UPDATE SKIP LOCKED, а сам код его не брал - два экземпляра
поллера отправили бы одно событие дважды. Заодно важно, где заканчивается
транзакция: SKIP LOCKED держит строки занятыми до её конца, поэтому коммит
сразу после выборки снимает защиту - блокировка живёт ровно на время чтения,
то есть ни на что.
CDC как альтернатива polling
Change Data Capture (CDC) - более продвинутый подход. Вместо polling инструмент вроде Debezium читает WAL (Write-Ahead Log) PostgreSQL и стримит изменения в Kafka. Outbox-таблица по-прежнему нужна, но poller - нет.
Преимущества CDC: нулевая задержка, нет нагрузки от polling-запросов. Недостаток: дополнительная инфраструктура (Debezium, Kafka Connect). Для начала достаточно polling - его проще запустить и отладить. Когда дойдёт до CDC, готовый compose с Kafka Connect, конфиг Debezium-коннектора и разбор «poller или CDC» лежат в уроке Transactional Outbox с Kafka.
Inbox для подписчиков
Outbox гарантирует, что событие будет отправлено. Но брокер может доставить его подписчику несколько раз (at-least-once). Inbox-таблица на стороне подписчика решает эту проблему - мы подробно разобрали её в предыдущем уроке.
Producer не теряет события (outbox). Consumer не обрабатывает дважды (inbox). Вместе они дают «ровно один раз» на уровне эффектов - при том, что доставка остаётся at-least-once. Гарантию даёт не транспорт, а пара «транзакция + ключ дедупликации»; где именно проходит граница.
Что связка действительно работает, видно только по метрикам: сколько записей висит с published_at IS NULL и как давно, какой lag у подписчиков, что накопилось в DLQ - это урок про наблюдаемость. И у пары outbox + inbox есть предел: она надёжно доставляет одно событие. Когда бизнес-процесс - это три шага во внешних системах и каждый нужно уметь откатить, одной доставкой не обойтись, начинается территория саги.
Очистка outbox
Отправленные события (published_at IS NOT NULL) можно удалять через несколько дней:
// CleanupPublished удаляет отправленные события старше ttl.
func (p *OutboxPoller) CleanupPublished(ctx context.Context, ttl time.Duration) (int64, error) {
res, err := p.db.ExecContext(ctx,
`DELETE FROM outbox
WHERE published_at IS NOT NULL
AND published_at < $1`,
time.Now().Add(-ttl),
)
if err != nil {
return 0, err
}
return res.RowsAffected()
}
<?php
// src/Application/Cron/CleanupOutboxCommand.php
declare(strict_types=1);
namespace App\Application\Cron;
use Doctrine\DBAL\Connection;
use Symfony\Component\Console\Attribute\AsCommand;
use Symfony\Component\Console\Command\Command;
use Symfony\Component\Console\Input\InputInterface;
use Symfony\Component\Console\Output\OutputInterface;
use DateTimeImmutable;
// Запуск из cron раз в сутки: `bin/console app:cleanup-outbox`.
#[AsCommand(name: 'app:cleanup-outbox')]
final class CleanupOutboxCommand extends Command
{
public function __construct(
private readonly Connection $connection,
) {
parent::__construct();
}
protected function execute(InputInterface $input, OutputInterface $output): int
{
$ttlDays = 7;
$threshold = (new DateTimeImmutable())->modify(sprintf('-%d days', $ttlDays));
$affected = $this->connection->executeStatement(
'DELETE FROM outbox
WHERE published_at IS NOT NULL AND published_at < :threshold',
['threshold' => $threshold->format('Y-m-d H:i:s')],
);
$output->writeln(sprintf('deleted %d outbox rows', $affected));
return Command::SUCCESS;
}
}
Как это проверено
Утверждение «producer не теряет события» проверяется прогоном:
examples/reliability/outbox (go test ./...).
Что именно перебирается:
- use case - сбой на каждой операции (запись прогресса → запись
в outbox →
COMMIT). После каждой проверяется, что число записей прогресса равно числу событий в outbox. Расхождение в любую сторону запрещено: «прогресс есть, события нет» - потеря, «события есть, прогресса нет» - подписчики узнают о том, чего не было; - поллер - сбой на каждой операции публикации. Запрещённое состояние одно: событие отмечено опубликованным, а в брокер не ушло. Обратное (ушло, но не отмечено) допустимо - это повторная отправка на следующем круге, то есть at-least-once, и именно поэтому подписчик обязан быть идемпотентным;
- два поллера внахлёст - второй запускается, пока транзакция первого открыта. Каждое событие должно уйти один раз.
Рядом с рабочим кодом лежат три сломанных версии - без них тесты ничего не доказывали бы:
| Версия | Что происходит |
|---|---|
| публикация внутри транзакции | транзакция откатилась, а подписчики событие уже получили |
COMMIT, потом Publish | упали в щели - прогресс есть, события нет и досылать нечего: в outbox пусто |
UPDATE внутри цикла rows.Next() | на пуле из одного соединения процесс встаёт; тест проверяет не ошибку, а зависание |
Типичные ошибки
-
Publish внутри транзакции до commit - публикация прошла, транзакция откатилась → подписчики видят событие на «не было». Канонический сценарий dual-write problem. Outbox: пиши событие в таблицу внутри транзакции, отправляй брокеру после commit.
-
Два Outbox Poller-а параллельно - оба читают одну и ту же не-опубликованную запись, оба отправляют. Брокер получает дубль (а у тебя может и не быть idempotency у consumer-а). Поллер - single instance, через leader election или
SELECT ... FOR UPDATE SKIP LOCKED. -
SELECT ... FOR UPDATEбезSKIP LOCKED- несколько poller-инстансов блокируются на одних рядах вместо распределения. Использовать именноSKIP LOCKED(PostgreSQL 9.5+). -
Polling каждые 100ms на больших объёмах - БД получает 600 пустых запросов/мин. Используй адаптивный интервал (если нашёл события - 100ms, если пусто - увеличивай до 5s) или LISTEN/NOTIFY (PostgreSQL).
-
Outbox без TTL/cleanup - таблица растёт до миллионов строк,
SELECT WHERE published_at IS NULLтормозит. Удаляй отправленные старше N дней + индекс на(published_at, created_at)для быстрого поиска не-отправленных. -
CDC + ручной Outbox publisher одновременно - оба читают изменения, дубли × 2. Выбирай один механизм: либо outbox poller, либо Debezium/CDC; не оба.
-
Inbox без bunded growth - таблица
processed_eventsрастёт вечно. Через год запросы тормозят, индекс не влезает в память. Партиционирование по дням + cleanup старше retention брокера. -
Outbox poller без retry &
attempts-лимита - broker лежит, poller бесконечно пытается отправить, забивает logs. Считайattempts; при превышении - отдельный «failed» статус и алёрт. -
CQRS - Transactional Outbox: событие как часть транзакции - пошаговая реализация Outbox Poller на Go с кодом и тестами
Мини-задание
- Создай таблицу
outboxс полямиid,event_id,event_type,payload,created_at,published_at,attempts - Напиши use-case, который сохраняет бизнес-данные и событие в одной транзакции
- Реализуй OutboxPoller с интервалом 500ms и batch size 10
- Добавь горутину очистки отправленных событий старше 7 дней
- Напиши тест: после вызова use-case в таблице outbox должна появиться запись с
published_at IS NULL