Каждая статья про 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: расхождение часов между подами делает только хуже.
Важно ли это — зависит от того, какой порядок вам на самом деле нужен. Порядок между агрегатами почти никогда не важен. Порядок внутри агрегата — все события одного кошелька, одного платежа, по порядку — обычно важен. Варианты — от простого к самому честному:
- Одна реплика publisher-а. Просто, на практике «достаточно упорядоченно», и один хорошо настроенный поллер тянет 1M+ событий в день не запыхавшись. Настоящая причина, по которой люди запускают несколько relay-ев, редко в масштабе; обычно это страх SPOF, а Kubernetes, перезапускающий ваш единственный поллер за секунды, — как правило, приемлемый ответ.
- Шардировать полл по агрегату: каждая реплика поллит
WHERE (hashtext(aggregate_id) & 2147483647) % $n = $i. Битовая маска не косметика:hashtext()возвращает знаковый int4 и примерно для половины входов отрицателен, а в Postgres-5 % 3 = -2— предикат без маски навсегда оставил бы события «отрицательных» агрегатов неопубликованными. Порядок внутри агрегата сохранён, горизонтальный масштаб оставлен, сложность оплачена. - Использовать 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 и есть то, что создаёт мёртвые кортежи, отравляющие горячий индекс. Реалистичные варианты:
- DELETE + агрессивные per-table настройки autovacuum. Работает на умеренных объёмах, если сильно опустить
autovacuum_vacuum_scale_factorдля этой таблицы и удалять мелкими батчами. Проще всего; с этого я бы начал. - Партиционирование по времени (
pg_partman+pg_cron), ретеншн черезDROP PARTITION. Дроп партиции — операция над метаданными: мгновенно, без мёртвых кортежей, место освобождается сразу. - Трюк 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 и доменная запись идут в одной транзакции — проверьте тестом, который убивает процесс между ними.
-
bigintidentity наidoutbox. Частичный индекс на неопубликованных строках. - 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 выполнился один раз.
Смежные материалы
- Как эта система заработала свои цифры: proof платёжных микросервисов — контекст платформы за заявлением про ноль двойных обработок.
- Когда хореография закончилась: Temporal SAGA для гео-распределённых депозитов — статья-близнец, продолжающая ровно там, где эта останавливается.