Kafka 4.1 tiered storage: переносим старые топики на MinIO и считаем экономию
Apache Kafka 4.1 выпустил tiered storage со стабильным Remote Log Manager. Переносим горячие данные на NVMe (7 дней), архив - на S3-совместимый MinIO. Считаем реальную экономию.
Apache Kafka 4.1 вышел с улучшенным tiered storage и graduated Remote Log Manager
Remote Log Manager в Kafka дебютировал ещё в 3.6 как early access, в 3.8 перешёл в GA, но API в ряде мест сохранял пометку @Unstable. В 4.1 эти оговорки убраны - API стабилизирован полностью, документация наконец описывает поведение при edge cases, а не только happy path. Самое время посмотреть, что это даёт на практике.
Что изменилось в 4.1 по сравнению с 4.0
Ключевые изменения, которые нас интересовали:
- Remote Log Manager API стабилизирован. В 4.0 API помечался как
@Unstableв нескольких местах; в 4.1 это убрано, что важно для плагинных реализаций (мы используем Kafka S3-compatible plugin через MinIO). - Починили race condition при leader election. В 4.0 при смене лидера партиции возникала ситуация, когда новый лидер мог начать загрузку в remote storage до завершения sync локального состояния. В 4.1 введена явная барьерная точка перед началом remote upload на новом лидере.
- Throttle для remote fetch. Теперь можно ограничивать пропускную способность на fetch из remote storage отдельно от локального I/O - это важно, когда S3-совместимое хранилище находится в той же сети, что и брокеры, и забирает пропускную способность у реплицирования.
- Улучшена метрика
remote-log-manager-task-queue-size. Раньше она не отражала задержку upload - теперь это осмысленный индикатор для алертинга.
Наша конфигурация: горячий tier + MinIO
Задача конкретная: у нас несколько топиков с retention 30 дней, которые занимают значительную часть NVMe на брокерах. Реально горячими являются примерно первые 7 дней - именно этот диапазон читают потребители в реальном времени и почти все replay-сценарии. Данные старше недели читают редко: иногда расследование инцидента, иногда перезапуск пайплайна.
MinIO у нас уже стоит в кластере - используем для других нужд. S3-compatible плагин для Kafka брали с open source репозитория, собирали сами ещё на 4.0; под 4.1 пришлось пересобрать с обновлённым API.
Конфигурация тиеринга в server.properties:
remote.log.storage.system.enable=true
log.remote.storage.manager.class.name=io.confluent.kafka.tieredstorage.RemoteStorageManager
# путь к плагину
log.remote.storage.manager.class.path=/opt/kafka/plugins/kafka-tiered-storage-s3.jar
remote.log.metadata.manager.class.name=org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManager
# параметры плагина
remote.log.storage.manager.impl.prefix=s3.
s3.endpoint.url=http://minio.internal:9000
s3.bucket.name=kafka-tiered
s3.credentials.provider=EnvVarCredentialsProvider
# throttle на remote fetch - 200 MiB/s на брокер
remote.fetch.max.wait.ms=500
На уровне топика выставляем:
kafka-configs.sh --alter --topic events-raw \
--add-config 'remote.storage.enable=true,local.retention.ms=604800000'
local.retention.ms=604800000 - это 7 дней в миллисекундах. Total retention.ms остаётся 30 дней, но данные старше 7 дней вытесняются в remote storage.
Что увидели после включения
Мигрировали постепенно: сначала один топик с самым предсказуемым трафиком, потом остальные. Процесс занял несколько дней - Kafka сама выгрузила исторические сегменты в MinIO в фоне, throttle не давал нагрузить сеть.
Объём на брокерах. Топики, переведённые на tiered storage, сократили локальный disk usage примерно в 3-4 раза (30 дней -> 7 дней локально). Суммарно по кластеру экономия вышла около 60% от общего объёма топиков с длинным retention. Это чуть больше, чем мы считали на бумаге, - несколько топиков имели неравномерный профиль записи, и исторический хвост у них оказался плотнее ожидаемого.
MinIO под нагрузкой. Первые сутки после включения upload-задачи создавали заметный трафик - выгрузка накопленных исторических сегментов. Потом устаканилось: текущая нагрузка - это просто вытеснение сегментов старше 7 дней по мере их создания, что несопоставимо меньше.
Replay из remote storage. Это тот самый вопрос, который был открытым ещё в 4.0. Потребители, которые читают исторические данные (старше 7 дней), получают данные из MinIO через брокера - Kafka прозрачно их подтягивает. Задержка на первый fetch заметно выше, чем при чтении с локального диска, - зависит от latency до MinIO и размера сегмента. На нашем кластере это добавляет 200-500 мс к первому fetch при replay. Для бэчевых пайплайнов это несущественно; для latency-sensitive consumers данные старше недели в принципе неактуальны, так что проблемы нет.
Race condition из 4.0. В 4.0 у нас был один задокументированный инцидент: после планового rolling restart несколько сегментов попали в remote storage дважды с разными offset-метаданными. Разобрались вручную. В 4.1 за несколько недель тестирования ничего подобного не воспроизвелось.
Что оставили на потом
Топики с retention меньше суток на tiered storage не переводили - смысла нет, там весь объём помещается в горячем tier. Топики compacted с ключами тоже пока оставили: документация по взаимодействию compaction и remote storage в 4.1 стала лучше, но у нас специфичный cleanup policy, хочется понять нюансы перед включением.
В рамках DWH/BI-проектов мы обычно смотрим на Kafka как на транспорт до хранилища, но объёмы в топиках нередко становятся самостоятельной проблемой. Tiered storage - понятный инструмент с теперь уже стабильным API, и с 4.1 можно убрать последние сомнения насчёт production-готовности.