Outbox + Inbox: как связать БД и брокер без потерь
Outbox + Inbox: как связать БД и брокер без потерь
Представьте: use-case завершает урок, записывает прогресс в PostgreSQL, а затем публикует событие lesson.completed в Kafka. Между этими двумя операциями приложение падает. Прогресс сохранён, но событие не отправлено. Achievement-сервис ничего не узнал. Analytics не записала метрику. Пользователь не получил бейдж.
Это называется dual-write problem - запись в два независимых хранилища без общей транзакции. Outbox-паттерн решает эту проблему элегантно и надёжно.
Dual-write problem
Две операции - запись в БД и публикация в брокер - не могут быть атомарными, потому что PostgreSQL и Kafka - разные системы. Нельзя обернуть их в одну транзакцию. Какой бы порядок вы ни выбрали, есть окно для сбоя:
- Сначала БД, потом брокер. Если упали после коммита в БД - данные изменились, но событие не ушло. Подписчики не узнали.
- Сначала брокер, потом БД. Если упали после публикации - подписчики получили событие о данных, которых нет в БД. Ещё хуже.
Dual-write нельзя решить retry-ами. Нужен другой подход.
Outbox: событие как часть транзакции
Идея: не отправляем событие в брокер напрямую. Вместо этого пишем его в специальную таблицу outbox в той же транзакции, что и бизнес-данные. Отдельный процесс (poller) читает outbox и отправляет события в брокер.
Если транзакция коммитится - и данные, и событие сохранены. Если откатывается - ничего не сохранено. Атомарность гарантирована средствами БД.
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-таблица на стороне подписчика решает эту проблему - мы подробно разобрали её в предыдущем уроке.
Producer не теряет события (outbox). Consumer не обрабатывает дважды (inbox). Вместе они дают exactly-once semantics на уровне приложения.
Очистка 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;
}
}
Типичные ошибки
-
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