Kafka не удаляла старые сегменты и забила диск за выходные
В субботу в два часа ночи прилетел алерт: диск на одном из брокеров 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.policy—delete, никаких следов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 ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.
Перейти в сообщество →