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) потом отправляет его в брокер.
┌─────────────────────────────────────────┐
│ Одна транзакция 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.
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]);
}
}
}
Полная картина
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