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 получил сообщение;
- локальный эффект consumer зафиксирован;
- offset подтверждён после эффекта;
- внешний или производный бизнес-результат достигнут.
Не каждый этап обязан жить в одной системе мониторинга. Важно, чтобы у него были наблюдаемый исход и владелец. Если внешний вызов выполняется отдельным рабочим процессом через второй 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, состояние партиции, исход обработчика и бизнес-сверку. Трассировка ускоряет диагностику конкретного события, но не заменяет контроль полноты.
На ревью достаточно трёх правил. У каждого долговечного перехода есть метрика возраста и окончательного исхода. Высококардинальные идентификаторы живут в журналах и трассировках, а не в метках. Последний сигнал проверяет бизнес-инвариант независимо от кода доставки. При свежей телеметрии и успешно выполненной сверке зелёное состояние означает отсутствие известных нарушений, а не только доступный брокер.