При корректной конфигурации Kafka чаще ломается не запись в журнал, а переходы вокруг брокера: бизнес-данные сохранились, но событие не ушло; обработчик вызвал внешний сервис, но не успел подтвердить offset; невалидное сообщение заблокировало партицию; повторная доставка создала второй заказ. Больше всего страдает не брокер, а бизнес-процесс, состояние которого оказалось разрезано между несколькими системами.
Надёжность здесь — не параметр producer и не обещание exactly-once. Это согласованная модель записи, доставки, повторной обработки и восстановления. Если хотя бы один переход не определён, сбой определит его за нас.
Контекст и границы задачи
Речь об интеграционных событиях между сервисами: заказ изменил статус, магазин стал недоступен, бронирование подтверждено, склад принял операцию. Источник владеет своим состоянием в PostgreSQL, Kafka переносит факт изменения, получатели обновляют собственные проекции или запускают следующий шаг процесса.
У такой схемы есть четыре ограничения.
Во-первых, транзакция базы источника и публикация в Kafka не атомарны. Последовательность «сначала COMMIT, затем Produce» оставляет окно потери. Обратный порядок оставляет фантомное событие: потребители его увидели, а изменение в базе откатилось.
Во-вторых, at-least-once означает нормальные повторы, а не редкую аномалию. Producer повторяет отправку после неоднозначного ответа, relay падает между публикацией и отметкой, consumer завершается после побочного эффекта, но до фиксации offset. Проектировать только счастливый путь — значит делегировать семантику повтора случайности.
В-третьих, порядок существует только внутри партиции. Его ещё нужно сохранить правильным ключом сообщения. Глобального порядка между заказами обычно не требуется; порядок изменений одного заказа — часто требуется.
В-четвёртых, Kafka не делает атомарными внешние эффекты. Транзакции Kafka полезны, когда чтение, преобразование и запись остаются внутри Kafka. Они не объединяют в одну транзакцию PostgreSQL, HTTP-запрос и действие сторонней платформы.
Поэтому сначала фиксирую бизнес-инвариант: что недопустимо потерять, что можно повторить, где нужен порядок и какое запаздывание приемлемо. Настройки брокера появляются после этого, не до.
Рабочая модель доставки
Базовая схема для событий из транзакционного сервиса выглядит так:
flowchart LR
A[Изменение бизнес-состояния] --> T[Транзакция источника]
T --> B[(Бизнес-таблицы)]
T --> O[(Outbox)]
O --> R[Relay]
R -->|event_id и aggregate_id| K[Kafka]
K --> C[Consumer]
C --> U[Транзакция получателя]
U --> I[(Inbox / dedup)]
U --> D[(Проекция)]
C -->|ошибка| Q[Retry или quarantine]
Она не обещает единственную доставку. Она делает повтор наблюдаемым и безопасным, а потерю — обнаруживаемой.
Событие как контракт
У события должен быть стабильный event_id, идентификатор сущности для ключа партиции, тип, версия схемы и время возникновения. Полезно разделять occurred_at и время публикации: первое относится к бизнес-факту, второе помогает диагностировать задержку outbox.
event_id создаётся один раз при записи outbox и не меняется при повторной публикации. Если relay генерирует новый идентификатор на каждой попытке, consumer не отличит повтор от нового факта.
Событийный контракт должен описывать факт, а не удалённую команду с неявными ожиданиями. OrderStatusChanged проще повторить и воспроизвести, чем UpdateEverythingForOrder. Асинхронная команда допустима, но у неё должны быть адресат, владелец решения и явная семантика результата. Событие при этом не обязано содержать всю модель. Достаточный снимок уменьшает связанность, но увеличивает размер и риск утечки лишних данных; ссылка на источник делает consumer зависимым от его доступности. Это выбор для конкретного потока.
Версия нужна не ради аккуратности. Producer и consumer обновляются независимо, а старые сообщения остаются в журнале и могут вернуться при replay. Совместимое добавление поля обычно дешевле, чем ветвление топиков. Несовместимое изменение лучше выпускать как новую версию контракта с явным периодом миграции.
Transactional outbox на стороне источника
Бизнес-изменение и строка outbox пишутся одной транзакцией PostgreSQL. Relay публикует сохранённое намерение и отмечает результат. Если он завершится после подтверждения Kafka, но до этой отметки, событие уйдёт повторно — это ожидаемая ветка, которую поглощает consumer.
Здесь важен принцип: outbox закрывает разрыв у источника, но не создаёт exactly-once для всего потока. Конкуренция relay, подтверждения брокера, aggregate_version, обслуживание таблицы и предел её ёмкости разобраны отдельно в материале «Transactional outbox в Go с атомарной записью и управляемыми повторами».
Идемпотентный consumer
Consumer проверяет event_id в своей базе и применяет локальное изменение в той же транзакции. Уникальный индекс превращает гонку двух обработчиков в контролируемый исход, а не в двойной эффект.
func (h *Handler) Handle(ctx context.Context, msg Message) error {
event, err := decodeAndValidate(msg.Value)
if err != nil {
return Permanent(err)
}
return h.db.WithTx(ctx, func(tx Tx) error {
claimed, err := h.inbox.TryClaim(ctx, tx, event.ID)
if err != nil || !claimed {
return err // nil означает уже обработанное событие
}
return h.projection.Apply(ctx, tx, event)
})
}
Offset подтверждается только после успешного COMMIT. Если процесс упадёт раньше подтверждения, сообщение придёт снова, TryClaim увидит его в inbox и завершит обработку без второго изменения.
Inbox — не единственный способ. Иногда идемпотентность естественно выражается состоянием: условное обновление по версии агрегата, UPSERT по бизнес-ключу, переход статуса только из допустимого предыдущего состояния. Такой инвариант сильнее отдельной таблицы dedup, потому что защищает и другие пути записи. Но проверка вида «если статус уже такой, ничего не делать» опасна, если обработчик должен выполнить ещё несколько независимых эффектов.
Вызов внешнего HTTP API нельзя откатить вместе с inbox. Если принимающая сторона поддерживает ключ идемпотентности, передаю стабильный идентификатор операции. Если нет, выношу внешний вызов в отдельный надёжный шаг: локально записываю команду в outbox, а специализированный worker выполняет её с явной сверкой состояния. Иногда полная автоматизация невозможна — тогда нужен статус «исход неизвестен» и процедура примирения, а не бесконечный retry.
Ошибки нужно классифицировать
Все ошибки в один retry — типичная причина скрытой остановки потока.
Временные ошибки — недоступная база, тайм-аут сети, исчерпанный пул соединений — повторяются с ограниченным числом быстрых попыток и увеличивающейся задержкой. Невалидная схема, неизвестная версия, нарушенный бизнес-инвариант и отсутствие обязательного поля от ожидания не исправятся. Их нужно выводить из автоматического повтора вместе с исходным payload, заголовками, причиной и координатами сообщения. DLQ здесь — транспортное место хранения, а quarantine — операционное состояние с владельцем и решением о возврате.
DLQ — не мусорная корзина. Для неё нужны владелец, оповещение, срок разбора и безопасный replay. Перед возвратом сообщение либо исправляется детерминированным преобразованием, либо устраняется причина в consumer. Простая кнопка «переотправить всё» воспроизводит аварию.
Если порядок по ключу критичен, перенос одного события в retry-топик может пропустить вперёд следующие изменения той же сущности. Тогда consumer должен приостановить партицию или хранить последовательность версий и отложить более новые события. Это снижает пропускную способность, но сохраняет инвариант. Если порядок не критичен, отдельный retry-топик освобождает основную партицию. Единой политики для всех потоков нет.
Как эта модель ломается
Неверный ключ партиции
Случайный ключ равномерно распределяет нагрузку, но перемешивает изменения одной сущности. Константный ключ сохраняет порядок ценой одной горячей партиции. Обычно ключом становится aggregate_id: он удерживает локальный порядок и даёт приемлемое распределение. Горячие сущности всё равно нужно видеть отдельно — Kafka не устранит перекос предметной области.
Слишком долгая обработка
Consumer делает сетевые вызовы, обработка превышает допустимое время, начинается rebalance, а сообщение получает другой участник группы. Результат — повторы и скачки lag. Лечат не увеличением всех тайм-аутов наугад, а разделением чтения и долгой работы, ограничением параллелизма, корректной настройкой poll loop и отказом от неограниченных внешних вызовов внутри обработчика.
Формально успешная, но неполная обработка
Обработчик проглотил ошибку, записал лог и вернул nil. Offset зафиксирован, событие потеряно для бизнес-процесса. Ошибка должна менять управляющий поток: retry, quarantine или явный компенсирующий статус. Лог без состояния восстановления — только свидетельство потери.
Dedup без срока жизни
Inbox растёт, индексы дорожают, обслуживание таблицы замедляется. Удалять записи можно только после расчёта максимального окна повторной доставки и replay. Если события могут переигрываться за год, семидневный dedup не защищает. Иногда проще сохранить компактный идентификатор навсегда, иногда — отделить обычную обработку от контролируемого исторического replay.
Эволюция схемы без проверки
Новый producer отправляет поле с другой семантикой, старый consumer принимает payload, но применяет неверное решение. Синтаксическая совместимость не равна смысловой. Контрактные тесты должны проверять примеры старых и новых версий, а rollout — учитывать одновременно работающие версии сервисов.
Что измерять
Один consumer lag не отвечает на вопрос, проходит ли бизнес-событие целиком. Я обычно собираю несколько слоёв сигналов.
- Для outbox: число ожидающих строк, возраст самой старой, частоту ошибок публикации и время от
occurred_atдо подтверждения relay. - Для Kafka: lag по партициям, возраст последнего обработанного события, частоту rebalance, ошибки produce/commit и перекос нагрузки между партициями.
- Для consumer: длительность обработки по типу события, долю повторов по inbox, число временных и постоянных ошибок, глубину retry и возраст старейшего сообщения в DLQ.
- Для бизнес-потока: расхождения между состояниями источника и проекции, число зависших переходов, время прохождения процесса от исходного факта до наблюдаемого результата.
Счётчик DLQ без возраста вводит в заблуждение: десять старых разобранных сообщений и одно новое критичное выглядят почти одинаково. Средняя задержка скрывает хвосты, поэтому нужны перцентили и максимум возраста. А нулевой lag при постоянно падающем обработчике возможен, если он ошибочно подтверждает offsets.
Технические метрики дополняю периодической сверкой. Например, источник публикует контрольные границы или отдельный процесс сравнивает агрегированные количества и версии. Это не замена корректной доставке, а способ обнаружить ошибку в самой реализации надёжности.
Когда Kafka и эта схема не нужны
Kafka оправдана, когда есть независимые потребители, журнал для replay, поток событий и требование переживать временную недоступность участников. Она добавляет стоимость: контракты, партиционирование, наблюдаемость, эксплуатацию consumer groups, DLQ и процедуры восстановления.
Для одного синхронного действия, где вызывающей стороне нужен немедленный ответ, идемпотентный HTTP или gRPC часто яснее. Для редкого фонового задания достаточно таблицы заданий в PostgreSQL. Для очереди команд с гибкой маршрутизацией и подтверждением каждого задания могут лучше подойти RabbitMQ или NATS — выбор зависит от требований, а не от привычки команды.
Не стоит внедрять outbox автоматически в каждый сервис. Если данные уже рождаются в Kafka и весь конвейер остаётся внутри неё, транзакционная обработка Kafka может дать более простую модель. Если событие допустимо восстановить из источника без потери смысла, периодическая синхронизация может быть дешевле. Если команда не готова владеть replay и DLQ, новый брокер лишь сделает отказ менее заметным.
Строгий глобальный порядок, мгновенная согласованность нескольких баз и ровно один внешний эффект Kafka сама не обеспечивает. Требования нужно либо ослабить, либо строить отдельный протокол координации — с соответствующей ценой доступности и сложности.
Что унести в ревью
Надёжная интеграция начинается с описания неоднозначных исходов. Между любыми двумя системами нужно ответить: что произойдёт после сбоя до подтверждения, после эффекта, но до фиксации, и при повторе через длительное время.
Практический минимум для ревью: бизнес-изменение и outbox атомарны; event_id стабилен; consumer идемпотентен на уровне локального инварианта; offset фиксируется после результата; временные и постоянные ошибки разведены; DLQ имеет владельца и процедуру replay; порядок обеспечивается только там, где он действительно нужен; технические метрики подкреплены бизнес-сверкой.
Exactly-once полезно обсуждать только после указания границ транзакции. В распределённой интеграции важнее не исключить каждый повтор, а сделать каждый допустимый повтор безопасным, каждый разрыв обнаруживаемым, а восстановление — штатной операцией.