MAATRIX / Блог / Kafka не удаляла старые сегменты и забила диск за выходные

Kafka не удаляла старые сегменты и забила диск за выходные

MAATRIX

В субботу в два часа ночи прилетел алерт: диск на одном из брокеров Kafka перевалил за 85% и продолжал расти линейно, без признаков остановки. Retention был настроен, cleanup.policy стоял правильный, а старые сегменты почему-то не удалялись. К понедельнику диск гарантированно заполнился бы под ноль. Ниже — как мы нашли причину и почему она не имела отношения ни к retention.ms, ни к репликации, ни к самому диску.

Что сломалось

Кластер — три брокера Kafka, обычная конфигурация: несколько топиков с событиями, retention на 7 дней (log.retention.hours=168), cleanup.policy=delete на всех «боевых» топиках, кроме служебных compacted-топиков вроде __consumer_offsets. В пятницу вечером выкатили новый сервис — форвардер аналитических событий, который читает данные из внешней системы и пишет их в Kafka с сохранением исходного времени события (это важно, вернёмся к этому позже).

В субботу ночью прометеевский алерт по node_filesystem_avail_bytes на одном из брокеров сработал на пороге в 85%. Дежурный посмотрел графики: место на разделе с данными Kafka таяло не рывками, а ровной прямой — по несколько гигабайт в час, без пауз. Топ-3 продюсеров по объёму не менялся, входящий трафик был в пределах нормы для выходных. То есть проблема была не «пишем слишком много», а «не удаляем то, что должны».

К утру воскресенья на месте оставалось меньше 10%, и стало ясно, что это не вопрос мониторинга, а вопрос устранения причины — иначе к понедельнику брокер просто перестанет принимать запись, а следом за ним могли посыпаться и остальные два узла кластера из-за перераспределения лидеров партиций.

Что видели в логах и метриках

Первым делом посмотрели, что физически лежит на диске:

du -sh /var/lib/kafka/data/* | sort -rh | head -20

Почти весь прирост давал один топик, а точнее — одна партиция: events-analytics-3. Внутри директории партиции лежали сегменты .log/.index/.timeindex с датами создания за последние несколько месяцев — притом что retention для топика был выставлен на 7 дней.

Проверили конфиг топика напрямую, чтобы исключить банальную опечатку:

kafka-configs.sh --bootstrap-server kafka-2:9092 \
  --describe --entity-type topics --entity-name events-analytics

Вывод подтвердил: retention.ms=604800000 (7 дней), cleanup.policy=delete. Никакого compact, никакого retention.ms=-1, ничего экзотического в overrides на уровне топика. Конфигурация была ровно такой, какой её ожидали увидеть.

Посмотрели логи самого брокера (server.log) за выходные — ни одной ошибки, связанной с LogCleaner или удалением сегментов. Планировщик ретеншна (kafka.log.LogManager) отрабатывал по расписанию (log.retention.check.interval.ms, по умолчанию 5 минут) без исключений — то есть процесс проверки запускался исправно, просто ничего не находил для удаления в этой партиции.

Проверили под-реплицированные партиции и состояние ISR:

kafka-topics.sh --bootstrap-server kafka-2:9092 \
  --describe --under-replicated-partitions

Список был пуст — с репликацией всё было в порядке, все три реплики партиции синхронны.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Версии, которые отбросили

По ходу разбора отбросили четыре гипотезы, каждую — с проверкой, а не «на глаз»:

  • Ретеншн настроен неправильно. Отбросили сразу — kafka-configs.sh --describe показал корректные retention.ms и cleanup.policy=delete и на уровне топика, и с учётом broker-defaults.
  • Кто-то по ошибке включил compaction. Проверили cleanup.policydelete, никаких следов compact в истории изменений конфигурации топика (смотрели через kafka-configs.sh --describe и changelog деплоя конфигов).
  • Репликация отстаёт и блокирует удаление. Под-реплицированных партиций не было, ISR полный на всех трёх брокерах. К тому же удаление сегментов у leader-партиции с cleanup.policy=delete в принципе не завязано на состояние реплик так, как это происходит с log compaction — это тоже сверили по документации, чтобы не тратить время на ложный след.
  • Диск забивают не данные Kafka, а что-то постороннее — логи демона, core-дампы, временные файлы. du -sh по корню раздела сразу показал, что 95%+ занятого места — это именно /var/lib/kafka/data, и внутри неё — конкретная партиция одного топика. Посторонний виновник отпал за одну команду.

После этого стало ясно: дело не в конфигурации и не в инфраструктуре вокруг, а в том, как Kafka сама решает, какой сегмент уже «протух» и годен к удалению.

Как нашли настоящую причину

Ключевой момент, который легко упустить: время жизни сегмента для time-based retention Kafka считает не по mtime файла на диске, а по максимальному timestamp записи внутри сегмента (значение хранится в индексе .timeindex и берётся из поля CreateTime, которое по умолчанию присылает клиент-продюсер). Формула по сути такая: «текущее время минус largest timestamp сегмента больше retention.ms — сегмент можно удалять». Если хотя бы одна запись в сегменте имеет заведомо будущий timestamp, вся эта разница обнуляется — сегмент выглядит для брокера «свежим», сколько бы реального времени ни прошло.

Чтобы проверить эту версию, вытащили содержимое одного из старых, но не удаляемых сегментов:

kafka-dump-log.sh --deep-iteration --print-data-log \
  --files /var/lib/kafka/data/events-analytics-3/00000000000012345678.log \
  | less

И нашли то, что искали: часть записей в сегменте имели нормальный CreateTime, а часть — timestamp, соответствующий дате на несколько десятков лет вперёд. Раз брокер доверял клиентскому времени (log.message.timestamp.type=CreateTime — значение по умолчанию), largest timestamp сегмента улетал в будущее вместе с этой записью, и сегмент переставал быть кандидатом на удаление — фактически навсегда, пока системные часы не «догонят» эту фиктивную дату.

Осталось понять, откуда взялись битые timestamp'ы. Виновником оказался как раз тот форвардер аналитических событий, который выкатили в пятницу. Он брал поле времени события из апстрим-системы и передавал его в Kafka как timestamp записи, ожидая миллисекунды — а апстрим отдавал микросекунды. Разница в три порядка превращала обычную дату в дату из XXI века вперёд на десятилетия. Проблема была не разовой: сервис писал так постоянно и на приличном объёме, поэтому каждый новый сегмент, начиная с пятничного деплоя, получал хотя бы одну «отравленную» запись — и переставал быть кандидатом на удаление сразу после закрытия. Именно поэтому диск рос стабильно всю ночь и весь день: удалялось всё меньше и меньше старых сегментов, а новые «неудаляемые» продолжали копиться.

Что изменили после инцидента

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

Немедленно. Отключили форвардер аналитики (тот самый пятничный деплой), чтобы прекратить приток новых «отравленных» сегментов, и временно снизили нагрузку на диск за счёт увеличения интервала между полными бэкапами других топиков — это дало запас по времени, а не решило проблему, но купило нам несколько часов без риска остановки записи.

На уровне топика. Для топиков, где точный клиентский timestamp не критичен для бизнес-логики (а events-analytics был именно таким — время события используется только в отчётах, где секунды точности роли не играют), переключили обработку времени на серверную:

kafka-configs.sh --bootstrap-server kafka-2:9092 \
  --alter --entity-type topics --entity-name events-analytics \
  --add-config log.message.timestamp.type=LogAppendTime

С LogAppendTime брокер сам проставляет время при получении записи, игнорируя то, что прислал клиент, — для retention это делает поведение предсказуемым независимо от того, что натворит очередной продюсер.

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

Исправили источник. У форвардера пофиксили единицы измерения времени (микросекунды → миллисекунды) и добавили валидацию: если timestamp события отклоняется от текущего времени больше чем на разумный порог, сервис логирует предупреждение и либо отбрасывает поле (давая брокеру проставить LogAppendTime), либо не публикует запись — в зависимости от топика.

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

log.message.timestamp.difference.max.ms=3600000

При превышении лимита брокер отклоняет запись с ошибкой на стороне продюсера — вместо того чтобы молча принять и «отравить» сегмент.

Как не наступить на те же грабли

Несколько практических выводов, которые пригодятся, даже если у вас Kafka на паре виртуалок, а не большой кластер:

  • Не доверяйте retention полностью клиентскому времени. Если продюсеров пишет несколько команд или сервисов, log.message.timestamp.type=LogAppendTime для «обычных» топиков — разумный дефолт. Клиентский CreateTime нужен точечно, там, где порядок событий во времени реально важен для бизнес-логики.
  • Держите retention.bytes как второй рубеж, а не только retention.ms. Он не спасёт от «отравленного» сегмента напрямую (Kafka всё равно ориентируется на комбинацию условий), но задаёт жёсткий потолок по объёму на партицию и не даст ситуации развиться до заполнения всего диска, пока вы разбираетесь в первопричине.
  • Мониторьте не только процент занятого диска, но и возраст самого старого сегмента на партицию. Процент диска на графике растёт линейно и выглядит «обычной утечкой» — а факт, что где-то лежит сегмент трёхмесячной давности при retention в 7 дней, сразу указывает на конкретную причину.
  • Синхронизируйте часы на продюсерах. Проблема с «уехавшим» timestamp почти всегда начинается либо с рассинхронизации системных часов на стороне клиента, либо с путаницы единиц измерения времени в коде — про правильную настройку синхронизации есть отдельный разбор в статье про настройку NTP на сервере.
  • Держите данные Kafka на отдельном разделе или диске, а не на одном томе с ОС и логами системы — тогда даже при повторении подобного инцидента у вас не откажет вся машина целиком, а деградирует только Kafka. Общие принципы разбирали в статье про частые ошибки Kafka на сервере.
  • Настройте мониторинг диска с запасом по времени реакции, а не только на пороге «почти всё, критично» — про то, как это сделать на VPS, есть отдельный материал про мониторинг диска на сервере.

Если вы разворачиваете Kafka с нуля и думаете о ресурсах заранее, стоит также прикинуть объём диска и памяти под retention-политику, которую реально планируете использовать, а не подгонять её постфактум под то, что уместилось.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Нужны сами нейросети для контента?

Генерируйте изображения, видео и озвучку нейросетями на falapi.io — десятки моделей в одном окне. Оплата картой РФ и по СБП.

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

Почему Kafka вообще доверяет времени, которое прислал клиент, а не смотрит на дату файла на диске?

Потому что mtime файла зависит от того, когда сегмент физически создан на конкретном брокере, а не от времени самого события — это ломается при пересборке реплики, миграции данных или восстановлении из бэкапа, когда файлы создаются заново, но данные в них старые. Разработчики Kafka сознательно выбрали ориентир на timestamp записи (CreateTime или LogAppendTime), а не на файловую систему — просто с CreateTime ответственность за корректность значения ложится на продюсера.

Разве log.retention.bytes не решает эту проблему автоматически?

Частично. Он задаёт лимит по объёму на партицию и рано или поздно начнёт удалять сегменты по размеру, даже если time-based retention «залип». Но это не защита от первопричины: во-первых, лимит нужно выставить достаточно тесно, чтобы он сработал раньше, чем кончится диск, во-вторых, вы всё равно потеряете часть свежих данных, если сработает size-based ограничение раньше, чем ожидалось. Это полезный второй рубеж, но не замена диагностике.

Как быстро проверить, нет ли у нас уже такой проблемы?

Сравните время последнего изменения самого старого .log-файла в директории партиции с настроенным retention: ls -la --time-style=full-iso /var/lib/kafka/data/<topic>-<partition>/*.log | sort | head -1. Если самый старый сегмент заметно старше retention.ms, а брокер его не удаляет — это тот самый симптом, и дальше стоит смотреть содержимое сегмента через kafka-dump-log.sh.

Можно ли было поймать это раньше, до выходных?

Да — если бы в мониторинге был отдельный алерт на «возраст старейшего сегмента на партицию» (метрика собирается кастомным экспортером поверх вывода kafka-log-dirs.sh или через JMX), проблема была бы видна в пятницу вечером сразу после деплоя форвардера, а не в субботу ночью по проценту занятого диска.

Стоит ли вообще разрешать продюсерам задавать собственный timestamp?

Разрешайте только там, где это реально нужно бизнес-логике (например, при восстановлении исторических данных с сохранением порядка событий), и обязательно с валидацией диапазона на стороне продюсера или через log.message.timestamp.difference.max.ms на брокере. Для всех остальных топиков LogAppendTime избавляет от целого класса подобных инцидентов раз и навсегда.

Обсудить статью, задать вопрос или начать новую тему

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

Перейти в сообщество →