Transactional Outbox с Kafka: БД и брокер без потерь
Transactional Outbox с Kafka: БД и брокер без потерь
Проблема dual-write: нельзя атомарно записать в PostgreSQL и опубликовать в Kafka. Если сначала записать в БД, потом отправить в Kafka - при краше между ними событие потеряно. Если наоборот - возможен phantom-event.
Outbox pattern решает это: событие сначала сохраняется в таблицу outbox в той же транзакции, что и бизнес-данные. Отдельный процесс читает outbox и публикует в Kafka.
Теорию и RabbitMQ-вариант разобрали в event-driven - Outbox и RabbitMQ Outbox. Здесь - Kafka-специфика и CDC через Debezium.
Схема outbox-таблицы
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
event_id UUID NOT NULL DEFAULT gen_random_uuid(),
topic TEXT NOT NULL, -- Kafka топик
partition_key TEXT, -- для routing
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
sent_at TIMESTAMPTZ,
attempts INT NOT NULL DEFAULT 0
);
-- Индекс на неотправленные события (для polling)
CREATE INDEX idx_outbox_unsent ON outbox (created_at)
WHERE sent_at IS NULL;
Use Case: запись в транзакции
func (uc *CompleteLessonUseCase) Execute(ctx context.Context, userID, lessonID string) error {
return uc.db.WithTx(ctx, func(tx *sql.Tx) error {
// 1. Бизнес-данные
if err := uc.progressRepo.MarkCompleted(ctx, tx, userID, lessonID); err != nil {
return err
}
// 2. Событие в outbox (в той же транзакции)
event := map[string]any{
"event": "lesson.completed",
"userId": userID,
"lessonId": lessonID,
"at": time.Now().UTC(),
}
payload, _ := json.Marshal(event)
_, err := tx.ExecContext(ctx, `
INSERT INTO outbox (topic, partition_key, payload)
VALUES ($1, $2, $3)`,
"lesson-events", userID, payload,
)
return err
})
}
Если транзакция откатится - outbox запись тоже не появится. Согласованность гарантирована.
Outbox Poller
type OutboxPoller struct {
db *sql.DB
producer *kgo.Client
interval time.Duration
batchSize int
}
func (p *OutboxPoller) Run(ctx context.Context) {
ticker := time.NewTicker(p.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := p.poll(ctx); err != nil {
log.Printf("outbox poll error: %v", err)
}
}
}
}
func (p *OutboxPoller) poll(ctx context.Context) error {
rows, err := p.db.QueryContext(ctx, `
SELECT id, event_id, topic, partition_key, payload
FROM outbox
WHERE sent_at IS NULL AND attempts < 5
ORDER BY created_at
LIMIT $1
FOR UPDATE SKIP LOCKED`,
p.batchSize,
)
if err != nil {
return err
}
defer rows.Close()
type outboxRow struct {
id, eventID, topic, partitionKey string
payload []byte
}
var records []outboxRow
for rows.Next() {
var r outboxRow
rows.Scan(&r.id, &r.eventID, &r.topic, &r.partitionKey, &r.payload)
records = append(records, r)
}
rows.Close()
if len(records) == 0 {
return nil
}
// Публикуем в Kafka
var kafkaRecords []*kgo.Record
for _, r := range records {
kafkaRecords = append(kafkaRecords, &kgo.Record{
Topic: r.topic,
Key: []byte(r.partitionKey),
Value: r.payload,
Headers: []kgo.RecordHeader{
{Key: "event-id", Value: []byte(r.eventID)},
},
})
}
results := p.producer.ProduceSync(ctx, kafkaRecords...)
if err := results.FirstErr(); err != nil {
// Инкрементируем attempts при ошибке
ids := make([]string, len(records))
for i, r := range records { ids[i] = r.id }
p.db.ExecContext(ctx,
`UPDATE outbox SET attempts = attempts + 1 WHERE id = ANY($1)`,
pq.Array(ids))
return err
}
// Помечаем как отправленные
ids := make([]string, len(records))
for i, r := range records { ids[i] = r.id }
_, err = p.db.ExecContext(ctx,
`UPDATE outbox SET sent_at = NOW() WHERE id = ANY($1)`,
pq.Array(ids))
return err
}
FOR UPDATE SKIP LOCKED позволяет нескольким poller-инстансам работать параллельно без конфликтов.
CDC через Debezium: нет polling overhead
Debezium читает PostgreSQL WAL (Write-Ahead Log) и публикует изменения в Kafka. Для outbox-таблицы это даёт:
- Нулевая задержка (событие в Kafka через миллисекунды после commit)
- Нет нагрузки polling-запросами на БД
- Гарантированный порядок (WAL строго последовательный)
# docker-compose.yml
connect:
image: debezium/connect:2.6
ports:
- "8083:8083"
environment:
BOOTSTRAP_SERVERS: kafka:29092
GROUP_ID: "debezium"
CONFIG_STORAGE_TOPIC: "debezium_configs"
OFFSET_STORAGE_TOPIC: "debezium_offsets"
STATUS_STORAGE_TOPIC: "debezium_statuses"
depends_on:
- kafka
- postgres
Регистрируем коннектор через REST API Kafka Connect:
curl -X POST http://localhost:8083/connectors -H 'Content-Type: application/json' -d '{
"name": "outbox-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "app",
"database.password": "secret",
"database.dbname": "app",
"table.include.list": "public.outbox",
"transforms": "outbox",
"transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
"transforms.outbox.route.topic.replacement": "${routedByValue}",
"transforms.outbox.table.field.event.id": "event_id",
"transforms.outbox.table.field.event.type": "topic",
"transforms.outbox.table.field.event.key": "partition_key",
"transforms.outbox.table.field.event.payload": "payload"
}
}'
Debezium Outbox Event Router автоматически маршрутизирует записи из одной outbox-таблицы в нужные Kafka топики по полю topic.
Outbox Poller vs CDC: когда что
| Outbox Poller | Debezium CDC | |
|---|---|---|
| Сложность | Просто (пара сотен строк Go) | Сложнее (Kafka Connect + коннектор) |
| Задержка | 100ms - 5s (polling interval) | < 100ms (WAL streaming) |
| Нагрузка на БД | Периодические SELECT | Repslot overhead (~0.1%) |
| Масштабируемость | Несколько pollers с SKIP LOCKED | Один CDC процесс на БД |
| PostgreSQL replication slot | Не нужен | Нужен (следи за lag!) |
Для стартапа - Outbox Poller, проще в ops. Для высоконагруженного сервиса (>1K событий/сек) - CDC.
Cleanup старых записей
-- Удалять отправленные записи старше 7 дней
DELETE FROM outbox
WHERE sent_at IS NOT NULL
AND sent_at < NOW() - INTERVAL '7 days';
// Запускать как background goroutine
func (p *OutboxPoller) runCleanup(ctx context.Context) {
ticker := time.NewTicker(1 * time.Hour)
defer ticker.Stop()
for {
select {
case <-ctx.Done(): return
case <-ticker.C:
p.db.ExecContext(ctx, `
DELETE FROM outbox
WHERE sent_at IS NOT NULL
AND sent_at < NOW() - INTERVAL '7 days'`)
}
}
}
Inbox pattern на Kafka-стороне
Consumer обеспечивает идемпотентность через processed_events таблицу:
func processWithInbox(ctx context.Context, db *sql.DB, r *kgo.Record) error {
eventID := getHeader(r, "event-id")
return withTx(ctx, db, func(tx *sql.Tx) error {
// Проверяем не обработано ли уже
var exists bool
tx.QueryRowContext(ctx,
`SELECT EXISTS(SELECT 1 FROM processed_events WHERE event_id = $1)`,
eventID,
).Scan(&exists)
if exists {
return nil // дубль - пропускаем
}
// Бизнес-логика
if err := applyEvent(ctx, tx, r); err != nil {
return err
}
// Фиксируем обработку
_, err := tx.ExecContext(ctx,
`INSERT INTO processed_events (event_id, processed_at) VALUES ($1, NOW())`,
eventID)
return err
})
}
Типичные ошибки
Poller без SKIP LOCKED - несколько poller-инстансов блокируют одни записи, вместо распределения нагрузки. SKIP LOCKED решает это.
Неограниченный рост outbox-таблицы - без cleanup за год могут накопиться миллионы строк. Индекс WHERE sent_at IS NULL тоже растёт со временем. Cleanup обязателен.
Replication slot без мониторинга (CDC) - если Debezium лагает или падает, PostgreSQL не может вычистить WAL пока slot не прочитает изменения. Диск переполняется. Мониторь pg_replication_slots lag.
Смотри также
- Event-Driven - Outbox + Inbox - теория и anti-patterns
- Message Brokers - RabbitMQ Outbox - тот же паттерн с RabbitMQ
- SQL - Транзакции - BEGIN/COMMIT и atomicity
Мини-задание
- Создай таблицу
outboxс полями из урока, добавь индекс наWHERE sent_at IS NULL - Напиши use case:
CompleteLessonсохраняет прогресс и событие в outbox в одной транзакции - Реализуй
OutboxPollerсFOR UPDATE SKIP LOCKEDи batch size 10 - Проверь: убей poller на середине batch - убедись что при перезапуске он обработает оставшиеся
- Добавь goroutine очистки sent_at старше 7 дней