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

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

Outbox - паттерн, который гарантирует: если транзакция в БД прошла, событие не потеряется. Без него возникает dual-write проблема. Подробнее об outbox/inbox - в 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)
            }
        }
    }
}

func (p *OutboxPublisher) processBatch(ctx context.Context) error {
    // Читаем пачку неотправленных
    rows, err := p.db.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 err
    }
    defer rows.Close()

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

        if err := p.publisher.Publish(ctx, eventType, payload); err != nil {
            // Запоминаем ошибку, увеличиваем счётчик
            p.db.ExecContext(ctx,
                `UPDATE outbox SET attempts = attempts + 1, last_error = $1 WHERE id = $2`,
                err.Error(), id,
            )
            continue
        }

        // Помечаем как отправленное
        p.db.ExecContext(ctx,
            `UPDATE outbox SET sent_at = now() WHERE id = $1`, id,
        )
    }
    return nil
}
<?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
    {
        // Читаем пачку неотправленных
        $rows = $this->db->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']);
                // Помечаем как отправленное
                $this->db->executeStatement(
                    'UPDATE outbox SET sent_at = NOW() WHERE id = :id',
                    ['id' => $row['id']],
                );
            } catch (Throwable $e) {
                // Запоминаем ошибку, увеличиваем счётчик
                $this->db->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.

Эта конструкция PostgreSQL позволяет нескольким publisher'ам работать параллельно: каждый берёт свою пачку, не блокируя остальных. Без неё - deadlock при конкурентном доступе.

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-таблицу вырастет.

Очистка 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

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