Outbox + Inbox: как связать БД и брокер без потерь

Outbox + Inbox: как связать БД и брокер без потерь

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

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

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 и публикует их в брокер:

// OutboxPoller читает outbox и публикует события в брокер.
type OutboxPoller struct {
    db        *sql.DB
    publisher EventPublisher
    interval  time.Duration
    batchSize int
}

// EventPublisher - порт для отправки событий (Kafka, NATS, in-memory).
type EventPublisher interface {
    Publish(ctx context.Context, eventType string, payload []byte) error
}

// Run запускает polling в бесконечном цикле. Останавливается по ctx.Done().
func (p *OutboxPoller) Run(ctx context.Context) {
    ticker := time.NewTicker(p.interval)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            slog.Info("outbox poller stopped")
            return
        case <-ticker.C:
            if err := p.pollBatch(ctx); err != nil {
                slog.ErrorContext(ctx, "outbox poll failed",
                    slog.String("err", err.Error()),
                )
            }
        }
    }
}

// pollBatch выбирает пачку неотправленных событий и публикует их.
func (p *OutboxPoller) pollBatch(ctx context.Context) error {
    rows, err := p.db.QueryContext(ctx,
        `SELECT id, event_id, event_type, payload
         FROM outbox
         WHERE published_at IS NULL
         ORDER BY created_at
         LIMIT $1`,
        p.batchSize,
    )
    if err != nil {
        return fmt.Errorf("query outbox: %w", err)
    }
    defer rows.Close()

    for rows.Next() {
        var id int64
        var eventID, eventType string
        var payload []byte

        if err := rows.Scan(&id, &eventID, &eventType, &payload); err != nil {
            return fmt.Errorf("scan row: %w", err)
        }

        // Публикуем в брокер
        if err := p.publisher.Publish(ctx, eventType, payload); err != nil {
            // Увеличиваем счётчик попыток, но не останавливаемся
            p.db.ExecContext(ctx,
                `UPDATE outbox SET attempts = attempts + 1 WHERE id = $1`, id)
            slog.ErrorContext(ctx, "publish failed",
                slog.String("event_id", eventID),
                slog.String("err", err.Error()),
            )
            continue
        }

        // Отмечаем как отправленное
        _, err := p.db.ExecContext(ctx,
            `UPDATE outbox SET published_at = now() WHERE id = $1`, id)
        if err != nil {
            slog.ErrorContext(ctx, "mark published failed",
                slog.String("event_id", eventID),
                slog.String("err", err.Error()),
            )
        }
    }
    return rows.Err()
}
<?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;

// В 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 слушает канал и просыпается мгновенно.

CDC как альтернатива polling

Change Data Capture (CDC) - более продвинутый подход. Вместо polling инструмент вроде Debezium читает WAL (Write-Ahead Log) PostgreSQL и стримит изменения в Kafka. Outbox-таблица по-прежнему нужна, но poller - нет.

Преимущества CDC: нулевая задержка, нет нагрузки от polling-запросов. Недостаток: дополнительная инфраструктура (Debezium, Kafka Connect). Для начала достаточно polling - его проще запустить и отладить.

Inbox для подписчиков

Outbox гарантирует, что событие будет отправлено. Но брокер может доставить его подписчику несколько раз (at-least-once). Inbox-таблица на стороне подписчика решает эту проблему - мы подробно разобрали её в предыдущем уроке.

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

Producer не теряет события (outbox). Consumer не обрабатывает дважды (inbox). Вместе они дают exactly-once semantics на уровне приложения.

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;

// Запуск из 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` старше нескольких минут - что-то не так. Настройте алерт на количество неотправленных событий старше порога.

Типичные ошибки

  • 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

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