asman.malikov_ EN

· ~16 мин чтения

Transactional Outbox в Go: обе половины effectively-once

Outbox на продюсере и inbox-дедупликация на консьюмере в Go и Postgres: порядок событий, очистка таблиц и 1M+ платёжных событий в день без дублей.

gopostgresqlkafkaoutbox-patterndistributed-systemspayments

Каждая статья про transactional outbox заканчивается в одном и том же месте: relay читает таблицу, публикует в Kafka, помечает строку отправленной — готово. Это половина паттерна. На платёжно-биллинговой платформе, над которой я работаю, та половина, которая реально уберегла от двойных списаний с клиентов, живёт по другую сторону брокера — в консьюмере. Этот текст про обе половины: таблицу outbox, которую знают все, и inbox-дедупликацию, решения про порядок событий, очистку и операционную рутину, о которых почти никто не пишет. Всё это работает на 1M+ платёжных событий в день с нулём двойных обработок при ретраях и повторной доставке. Вот что для этого понадобилось.

TL;DR

  • Запись в БД и publish в Kafka — две независимые записи. Рано или поздно одна из них упадёт в одиночку, и оба порядка отказа стоят денег.
  • Половина продюсера: событие пишется в таблицу outbox в той же транзакции, что и изменение состояния. Поллящий publisher с FOR UPDATE SKIP LOCKED перекладывает его в Kafka.
  • SKIP LOCKED с несколькими репликами publisher-а молча уничтожает порядок событий внутри агрегата. Решите, важно ли это вам, до того как отмасштабируетесь до двух подов.
  • «Exactly-once» в Kafka заканчивается на брокере. Запись консьюмера в Postgres или вызов платёжного провайдера повторится при redelivery. Поэтому — половина консьюмера: таблица inbox с INSERT ... ON CONFLICT DO NOTHING в одной транзакции с side effect-ом.
  • Обе таблицы растут бесконечно. Очистка — это архитектурное решение первого дня, а не тушение пожара на шестом месяце.
  • Outbox — это хореография. Когда флоу требует многошагового отката, хореография заканчивается — про это статья о Temporal.

Баг двойной записи теряет деньги в обе стороны

Наивная реализация «сохрани платёж и расскажи всем о нём» — это две записи:

if err := repo.SavePayment(ctx, p); err != nil { // write 1: Postgres
    return err
}
if err := producer.Publish(ctx, evt); err != nil { // write 2: Kafka
    return err // ...and now what?
}

Транзакции, накрывающей Postgres и Kafka одновременно, не существует. Какой порядок ни выбери, зазор между двумя записями — это режим отказа, и я видел оба варианта в реальном коде, а не в учебнике:

Ordering A: commit DB first, publish second
┌─────────┐  1. COMMIT ok   ┌──────────┐
│ service │ ──────────────► │ Postgres │
└────┬────┘                 └──────────┘
     │ 2. publish
     ✖  crash / broker unavailable / pod OOM-killed

  [ Kafka ]     payment saved, event LOST
                downstream never hears about the money

Ordering B: publish first, commit DB second
┌─────────┐  1. publish ok  ┌──────────┐
│ service │ ──────────────► │  Kafka   │
└────┬────┘                 └──────────┘
     │ 2. COMMIT
     ✖  serialization failure / deadlock / crash

 [ Postgres ]   event announced, payment NEVER SAVED
                downstream acts on money that doesn't exist

Порядок A означает, что сервис нотификаций так и не отправит чек, а пайплайн отчётности недосчитается выручки. Порядок B хуже: консьюмер зачислит баланс за депозит, который ваша же база откатила. На платформе с дюжиной с лишним интеграций платёжных провайдеров их колбэки ретраятся агрессивно, так что окно не теоретическое — в него попадают.

Обернуть вторую запись в ретраи — не решение. Ретраи сужают окно, но не закрывают его: процесс может умереть между двумя записями. Единственное настоящее решение — превратить две записи в одну.

Половина продюсера: одна транзакция, одна таблица

Паттерн outbox превращает две записи в одну: событие ложится в таблицу в той же базе, в той же транзакции, что и изменение состояния. ACID-гарантии Postgres делают ту координацию, которую не даст никакой объём retry-логики.

CREATE TABLE outbox (
    id            bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    aggregate_id  text        NOT NULL,  -- becomes the Kafka partition key
    topic         text        NOT NULL,
    event_type    text        NOT NULL,
    payload       jsonb       NOT NULL,  -- envelope includes a producer-side event UUID
    created_at    timestamptz NOT NULL DEFAULT now(),
    published_at  timestamptz,
    retry_count   int         NOT NULL DEFAULT 0
);

-- partial index: the publisher only ever scans unpublished rows
CREATE INDEX outbox_unpublished_idx ON outbox (id)
    WHERE published_at IS NULL;

Две детали в этом DDL несут нагрузку. bigint, а не serial: при 1M+ событий в день sequence на int4 переполнится примерно за пять лет — ровно достаточно, чтобы все, кто об этом знал, успели уволиться. И retry_count: «отравленная» строка — payload, который не сериализуется, или сообщение больше лимита брокера — иначе заблокирует свой батч навсегда. Нужен запасной выход в dead-letter.

Доменная запись превращается в:

err := pgx.BeginFunc(ctx, pool, func(tx pgx.Tx) error {
    if err := savePayment(ctx, tx, p); err != nil {
        return err
    }
    _, err := tx.Exec(ctx, `
        INSERT INTO outbox (aggregate_id, topic, event_type, payload)
        VALUES ($1, $2, $3, $4)`,
        p.WalletID, "payments.events", "payment.captured", evt)
    return err
})

Либо существуют обе строки, либо ни одной. Баг двойной записи с этого момента структурно невозможен.

Publisher — это цикл. Наш поллит с FOR UPDATE SKIP LOCKED, чтобы несколько реплик не наступали друг другу на ноги:

func (p *Publisher) pollOnce(ctx context.Context) (int, error) {
    tx, err := p.pool.Begin(ctx)
    if err != nil {
        return 0, err
    }
    defer tx.Rollback(ctx)

    rows, err := tx.Query(ctx, `
        SELECT id, topic, aggregate_id, payload
        FROM outbox
        WHERE published_at IS NULL
        ORDER BY id
        LIMIT 100
        FOR UPDATE SKIP LOCKED`)
    if err != nil {
        return 0, err
    }
    batch, ids, err := collectRecords(rows) // []*kgo.Record keyed by aggregate_id
    if err != nil {
        return 0, err
    }
    if len(ids) == 0 {
        return 0, tx.Commit(ctx)
    }

    // franz-go; idempotent producer and acks=all are the defaults
    if err := p.client.ProduceSync(ctx, batch...).FirstErr(); err != nil {
        return 0, fmt.Errorf("produce batch: %w", err) // rollback: rows stay unpublished
    }

    if _, err := tx.Exec(ctx,
        `UPDATE outbox SET published_at = now() WHERE id = ANY($1)`, ids); err != nil {
        return 0, err
    }
    return len(ids), tx.Commit(ctx)
}

Порядок операций внутри этой функции — и есть весь смысл. Помечаем строки опубликованными после ack-а от брокера. Если publisher упал посреди батча, транзакция откатывается, блокировки снимаются, и другая реплика подбирает строки заново. Значит, одно и то же событие может быть отправлено дважды — и это нормально, потому что мы и не притворяемся, что это exactly-once delivery. Это at-least-once, намеренно. Дедупликация живёт ниже по течению.

Размер батча важнее, чем кажется. 100 строк на транзакцию — разумная отправная точка; рабочий диапазон — примерно 50–200. Возьмёте сильно больше — будете держать блокировки строк в длинной транзакции, которая заодно удерживает горизонт xmin и блокирует vacuum на таблице, которой vacuum отчаянно нужен (об этом ниже).

Если не хочется писать это руками — есть Forwarder в Watermill и библиотеки вроде oagudo/outbox, и они нормальные. Мы всё равно написали сами: на платёжной платформе в итоге хочется контроля над батчингом, порядком по шардам и гранулярностью метрик, которых generic-библиотеки не дают.

SKIP LOCKED даёт горизонтальный масштаб и тихо забирает порядок

Вот ловушка, которую туториалы пропускают. ORDER BY id внутри одного батча не даёт глобального порядка, как только у вас две реплики publisher-а. Реплика A хватает строки 1–100, реплика B — 101–200 (SKIP LOCKED означает, что B пропускает залоченные A строки), и B может добраться до Kafka первой. События 101–200 попадают в топик раньше 1–100. Даже с одной репликой значения sequence коммитятся не по порядку при конкурентных писателях — порядок по id не равен порядку коммитов. А если вы подумали про сортировку по created_at: расхождение часов между подами делает только хуже.

Важно ли это — зависит от того, какой порядок вам на самом деле нужен. Порядок между агрегатами почти никогда не важен. Порядок внутри агрегата — все события одного кошелька, одного платежа, по порядку — обычно важен. Варианты — от простого к самому честному:

  1. Одна реплика publisher-а. Просто, на практике «достаточно упорядоченно», и один хорошо настроенный поллер тянет 1M+ событий в день не запыхавшись. Настоящая причина, по которой люди запускают несколько relay-ев, редко в масштабе; обычно это страх SPOF, а Kubernetes, перезапускающий ваш единственный поллер за секунды, — как правило, приемлемый ответ.
  2. Шардировать полл по агрегату: каждая реплика поллит WHERE (hashtext(aggregate_id) & 2147483647) % $n = $i. Битовая маска не косметика: hashtext() возвращает знаковый int4 и примерно для половины входов отрицателен, а в Postgres -5 % 3 = -2 — предикат без маски навсегда оставил бы события «отрицательных» агрегатов неопубликованными. Порядок внутри агрегата сохранён, горизонтальный масштаб оставлен, сложность оплачена.
  3. Использовать CDC, который сохраняет порядок коммитов бесплатно.

Дальше вторая половина порядка: ключ партиционирования Kafka. Ключ каждой записи — aggregate_id. Один кошелёк → одна партиция → консьюмеры видят события этого кошелька по порядку. Возьмёте ключом ID события или оставите ключ nil — и round-robin размажет authorized/captured/failed одного платежа по партициям, и никакая хитрость на стороне консьюмера дёшево это не восстановит.

Поллинг vs Debezium: у меня в проде оба, и явного победителя нет

В нашем стеке ещё крутится Debezium CDC — из WAL PostgreSQL в Kafka, — так что это сравнение из эксплуатации обоих, а не из чтения двух README.

Обещание CDC настоящее: никаких поллящих запросов, долбящих таблицу, доставка в порядке коммитов бесплатно, почти нулевая задержка relay, а с Outbox Event Router SMT можно даже не помечать строки опубликованными. Когда он здоров — это более элегантная машина.

Кусается счёт за эксплуатацию. Коннектор Debezium держит логический replication slot, а replication slot — это обещание, что Postgres будет хранить WAL, пока коннектор его не прочитает. Коннектор, лежащий все выходные, означает, что WAL копится, пока не кончится диск, а Postgres без диска — инцидент куда хуже, чем несвежий outbox. Это митигируется через max_slot_wal_keep_size, heartbeat-ы для тихих таблиц и алерты на лаг в pg_replication_slots — но это новый класс отказов, в котором команда теперь обязана свободно ориентироваться, плюс сам Kafka Connect как дополнительная распределённая система, которую надо гонять и обновлять.

Поллящий publisher — скучный цикл на Go. Его режимы отказа — это режимы отказа Postgres, которые моя команда и так знает: bloat, борьба за блокировки, vacuum. Он добавляет постоянную нагрузку запросами и задержку в один poll-интервал (обычно это низкие сотни миллисекунд) и стоит вам порядка коммитов, если не шардировать.

Моё эмпирическое правило из эксплуатации обоих: если организация уже уверенно оперирует Kafka Connect и надо стримить много таблиц — Debezium окупается. Если outbox — ваш единственный кейс CDC, скучность поллящего publisher-а — это фича. Для событий, двигающих деньги, победила скука.

«Exactly-once» в Kafka заканчивается на границе брокера

Раз в несколько месяцев кто-нибудь читает про Kafka EOS и спрашивает, зачем мы возимся с дедупликацией. Затем, что граница гарантии уже, чем маркетинг:

  • Идемпотентный продюсер (по умолчанию с Kafka 3.0) дедуплицирует ретраи внутри одной сессии продюсера. Перезапустите publisher — а это ровно то, что происходит в сценарии «упал посреди батча» выше, — и новая сессия с радостью отправит ту же строку outbox ещё раз.
  • Транзакции Kafka делают consume-transform-produce атомарным внутри Kafka. В момент, когда консьюмер трогает что-то внешнее — UPDATE в Postgres, вызов платёжного провайдера, — вы вне транзакции. Консьюмер обработал сообщение, записал в свою БД, упал до коммита оффсета: на rebalance сообщение доставится повторно, и side effect выполнится дважды.

Честный контракт end-to-end: at-least-once delivery плюс дедупликация на стороне консьюмера равно effectively-once processing. Вторая половина этого предложения — ваша забота. Я видел, как консьюмеры дважды обрабатывали redelivery ровно по этому сценарию; это не редкость, это поведение по умолчанию для корректной at-least-once системы.

Половина консьюмера: таблица inbox, о которой никто не пишет

Фикс на стороне консьюмера симметричен продюсерскому и примерно того же объёма:

CREATE TABLE inbox (
    consumer     text        NOT NULL,   -- consumer group / logical consumer name
    message_id   uuid        NOT NULL,   -- producer-generated event ID from the envelope
    processed_at timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer, message_id)
);

Выбор ключа осознанный: UUID, сгенерированный продюсером и лежащий в envelope события. Не оффсет Kafka — оффсеты меняются, если топик пересоздали или события реплеятся из перестроенного топика, а это ровно те моменты, когда дедупликация нужна больше всего.

Консьюмер делает проверку дедупликации и side effect в одной транзакции базы:

func (c *Consumer) handle(ctx context.Context, rec *kgo.Record) error {
    var evt PaymentEvent
    if err := json.Unmarshal(rec.Value, &evt); err != nil {
        return c.toDLQ(ctx, rec, err) // poison message, don't block the partition
    }

    return pgx.BeginFunc(ctx, c.pool, func(tx pgx.Tx) error {
        tag, err := tx.Exec(ctx, `
            INSERT INTO inbox (consumer, message_id)
            VALUES ($1, $2)
            ON CONFLICT DO NOTHING`,
            c.name, evt.ID)
        if err != nil {
            return err
        }
        if tag.RowsAffected() == 0 {
            return nil // duplicate delivery: someone already processed it, commit no-op
        }
        return c.applyEffect(ctx, tx, evt) // credit balance, mark invoice, etc.
    })
}

Три вещи, которые здесь часто делают неправильно. Первое: check-then-insert (SELECT, потом INSERT) гонится при конкурентной повторной доставке; ON CONFLICT DO NOTHING плюс rows-affected — атомарная версия. Второе: вставка в inbox и side effect обязаны идти в одной транзакции — дедупликация в Redis или в отдельной транзакции заново открывает тот самый зазор, который вы закрыли на стороне продюсера. Третье: это покрывает только side effect-ы, живущие в собственном Postgres консьюмера. Если side effect — внешний вызов, скажем, запрос выплаты к провайдеру, inbox говорит, что вы начали, а не что закончили, и на самом вызове провайдера нужен idempotency key. Мы вешаем idempotency keys на каждую money-critical операцию ровно поэтому: inbox и ключ идемпотентности закрывают разные окна отказа, и нужны оба.

В сборе полный путь одного платёжного события:

┌────────────┐ BEGIN                             ┌──────────┐
│ billing    │  INSERT payment                   │ Postgres │
│ service    │  INSERT outbox ────────────────►  │ (svc A)  │
└────────────┘ COMMIT  (atomic)                  └────┬─────┘
                                                      │ poll (SKIP LOCKED)
                                                ┌─────▼─────┐
                                                │ publisher │  at-least-once
                                                └─────┬─────┘
                                                      │ produce, key=aggregate_id
                                                ┌─────▼─────┐
                                                │   Kafka   │  (dupes possible)
                                                └─────┬─────┘
                                                      │ consume (redelivery possible)
┌─────────────────────────────────────┐              │
│ consumer          BEGIN             │◄─────────────┘
│  INSERT inbox ON CONFLICT DO NOTHING│   0 rows → skip, commit
│  apply side effect                  │   1 row  → process
│                   COMMIT  (atomic)  │
└─────────────────────────────────────┘
        = effectively-once processing

Дубли разрешены везде в середине. Корректность обеспечивается на обоих краях, атомарно. Это весь дизайн в одном предложении, и именно поэтому мы можем говорить «ноль двойных обработок» с серьёзным лицом: не потому что дублей не бывает, а потому что они поглощаются.

Обе таблицы растут вечно, если не решить иначе в первый день

При 1M+ событий в день outbox прибавляет миллион строк в сутки, а inbox — миллион на каждого консьюмера. Я видел, во что превращается неограниченный outbox, и отказ коварнее, чем «кончился диск». Горячий запрос publisher-а — скан по частичному индексу на неопубликованных строках — деградирует по мере накопления мёртвых кортежей: index-only scan скатывается в проверки видимости по heap против миллионов мёртвых строк. Разбор Sadeq Dousti документирует финал: тот же запрос уходит с 0.13 мс до 18.5 секунд. А дефолтный autovacuum (autovacuum_vacuum_scale_factor = 0.2) даже не запустится, пока 20% огромной таблицы не станет мёртвыми.

Очевидный фикс — ночной DELETE FROM outbox WHERE published_at < now() - interval '7 days' — делает хуже. Массовый DELETE и есть то, что создаёт мёртвые кортежи, отравляющие горячий индекс. Реалистичные варианты:

  1. DELETE + агрессивные per-table настройки autovacuum. Работает на умеренных объёмах, если сильно опустить autovacuum_vacuum_scale_factor для этой таблицы и удалять мелкими батчами. Проще всего; с этого я бы начал.
  2. Партиционирование по времени (pg_partman + pg_cron), ретеншн через DROP PARTITION. Дроп партиции — операция над метаданными: мгновенно, без мёртвых кортежей, место освобождается сразу.
  3. Трюк Dousti: LIST-партиционировать outbox по статусу публикации, чтобы неопубликованные строки жили в маленькой горячей партиции, а опубликованные падали в партицию, которую периодически делают TRUNCATE.

Inbox требует того же обращения, с одним ограничением, которое упускают: ретеншн inbox должен превышать максимальное окно replay. Если топик хранит события 14 дней, а строки inbox вы держите 7, replay на десятый день проскочит мимо вашей дедупликации и дважды зачислит всё, чего коснётся. Свяжите эти два числа ретеншна в одном месте, с комментарием почему.

И один операционный шрам, который относится к обеим таблицам: любое изменение индексов на горячем outbox — только через CREATE INDEX CONCURRENTLY. Обычный CREATE INDEX берёт блокировку, останавливающую каждую запись, — а на этой таблице это значит остановить каждую доменную транзакцию сервиса, потому что вставка в outbox сидит внутри них. Это та миграционная ошибка, которая превращает «добавить индекс» в частичный простой. CONCURRENTLY может упасть и оставить за собой INVALID-индекс; его надо дропнуть и повторить, а не игнорировать.

Метрика, которая важна, — лаг от insert до publish

Consumer lag Kafka — метрика, которая у всех уже есть. Здесь она вторая по полезности. Дашборд, который реально ловит инциденты с outbox рано:

  • Глубина outbox: count(*) WHERE published_at IS NULL. Должна держаться около нуля.
  • Возраст самой старой неопубликованной строки: now() - min(created_at) WHERE published_at IS NULL. Это ваша настоящая end-to-end граница свежести. Вешайте алерт: заклинивший publisher появится здесь за минуты до того, как заметит хоть один пользователь.
  • Батч-метрики publisher-а: строк за полл, латентность produce, error rate и количество строк с retry_count > 0 (кандидаты в «отравленные»).
  • Здоровье Postgres по обеим таблицам: мёртвые кортежи из pg_stat_user_tables, время последнего autovacuum.
  • Если у вас Debezium: лаг replication slot-а в байтах, с алертом сильно ниже max_slot_wal_keep_size.

Каждая проблема с outbox, с которой я имел дело, объявляла о себе в этих числах до того, как становилась видна ниже по течению. Сервис отчётности на ClickHouse, потребляющий те же платёжные события из Kafka, — тоже приличная канарейка: аналитика, заметившая дыру, — это стыдно, но это бесплатный мониторинг, а держать аналитику на Kafka-стороне пайплайна означает, что её нагрузка никогда не касается транзакционного пути.

Где паттерн заканчивается

Честные пределы, потому что outbox — не универсальный ответ.

Не стройте это для модульного монолита с одной базой (транзакция уже всё покрывает), для по-настоящему fire-and-forget событий, где потерянное сообщение ничего не стоит, и до того, как вы приняли налог: каждое событие теперь стоит дополнительный insert, хоп задержки на relay и две таблицы, которые надо обслуживать. Чувствительный к задержкам request/response в нашем стеке остаётся на gRPC; outbox — для фактов, которые нельзя потерять, а не для вопросов, на которые нужен ответ сейчас.

Более глубокий предел — архитектурный. Outbox — это хореография: сервисы публикуют факты и реагируют на факты. Для fan-out это ровно то, что нужно: платёж случился, и нотификации, отчётность и реконсиляция независимо делают каждый своё. Это перестаёт быть правильным, когда флоу требует многошаговых записей через несколько сервисов с семантикой отката. Мы упёрлись в это с гео-распределёнными депозитными флоу: многошаговые, кросс-сервисные, с таймаутами и компенсациями при падении шага. Выстроенные хореографией поверх топиков, они превращаются в неявную стейт-машину, размазанную по сервисам: ни одно место не отвечает на вопрос «где застрял этот депозит и что откатывать, если упал шаг 4?». Прикручивать компенсации к событийной хореографии — значит случайно строить хрупкий оркестратор, по одной дедуп-таблице и одному таймаут-топику за раз.

Мы перевели эти флоу на Temporal SAGA, а outbox оставили для всего остального. Правило, которое я готов защищать: outbox — для фактов, оркестрация — для процессов. Как прошла та миграция — компенсации, таймеры, версионирование workflow в проде — отдельная статья, ссылка ниже.

Чек-лист перед выкаткой

  • Вставка в outbox и доменная запись идут в одной транзакции — проверьте тестом, который убивает процесс между ними.
  • bigint identity на id outbox. Частичный индекс на неопубликованных строках.
  • Publisher помечает строки опубликованными только после ack-а брокера. Уроните его в тесте посреди батча и посмотрите, как строки подбираются заново.
  • Kafka partition key = ID агрегата, выбранный под ваше реальное требование к порядку.
  • Решено: один publisher, hash-шардированные поллеры или CDC — и вы можете объяснить почему.
  • retry_count + dead-letter путь для «отравленных» строк, на стороне и продюсера, и консьюмера.
  • Inbox с ON CONFLICT DO NOTHING в одной транзакции с side effect-ом; ключ — сгенерированный продюсером ID события, а не оффсет.
  • Idempotency keys на внешних вызовах, двигающих деньги, — inbox их не покрывает.
  • Очистка спроектирована сейчас: партиционирование или тюнингованный autovacuum, для обеих таблиц; ретеншн inbox > окна replay.
  • Вся будущая работа с индексами на этих таблицах: CREATE INDEX CONCURRENTLY, с проверкой на оставшиеся INVALID.
  • Алерт на возраст самой старой неопубликованной строки и (если CDC) на лаг replication slot-а.
  • Тест на дублированную доставку в CI: доставьте каждое событие дважды, проверьте, что side effect выполнился один раз.

Смежные материалы

← Назад в блог