ARTICLE / GO BACKEND

All articles

Transactional Outbox in Go with Atomic Writes and Controlled Retries

How to connect PostgreSQL changes to Kafka publishing, survive relay retries and failures, and treat at-least-once delivery as an explicit design constraint.

The full text is in Russian.

GoPostgreSQLKafkaDistributed SystemsIntegrations

Сервис сохранил заказ в PostgreSQL, но не опубликовал событие в Kafka. Для клиента операция завершилась, а следующий сервис о ней не узнал. Если поменять действия местами, возникает зеркальная поломка: событие уже ушло, а транзакция откатилась. Между базой и брокером нет общей атомарной фиксации, поэтому «добавим retry» закрывает только часть разрыва.

Transactional outbox переносит этот разрыв в место, которым можно управлять. Он не даёт exactly-once и не отменяет идемпотентность. Зато делает намерение опубликовать событие частью той же транзакции, что и изменение бизнес-состояния.

Контекст и границы решения

Outbox нужен, когда PostgreSQL хранит исходное состояние, а Kafka распространяет факт его изменения. Для такого потока обычно важны четыре условия:

Последний пункт определяет стоимость решения. Глобальный порядок заставляет все изменения проходить через одно узкое место. На практике чаще нужен порядок событий одного заказа, брони или задания. Тогда aggregate_id можно использовать как ключ партиционирования, а разные агрегаты обрабатывать параллельно.

Outbox не решает семантику события. Если сервис публикует снимок случайного набора полей своей таблицы, потребители всё равно зависят от его внутренней модели. Событие должно фиксировать состоявшийся бизнес-факт, иметь стабильный идентификатор и явную версию схемы.

Рабочая модель: две атомарные границы вместо распределённой транзакции

Путь события состоит из трёх независимых частей: запись намерения, публикация и применение у потребителя.

flowchart LR
    A[Команда] --> B[Go-сервис]
    B -->|одна транзакция| C[(Бизнес-таблицы)]
    B -->|та же транзакция| D[(Outbox)]
    D --> E[Relay]
    E -->|at-least-once| F[Kafka]
    F --> G[Consumer]
    G -->|одна транзакция| H[(Inbox / бизнес-состояние)]
    H -->|после фиксации PostgreSQL| I[Фиксация смещения Kafka]

Записать состояние и намерение атомарно

Обработчик открывает транзакцию PostgreSQL, меняет бизнес-состояние и добавляет строку в outbox. Если любое действие не удалось, откатываются оба. Сетевого обращения к Kafka внутри этой транзакции нет.

Минимальная запись outbox содержит event_id, тип, версию схемы, идентификатор агрегата, полезную нагрузку и время создания. Поля для попыток публикации относятся к механике relay и не должны попадать в контракт события.

func (s *Service) confirm(ctx context.Context, orderID uuid.UUID) error {
	return s.db.WithTx(ctx, func(tx *sql.Tx) error {
		var aggregateVersion int64
		err := tx.QueryRowContext(ctx, `
			UPDATE orders
			SET status = 'confirmed', version = version + 1
			WHERE id = $1 AND status = 'pending'
			RETURNING version`, orderID).Scan(&aggregateVersion)
		if errors.Is(err, sql.ErrNoRows) {
			return ErrInvalidTransition
		}
		if err != nil {
			return fmt.Errorf("confirm order: %w", err)
		}

		payload, err := json.Marshal(OrderConfirmed{OrderID: orderID})
		if err != nil {
			return fmt.Errorf("marshal event: %w", err)
		}

		_, err = tx.ExecContext(ctx, `
			INSERT INTO outbox(
				event_id, aggregate_id, aggregate_version,
				event_type, schema_version, payload
			)
			VALUES ($1, $2, $3, 'order.confirmed', 1, $4)`,
			uuid.New(), orderID, aggregateVersion, payload)
		if err != nil {
			return fmt.Errorf("append outbox event: %w", err)
		}
		return nil
	})
}

aggregate_version описывает последовательность изменений сущности, а schema_version — формат контракта; заменять одно другим нельзя. Условный UPDATE ... RETURNING не даёт записать событие о переходе, которого не было. Условие перехода, новая версия и запись outbox описывают один факт.

Публиковать с повтором и ограниченной конкуренцией

Relay выбирает небольшую пачку неопубликованных строк, отправляет их в Kafka и отмечает результат. Для нескольких экземпляров есть три разные модели: держать транзакцию с FOR UPDATE SKIP LOCKED до ответа брокера; в короткой транзакции сохранить аренду записи, а затем отпустить блокировку; либо читать журнал изменений через CDC. Выбрать строки с SKIP LOCKED, закрыть транзакцию и отправлять без сохранённой аренды недостаточно — второй процесс увидит те же записи.

Если держать транзакцию и блокировки во время обращения к Kafka, модель проще, но медленный брокер занимает соединения и растит время блокировок. Если сначала пометить строки как взятые в работу и завершить транзакцию, появляется аренда: процесс может умереть после её фиксации, и запись должна вернуться в выборку после окончания срока. CDC убирает опрос таблицы из приложения, но добавляет инфраструктуру и операционные зависимости.

published_at ставится только после подтверждения брокера, а не после помещения сообщения во внутренний буфер клиента. Для критичного потока подтверждение всех доступных синхронных реплик (acks=all) должно сочетаться с осознанным min.insync.replicas; idempotent producer убирает часть дублей своих повторов. Но независимо от настроек остаётся окно: Kafka подтвердила сообщение, а отметка published_at не сохранилась. После перезапуска оно уйдёт ещё раз. Это штатный сценарий at-least-once.

Поглощать дубликаты там, где возникает эффект

Потребитель не должен рассчитывать, что дубликаты отфильтрует брокер. Надёжная граница — его локальная транзакция. Сначала INSERT ... ON CONFLICT (consumer, event_id) DO NOTHING RETURNING регистрирует событие. Если строка вернулась, обработчик применяет бизнес-изменение и проверяет ожидаемое число изменённых строк. Если не вернулась, это уже завершённый дубликат и эффект не повторяется. Конфликт версии после успешной регистрации требует отката всей транзакции: помечать неприменённое событие обработанным нельзя.

Автоматическая фиксация смещений consumer здесь отключена. Сначала завершается транзакция PostgreSQL, только затем фиксируется offset. Сбой между этими действиями даст повтор, который поглотит inbox; обратный порядок способен потерять обработку. При параллельной работе внутри одной партиции нельзя фиксировать более позднее смещение, пока более раннее не получило окончательный результат.

Если эффект — запрос во внешний API, локальная таблица processed_events атомарности не создаёт. Outbox надёжно сохраняет намерение сделать следующий вызов, но сам вызов всё равно требует стабильного idempotency key у получателя. Если получатель его не поддерживает, остаются сверка неизвестного исхода и компенсация.

Где схема ломается

Событие создано без фактического перехода

Без условного UPDATE, блокировки версии или проверки результата две конкурентные команды могут записать два взаимоисключающих факта. Уникальный event_id здесь не помогает: оба события технически разные. Ловится это инвариантами состояния, тестами конкурентных переходов и метрикой конфликтов версий.

Relay превращает базу в очередь без ограничений

Постоянный опрос, крупные пачки и отсутствие паузы нагружают индексы и пул соединений. При недоступной Kafka накопившиеся записи растут вместе с основной базой. Нужны ограничение размера пачки, backoff с разбросом, подходящий частичный индекс по неопубликованным строкам и политика удаления уже доставленных записей. Ёмкость диска задаёт конечное окно отказа: заранее нужны прогноз времени до заполнения и решение, когда включать обратное давление или ограничивать некритичные операции.

Удаление тоже не должно идти одной большой транзакцией. Иначе обслуживание outbox создаёт всплески WAL, удерживает блокировки и мешает autovacuum. Обычно записи удаляют или архивируют небольшими порциями после достаточного периода диагностики.

«Ядовитое» событие блокирует поток

Ошибка сериализации должна обнаруживаться до фиксации бизнес-транзакции. Но несовместимая схема, превышение лимита сообщения или неверная конфигурация топика проявятся позже. Бесконечный retry одной записи либо задержит весь агрегат, либо создаст горячий цикл.

Полезно разделять временные и постоянные ошибки. Временные получают backoff. Постоянные переводят запись в диагностируемое состояние с причиной, поднимают сигнал и требуют решения: исправить контракт, переиздать событие или явно отказаться от него. Автоматически пропускать запись опасно, если следующие события зависят от её порядка.

Порядок обещан шире, чем обеспечен

Порядок вставки в таблицу, чтения relay и доставки в Kafka — разные вещи. Один ключ сохраняет порядок только внутри одной партиции топика; несколько relay-процессов могут отправить соседние версии агрегата наоборот, а изменение числа партиций — поменять размещение ключа для будущих событий. Конфигурация producer тоже должна сохранять порядок при повторах.

aggregate_version обнаруживает пропуск, но не предотвращает его. Для строгой последовательности приходится исключать одновременный захват событий одного агрегата, сериализовать их публикацию либо хранить у потребителя более новые версии до восстановления пропущенной. Проверка версии может заменить inbox только для простой монотонной проекции. Если обработчик создаёт связанные записи, считает приращение или публикует следующий эффект, нужна отдельная защита каждого действия.

Сроки хранения не совпадают с окном повторов

Если processed_events очищается через неделю, а старое событие можно переиграть через месяц, удалённый идентификатор перестаёт защищать эффект. Нужно согласовать сроки outbox, Kafka retention, максимального consumer lag, ручного replay и дедупликации. Если consumer отстал дальше доступного журнала Kafka, сохранённый когда-то outbox уже не восстановит поток. Для compacted topic отдельно проверяют семантику: промежуточные факты можно удалять только тогда, когда их действительно заменяет последнее состояние.

Что измерять

Главная метрика relay — возраст самой старой неопубликованной записи. Размер накопившейся очереди полезен, но тысяча свежих событий и тысяча событий часовой давности означают разные проблемы. Рядом нужны:

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

В журнале достаточно event_id, aggregate_id, типа, попытки и этапа. Полную полезную нагрузку писать рискованно: она увеличивает объём, может содержать чувствительные данные и усложняет удаление.

Когда outbox не стоит применять

Outbox добавляет таблицу, relay, очистку, новые метрики и обязательную идемпотентность потребителей. Для операции внутри одного сервиса и одной базы обычная транзакция проще и надёжнее. Если событие не связано с изменением PostgreSQL, паттерн может не дать атомарной границы вообще.

Опрос таблицы не подходит, когда нужна очень малая и предсказуемая задержка публикации, а допустимый интервал опроса её превышает. Тогда стоит оценить CDC, но учитывать эксплуатацию коннектора, репликационные слоты и эволюцию схемы. Если команда не готова наблюдать накопившиеся записи и разбирать зависшие, outbox лишь сделает потерю события менее заметной.

Не стоит продавать этот подход как exactly-once. В Kafka-only потоке транзакция может атомарно записать выходные сообщения и offsets, а consumer с read_committed — не видеть незавершённый результат. PostgreSQL и внешний API в эту границу не входят. Для outbox-потока формулировка «at-least-once плюс идемпотентный эффект» точнее и сразу задаёт нужные тесты.

Вывод для архитектурного ревью

У рабочего outbox есть две атомарные границы: бизнес-изменение вместе с записью события у производителя и дедупликация вместе с эффектом у потребителя. Между ними повторы разрешены и наблюдаемы.

На ревью я проверяю не наличие таблицы outbox, а четыре свойства: событие соответствует реальному переходу, relay переживает сбой в любом месте, смещение Kafka фиксируется только после локального результата, а потребитель безопасно применяет одно событие повторно. Порядок при этом определён только там, где он нужен. Без этих условий outbox переносит проблему, но не закрывает её.