Outbox + Inbox: событие и данные в одной транзакции

Outbox + Inbox: событие и данные в одной транзакции

Представьте: use-case завершает урок, записывает прогресс в PostgreSQL, а затем публикует событие lesson.completed в Kafka. Между этими двумя операциями приложение падает. Прогресс сохранён, но событие не отправлено. Achievement-сервис ничего не узнал. Analytics не записала метрику. Пользователь не получил бейдж.

Это называется dual-write problem - запись в два независимых хранилища без общей транзакции. Outbox закрывает именно эту щель, и важно понимать, что он закрывает только её.

Формулировку «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 - разные системы. Нельзя обернуть их в одну транзакцию. Какой бы порядок вы ни выбрали, есть окно для сбоя:

  1. Сначала БД, потом брокер. Если упали после коммита в БД - данные изменились, но событие не ушло. Подписчики не узнали.
  2. Сначала брокер, потом БД. Если упали после публикации - подписчики получили событие о данных, которых нет в БД. Ещё хуже.

Dual-write нельзя решить retry-ами. Нужен другой подход.

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

Идея: не отправляем событие в брокер напрямую. Вместо этого пишем его в специальную таблицу outbox в той же транзакции, что и бизнес-данные. Отдельный процесс (poller) читает outbox и отправляет события в брокер.

Если транзакция коммитится - и данные, и событие сохранены. Если откатывается - ничего не сохранено. Атомарность гарантирована средствами БД.

Без outbox - окно сбоя между DB commit и broker publish; с outbox - всё в одной транзакции, poller досылает в broker

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, собран в мини-проекте трека.

Первая версия этого поллера выглядела короче и содержала два дефекта - оба воспроизведены тестами в `examples/reliability/outbox`.

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-таблица на стороне подписчика решает эту проблему - мы подробно разобрали её в предыдущем уроке.

Inbox: consumer проверяет таблицу inbox по event_id; если запись есть - skip, если нет - обработать и записать в inbox в одной транзакции, потом подтвердить брокеру

Producer не теряет события (outbox). Consumer не обрабатывает дважды (inbox). Вместе они дают «ровно один раз» на уровне эффектов - при том, что доставка остаётся at-least-once. Гарантию даёт не транспорт, а пара «транзакция + ключ дедупликации»; где именно проходит граница.

Что связка действительно работает, видно только по метрикам: сколько записей висит с published_at IS NULL и как давно, какой lag у подписчиков, что накопилось в DLQ - это урок про наблюдаемость. И у пары outbox + inbox есть предел: она надёжно доставляет одно событие. Когда бизнес-процесс - это три шага во внешних системах и каждый нужно уметь откатить, одной доставкой не обойтись, начинается территория саги.

Outbox превращает dual-write в single-write. Транзакция БД - единственный источник истины. Poller может упасть и подняться - непубликованные события никуда не денутся, они в той же БД, что и бизнес-данные.

Очистка 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;
    }
}
Если в outbox есть записи с `published_at IS NULL` и `created_at` старше нескольких минут - что-то не так. Настройте алерт на количество неотправленных событий старше порога.

Как это проверено

Утверждение «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

Зарегистрируйтесь бесплатно, чтобы пройти квиз, вести прогресс.