Consumer lag: что это и как с ним работать

Лаг консьюмера — расстояние между High Watermark и закоммиченным оффсетом. Почему он растёт и что писать в требованиях.

Consumer lag — это разница между тем, сколько записей уже доступно в партиции, и тем, до какого оффсета группа успела дойти. Формально: lag = High Watermark − committed offset. Это единственная метрика Kafka, на которую в проде смотрят каждый день.

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

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

В требования к интеграции лаг попадает тремя числами: целевой RPS, допустимое время обработки одного сообщения и максимально допустимый лаг — в записях или в секундах отставания. Отдельным пунктом — алерт: без него о проблеме узнают от бизнеса, а не от мониторинга.

Отдельно стоит помнить: лаг считается от High Watermark, а не от конца лога. Записи выше High Watermark ещё не реплицированы на все синхронные реплики и консьюмеру недоступны в принципе.

Разобрать это в симуляторе Интерактивный сценарий: настраиваешь кластер, ломаешь его и смотришь, что происходит

Частые вопросы

Лаг растёт, добавили консьюмеров — не помогло. Почему?

Скорее всего, консьюмеров уже столько же, сколько партиций. Лишние простаивают. Нужно ускорять обработку или увеличивать число партиций.

Какой лаг считать нормальным?

Тот, который стабилен и укладывается в бизнес-требование по задержке. Постоянный лаг в тысячу записей при обработке в 10 000 rps — это 0,1 секунды и обычно нормально.

Чем лаг грозит, кроме задержки?

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

Дальше по теме

Что такое партиция в Kafka acks в Kafka: 0, 1 и all ISR и min.insync.replicas Ребалансировка consumer group Retry-топик и DLQ в Kafka