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

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

Outbox - паттерн, который атомарно фиксирует бизнес-изменение и намерение опубликовать событие: если транзакция в БД прошла, событие точно есть что отправить. Без него возникает dual-write проблема.

Дальше событие досылает отдельный процесс, и вот здесь атомарности уже нет: доставка из outbox остаётся at-least-once и требует повторов плюс идемпотентного подписчика. Разбор границы - в event-driven.

Dual-Write Problem

Типичная ошибка - писать в БД и в брокер отдельно:

// ОПАСНО: dual write
func (uc *CompleteLessonUseCase) Execute(ctx context.Context, cmd CompleteLesson) error {
    // Шаг 1: пишем в БД
    err := uc.repo.MarkCompleted(ctx, cmd.UserID, cmd.LessonID)
    if err != nil {
        return err
    }

    // Шаг 2: публикуем событие
    err = uc.publisher.Publish(ctx, "lesson.completed", event)
    if err != nil {
        // БД обновлена, но событие не отправлено!
        // Откатить БД? А если откат тоже упадёт?
        return err
    }

    return nil
}
<?php
declare(strict_types=1);

// ОПАСНО: dual write
final class CompleteLessonUseCase
{
    public function __construct(
        private readonly ProgressRepository $repo,
        private readonly EventPublisher $publisher,
    ) {}

    public function execute(CompleteLesson $cmd): void
    {
        // Шаг 1: пишем в БД
        $this->repo->markCompleted($cmd->userId, $cmd->lessonId);

        // Шаг 2: публикуем событие
        try {
            $this->publisher->publish('lesson.completed', $cmd);
        } catch (Throwable $e) {
            // БД обновлена, но событие не отправлено!
            // Откатить БД? А если откат тоже упадёт?
            throw $e;
        }
    }
}
Что может пойти не так:

Сценарий 1: БД ✓, Publish ✗ → данные есть, событие потеряно
Сценарий 2: БД ✓, Publish ✓, но ack от брокера не дошёл → дубль
Сценарий 3: Процесс упал между шагами → несогласованность

Никакая комбинация try/catch не решит эту проблему.

Решение: Transactional Outbox

Событие записывается в outbox-таблицу в той же транзакции, что и бизнес-данные. Отдельный процесс (publisher) потом отправляет его в брокер.

Transactional Outbox: UPDATE и INSERT INTO outbox в одной транзакции, отдельный publisher отправляет в RabbitMQ

┌─────────────────────────────────────────┐
│ Одна транзакция PostgreSQL              │
│                                         │
│ UPDATE progress SET completed = true    │
│ INSERT INTO outbox (event_id, ...)      │
│                                         │
│ COMMIT → обе операции или ни одной     │
└─────────────────────────────────────────┘
         │
         ▼
┌──────────────────┐     ┌──────────────┐
│ Outbox Publisher │ ──→ │  RabbitMQ    │
│ (отдельный       │     │              │
│  процесс/горутина)│    └──────────────┘
└──────────────────┘

Outbox-таблица

CREATE TABLE outbox (
    id         BIGSERIAL PRIMARY KEY,
    event_id   TEXT NOT NULL UNIQUE,
    event_type TEXT NOT NULL,
    payload    JSONB NOT NULL,
    created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
    sent_at    TIMESTAMPTZ, -- NULL = не отправлено
    attempts   INT NOT NULL DEFAULT 0,
    last_error TEXT
);

CREATE INDEX idx_outbox_unsent ON outbox(created_at) WHERE sent_at IS NULL;

Запись в outbox из use case

type CompleteLessonUseCase struct {
    db *gorm.DB
}

func (uc *CompleteLessonUseCase) Execute(ctx context.Context, cmd CompleteLesson) error {
    return uc.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
        // Бизнес-операция
        if err := tx.Model(&Progress{}).
            Where("user_id = ? AND lesson_id = ?", cmd.UserID, cmd.LessonID).
            Update("completed", true).Error; err != nil {
            return fmt.Errorf("mark completed: %w", err)
        }

        // Событие в outbox - в той же транзакции
        event := OutboxEvent{
            EventID:   uuid.New().String(),
            EventType: "lesson.completed",
            Payload: map[string]any{
                "user_id":   cmd.UserID,
                "lesson_id": cmd.LessonID,
                "track_id":  cmd.TrackID,
            },
        }

        payload, _ := json.Marshal(event.Payload)
        if err := tx.Exec(
            `INSERT INTO outbox (event_id, event_type, payload) VALUES (?, ?, ?)`,
            event.EventID, event.EventType, payload,
        ).Error; err != nil {
            return fmt.Errorf("outbox insert: %w", err)
        }

        return nil
    })
}
<?php
declare(strict_types=1);

use Doctrine\DBAL\Connection;
use Symfony\Component\Uid\Uuid;

final class CompleteLessonUseCase
{
    public function __construct(
        private readonly Connection $db,
    ) {}

    public function execute(CompleteLesson $cmd): void
    {
        $this->db->transactional(function (Connection $tx) use ($cmd): void {
            // Бизнес-операция
            $tx->executeStatement(
                'UPDATE progress SET completed = TRUE
                 WHERE user_id = :uid AND lesson_id = :lid',
                ['uid' => $cmd->userId, 'lid' => $cmd->lessonId],
            );

            // Событие в outbox - в той же транзакции
            $payload = json_encode([
                'user_id'   => $cmd->userId,
                'lesson_id' => $cmd->lessonId,
                'track_id'  => $cmd->trackId,
            ], JSON_THROW_ON_ERROR);

            $tx->executeStatement(
                'INSERT INTO outbox (event_id, event_type, payload)
                 VALUES (:id, :type, :payload)',
                [
                    'id'      => Uuid::v4()->toRfc4122(),
                    'type'    => 'lesson.completed',
                    'payload' => $payload,
                ],
            );
        });
    }
}

Outbox Publisher: Polling

Самый простой способ - polling: периодически читаем неотправленные события:

type OutboxPublisher struct {
    db        *sql.DB
    publisher EventPublisher
    interval  time.Duration
}

func (p *OutboxPublisher) Run(ctx context.Context) error {
    ticker := time.NewTicker(p.interval)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return nil
        case <-ticker.C:
            if err := p.processBatch(ctx); err != nil {
                log.Printf("outbox batch error: %v", err)
            }
        }
    }
}

type outboxEvent struct {
    id        int64
    eventID   string
    eventType string
    payload   []byte
}

func (p *OutboxPublisher) processBatch(ctx context.Context) error {
    // Транзакция на весь батч: FOR UPDATE SKIP LOCKED держит строки
    // занятыми только до её конца - см. врезку ниже.
    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)
    if err != nil {
        return err
    }

    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, last_error = $1 WHERE id = $2`,
                err.Error(), e.id,
            ); uErr != nil {
                return fmt.Errorf("bump attempts for %s: %w", e.eventID, uErr)
            }
            continue
        }

        if _, err := tx.ExecContext(ctx,
            `UPDATE outbox SET sent_at = now() WHERE id = $1`, e.id,
        ); err != nil {
            return fmt.Errorf("mark sent %s: %w", e.eventID, err)
        }
    }

    return tx.Commit()
}

// claimBatch вычитывает пачку целиком и закрывает rows - только после этого
// можно делать другие запросы. Почему это обязательно - во врезке ниже.
func claimBatch(ctx context.Context, tx *sql.Tx) ([]outboxEvent, error) {
    rows, err := tx.QueryContext(ctx,
        `SELECT id, event_id, event_type, payload
         FROM outbox
         WHERE sent_at IS NULL AND attempts < 10
         ORDER BY created_at
         LIMIT 100
         FOR UPDATE SKIP LOCKED`,  // конкурентные publisher'ы не мешают
    )
    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)
    }
    return events, rows.Err()
}
<?php
declare(strict_types=1);

use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;

final class OutboxPublisher
{
    public function __construct(
        private readonly Connection $db,
        private readonly EventPublisher $publisher,
        private readonly LoggerInterface $logger,
        private readonly int $intervalSeconds,
    ) {}

    public function run(): void
    {
        while (true) {
            try {
                $this->processBatch();
            } catch (Throwable $e) {
                $this->logger->error('outbox batch error: {err}', ['err' => $e->getMessage()]);
            }
            sleep($this->intervalSeconds);
        }
    }

    private function processBatch(): void
    {
        // Транзакция на весь батч: FOR UPDATE SKIP LOCKED держит строки
        // занятыми только до её конца - см. врезку ниже.
        $this->db->transactional(function (Connection $tx): void {
            // fetchAllAssociative вычитывает пачку целиком - других запросов
            // при открытом курсоре не будет.
            $rows = $tx->fetchAllAssociative(
                'SELECT id, event_id, event_type, payload
                 FROM outbox
                 WHERE sent_at IS NULL AND attempts < 10
                 ORDER BY created_at
                 LIMIT 100
                 FOR UPDATE SKIP LOCKED'  // конкурентные publisher'ы не мешают
            );

            foreach ($rows as $row) {
                try {
                    $this->publisher->publish($row['event_type'], $row['payload']);
                    // Помечаем как отправленное
                    $tx->executeStatement(
                        'UPDATE outbox SET sent_at = NOW() WHERE id = :id',
                        ['id' => $row['id']],
                    );
                } catch (Throwable $e) {
                    // Запоминаем ошибку, увеличиваем счётчик. Событие
                    // останется неотправленным и попадёт в следующую пачку.
                    $tx->executeStatement(
                        'UPDATE outbox SET attempts = attempts + 1, last_error = :err WHERE id = :id',
                        ['err' => $e->getMessage(), 'id' => $row['id']],
                    );
                }
            }
        });
    }
}

В Symfony Messenger тот же эффект даёт doctrine транспорт + messenger:consume - таблица messenger_messages играет роль outbox, а worker сам делает FOR UPDATE SKIP LOCKED.

Конструкция позволяет нескольким publisher'ам работать параллельно: каждый берёт свою пачку, не блокируя остальных. Но сама по себе строчка в запросе ничего не гарантирует - нужны три вещи. В первой версии этого урока не выполнялись первые две, а в Go-версии ещё и третья.

1. Запрос обязан идти в транзакции. db.QueryContext вне транзакции выполняется в неявной, которая закрывается по завершении запроса - блокировка снимается сразу, и второй publisher видит те же строки. FOR UPDATE в таком запросе декоративен.

2. Транзакция обязана жить до отметки об отправке. Закоммить её сразу после выборки - тот же результат: блокировка защищала ровно на время чтения.

3. Пачку нужно вычитать до конца и закрыть rows перед другими запросами. Пока набор строк открыт, он держит соединение; UPDATE через db попросит второе - и на пуле из одного соединения процесс встанет навсегда, без ошибки и паники. На машине разработчика с пулом по умолчанию это работает, поэтому такой код спокойно доезжает до прода и ложится там под нагрузкой.

Проверяется всё это прогоном: examples/reliability/outbox перебирает сбои по каждой операции и запускает второй поллер, пока транзакция первого открыта. Сломанные версии лежат рядом - в том числе та, где тест проверяет не ошибку, а зависание.

Polling vs CDC (Change Data Capture)

Подход     Как работает                     Плюсы / Минусы
────────   ──────────────────────           ──────────────────────────
Polling    SELECT WHERE sent_at IS NULL     + Просто реализовать
           каждые N секунд - Задержка до interval
 - Нагрузка на БД

CDC        Debezium читает WAL             + Мгновенная доставка
           PostgreSQL                       + Нет нагрузки на таблицу
 - Сложный инфра-сетап
 - Нужен Kafka Connect

Polling - правильный выбор для старта. Переходи на CDC (Debezium), когда задержка polling'а станет проблемой или нагрузка на outbox-таблицу вырастет. Как это выглядит в деле - compose с Kafka Connect, конфиг коннектора и мониторинг replication slot - в уроке Transactional Outbox с Kafka.

Очистка outbox

Отправленные события нужно удалять, иначе таблица вырастет бесконечно:

// Удаляем отправленные события старше 7 дней
func (p *OutboxPublisher) Cleanup(ctx context.Context) error {
    result, err := p.db.ExecContext(ctx,
        `DELETE FROM outbox WHERE sent_at IS NOT NULL AND sent_at < now() - interval '7 days'`,
    )
    if err != nil {
        return err
    }
    rows, _ := result.RowsAffected()
    if rows > 0 {
        log.Printf("cleaned up %d outbox events", rows)
    }
    return nil
}
<?php
declare(strict_types=1);

use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;

final class OutboxCleanup
{
    public function __construct(
        private readonly Connection $db,
        private readonly LoggerInterface $logger,
    ) {}

    // Удаляем отправленные события старше 7 дней
    public function cleanup(): void
    {
        $deleted = $this->db->executeStatement(
            "DELETE FROM outbox
             WHERE sent_at IS NOT NULL
               AND sent_at < NOW() - INTERVAL '7 days'"
        );

        if ($deleted > 0) {
            $this->logger->info('cleaned up {n} outbox events', ['n' => $deleted]);
        }
    }
}
Outbox publisher может отправить одно сообщение дважды (упал после publish, но до UPDATE sent_at). Поэтому consumer на другой стороне **обязан** быть идемпотентным (inbox-таблица из прошлого урока).

Полная картина

Use Case                    Outbox Publisher              RabbitMQ
─────────                   ────────────────              ────────
BEGIN TX
  UPDATE progress
  INSERT INTO outbox
COMMIT
                            SELECT unsent
                            Publish(event)  ──────────→  Exchange
                            UPDATE sent_at               │
                                                         ▼
                                                    Consumer (с inbox)

Гарантии:

  1. Если транзакция откатилась - событие не попадёт в outbox
  2. Если publisher упал - событие останется в outbox и будет отправлено при следующем polling
  3. Если consumer получил дубль - inbox отфильтрует

Мини-задание

  • Создай таблицу outbox с индексом на неотправленные события
  • Перепиши use case: бизнес-операция + INSERT INTO outbox в одной транзакции
  • Напиши Outbox Publisher с polling каждые 5 секунд
  • Проверь: останови publisher, выполни use case 3 раза, запусти publisher - все 3 события должны уйти в RabbitMQ

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