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) потом отправляет его в брокер.
┌─────────────────────────────────────────┐
│ Одна транзакция 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.
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]);
}
}
}
Полная картина
Use Case Outbox Publisher RabbitMQ
───────── ──────────────── ────────
BEGIN TX
UPDATE progress
INSERT INTO outbox
COMMIT
SELECT unsent
Publish(event) ──────────→ Exchange
UPDATE sent_at │
▼
Consumer (с inbox)
Гарантии:
- Если транзакция откатилась - событие не попадёт в outbox
- Если publisher упал - событие останется в outbox и будет отправлено при следующем polling
- Если consumer получил дубль - inbox отфильтрует
Мини-задание
- Создай таблицу outbox с индексом на неотправленные события
- Перепиши use case: бизнес-операция + INSERT INTO outbox в одной транзакции
- Напиши Outbox Publisher с polling каждые 5 секунд
- Проверь: останови publisher, выполни use case 3 раза, запусти publisher - все 3 события должны уйти в RabbitMQ