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 PollerDebezium CDC
СложностьПросто (пара сотен строк Go)Сложнее (Kafka Connect + коннектор)
Задержка100ms - 5s (polling interval)< 100ms (WAL streaming)
Нагрузка на БДПериодические SELECTRepslot 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.

Смотри также

Мини-задание

  • Создай таблицу outbox с полями из урока, добавь индекс на WHERE sent_at IS NULL
  • Напиши use case: CompleteLesson сохраняет прогресс и событие в outbox в одной транзакции
  • Реализуй OutboxPoller с FOR UPDATE SKIP LOCKED и batch size 10
  • Проверь: убей poller на середине batch - убедись что при перезапуске он обработает оставшиеся
  • Добавь goroutine очистки sent_at старше 7 дней

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