ARTICLE / GO BACKEND

All articles

Observing an Asynchronous Flow from Consumer Lag to Business Invariants

How to connect Kafka, outbox, consumer, and domain metrics to distinguish temporary lag from event loss and silently stalled processing.

The full text is in Russian.

KafkaObservabilityDistributed SystemsArchitectureHighload

Kafka доступна, consumer lag равен нулю, DLQ пуста — а состояние в целевом сервисе не изменилось. Обработчик мог проглотить ошибку и подтвердить offset, событие могло застрять в outbox до брокера, а формально успешный consumer — применить только часть эффекта. Зелёная инфраструктура в таком потоке не подтверждает завершение бизнес-процесса.

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

Контекст и ограничения

Рассмотрим общий путь: сервис фиксирует бизнес-переход в PostgreSQL и outbox, relay публикует событие в Kafka, consumer обновляет свою проекцию или запускает внешний шаг. Между исходным фактом и наблюдаемым результатом несколько независимых фиксаций. У каждой свои повторы, задержки и неоднозначные исходы.

Consumer lag показывает только положение consumer group относительно журнала Kafka. Конкретный способ расчёта зависит от сборщика: он может сравнивать конечный offset партиции с подтверждённым смещением группы или с текущей позицией клиента. В обоих случаях метрика не знает, корректно ли выполнено бизнес-изменение после чтения сообщения.

Задержка тоже неоднозначна. Количество записей в lag не равно времени: одна и та же очередь при разной скорости входа и обработки даёт разный возраст. Время occurred_at, записанное источником, подходит для сквозной оценки только при согласованных часах. Внутри процесса длительность лучше измерять монотонным таймером, а между сервисами — хранить отдельные временные отметки и учитывать погрешность синхронизации часов.

Поэтому единицей наблюдения я считаю не топик и не экземпляр сервиса, а именованный бизнес-поток: например, «подтверждённый заказ стал доступен интеграционному адаптеру». У него есть источник, ожидаемый эффект, допустимое запаздывание и процедура восстановления.

Рабочая модель: сигналы на каждом переходе

Полезная схема связывает инфраструктурные метрики с состояниями одного потока.

flowchart LR
    A[Бизнес-переход] -->|created_at / event_id| O[(Outbox)]
    O -->|возраст и размер очереди| R[Relay]
    R -->|подтверждение / trace context| K[(Kafka)]
    K -->|lag и возраст сообщения| C[Consumer]
    C -->|исход и длительность| P[(Проекция / inbox)]
    P -->|ожидаемый эффект| B[Бизнес-инвариант]
    B -. периодическая сверка .-> A

Назвать этапы до выбора метрик

Для каждого потока я сначала фиксирую конечный автомат доставки. Минимальный набор этапов выглядит так:

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

event_id связывает журналы и трассировки, но не должен становиться меткой Prometheus: число временных рядов будет расти вместе с трафиком. Для метрик подходят ограниченные измерения — имя потока, тип события, consumer group, этап и класс исхода. Идентификаторы заказа, пользователя, магазина и самого события остаются в структурированном журнале или трассировке.

Источник: увидеть событие до Kafka

Outbox — первый наблюдаемый переход до Kafka. Здесь достаточно связать возраст старейшей неопубликованной записи с движением источника: если бизнес-переходы продолжаются, а published_at не появляется, остановка находится до брокера. Если outbox пуста и источник молчит, тревога по отсутствию публикаций будет шумом. Метрики попыток, ошибок и ёмкости подробнее разобраны в заметке Transactional outbox в Go.

Kafka и consumer: разделить очередь и обработку

Для Kafka полезны lag по партициям, скорость чтения, частота rebalance и ошибки фиксации смещений. Java consumer Kafka 4.2 экспортирует records-lag-max от текущей позиции и last-poll-seconds-ago — время с последнего вызова poll(). Текущая позиция может продвинуться до завершения обработчика; внешний сборщик часто считает расстояние от подтверждённого offset группы. Go-клиент или другой сборщик может экспортировать эквивалент под иным именем и с иной семантикой, поэтому источник lag указывают прямо на панели. Максимум выявляет наиболее отставшую партицию; горячий ключ — только одна из гипотез наряду с медленным обработчиком, зависимостью или назначением партиций. Сумма показывает общий объём накопленной работы, но скрывает локальный перекос.

Возраст тоже нужно назвать точно. occurred_at даёт бизнес-возраст при согласованных часах, timestamp записи Kafka — время нахождения в брокере, а возраст уже полученного, но не применённого сообщения отслеживает само приложение. Все три ближе к пользовательской задержке, чем количество offsets, но отвечают на разные вопросы. Если скорость поступления выше устойчивой скорости обработки, lag — уже не авария одного consumer, а прогноз исчерпания срока хранения. Здесь нужна оценка времени до потери доступного окна replay, а не только красная линия на графике.

Получение сообщения и его применение — разные операции. Kafka публикует клиентские метрики чтения, poll и rebalance, но бизнес-обработку должен измерять сервис. В OpenTelemetry Semantic Conventions 1.44 messaging.client.operation.duration измеряет операцию клиента, а messaging.process.duration — обработку. Счётчик messaging.client.consumed.messages корректно считает сообщения, доставленные приложению, но не подтверждает бизнес-эффект. Прикладные исходы applied, rejected и deferred нужно считать отдельно после фиксации результата; рядом нужны длительность, повторы, дубликаты и конфликты версий.

Бизнес-инвариант: проверить обещанный результат

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

Сверку можно выполнять по временным корзинам или диапазонам ключей, чтобы не сканировать всё состояние одновременно. Счётчик расхождений и возраст старейшего нарушения идут в метрики. Конкретные event_id остаются в отчёте сверки: по ним выбирают кандидатов для управляемого replay или ручного разбора. Безопасность повтора отдельно зависит от идемпотентности, текущего состояния и проверяемых предусловий.

Равенство агрегированных количеств — быстрый сигнал, но не доказательство: одна пропущенная запись и один дубликат взаимно сократятся. Для критичных потоков нужен выборочный или полный контроль идентификаторов и версий.

Трассировка объясняет, но не доказывает полноту

Контекст трассировки стоит передавать в заголовках сообщения. Это связывает создание, отправку, получение и обработку без попытки держать один span открытым всё время нахождения события в Kafka. При разветвлении и пакетной обработке зависимости честнее выражают span links: у одного span не может быть несколько родителей. Отдельные участки трассировки показывают локальную работу каждого этапа.

У трассировки есть пределы. Выборка может убрать именно нужное событие, трассировка может жить дольше срока хранения системы наблюдаемости, а пакетная обработка связывает одно действие с несколькими сообщениями. Поэтому отсутствие трассировки не означает потерю события. Она отвечает на вопрос «почему этот экземпляр шёл долго», а сверка проверяет полноту потока в объявленном окне.

OpenTelemetry Semantic Conventions 1.44 разделяют операции send, receive, process и settle, но правила для сообщений остаются в статусе Development. Внутреннюю модель этапов лучше держать устойчивой, а версию экспортируемых атрибутов — фиксировать и мигрировать осознанно. Иначе обновление инструмента переименует временные ряды быстрее, чем изменится сама система.

Типичные поломки и как их ловят

Нулевой lag после потерянного эффекта

Consumer подтвердил offset до фиксации транзакции PostgreSQL или вернул успех после ошибки. Подтверждённое смещение группы оказалось за этой записью: после перезапуска обычное чтение продолжится со следующей позиции, хотя эффект не зафиксирован. Это ловят счётчик окончательных исходов, сверка бизнес-инварианта и правило: offset подтверждается только после локального эффекта.

Lag растёт, но причина не в производительности обработчика

Частые rebalance, остановившийся poll, недоступная база или одна горячая партиция выглядят как одинаковая очередь. Диагностика начинается со сравнения скорости входа и выхода, затем проверяет распределение по партициям, время последнего poll, ожидания зависимостей и насыщение пула. Увеличивать число consumer до выяснения причины опасно: оно может усилить конкуренцию за базу.

Повторы выдают себя за полезную работу

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

DLQ существует, но никто ею не владеет

Размер DLQ без возраста и статуса разбора мало полезен. Одно новое критичное сообщение важнее старых подтверждённых записей. Для карантина нужны причина с ограниченной кардинальностью, возраст старейшего сообщения, ответственный поток и результат последнего replay.

«Нет данных» превращается в ноль

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

Метки превращаются в журнал событий

Добавление event_id, order_id или текста ошибки в метки резко увеличивает кардинальность и стоимость запросов. Ошибки классифицируют стабильным кодом, динамический контекст оставляют в журналах и трассировках, а связь выполняют через идентификатор в записи, не в имени временного ряда.

Что измерять

Набор сигналов зависит от обещания конкретного потока, но на ревью я ожидаю четыре слоя:

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

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

Когда такой контур не стоит строить

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

Телеметрия не подходит для регуляторного аудита: метрики агрегируются, трассировки могут отбираться, а оба сигнала имеют ограниченный срок хранения. Для доказательства судьбы каждого события нужен долговечный журнал аудита с собственной моделью доступа и удаления.

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

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

Consumer lag отвечает на вопрос о положении группы в Kafka, но не о завершении процесса. Рабочий контур связывает возраст outbox, состояние партиции, исход обработчика и бизнес-сверку. Трассировка ускоряет диагностику конкретного события, но не заменяет контроль полноты.

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