Что такое retention policy
В отличие от классических очередей сообщений (RabbitMQ, IBM MQ), где сообщение безвозвратно удаляется сразу после подтверждения консьюмером (ack), Kafka хранит сообщения всегда,...
🟢 Junior Level
Что такое Retention Policy простыми словами
Retention Policy (политика удержания данных) в Apache Kafka — это набор правил, определяющих, как долго и в каком объеме брокеры сохраняют сообщения в партициях топика, а также условия, при которых старые данные подлежат очистке.
В отличие от классических очередей сообщений (RabbitMQ, IBM MQ), где сообщение безвозвратно удаляется сразу после подтверждения консьюмером (ack), Kafka хранит сообщения всегда, независимо от того, сколько консьюмеров их прочитало и прочитало ли вообще.
Лог партиции на диске:
[Сегмент 1: 8 дней назад] [Сегмент 2: 4 дня назад] [Сегмент 3: 1 день назад] [Сегмент 4: Активный]
│
▼ (retention.ms = 7 дней)
[КАНДИДАТ НА УДАЛЕНИЕ] ──────> Удаляется целиком фоновым процессом брокера
Основные стратегии очистки (Cleanup Policies)
Параметр топика cleanup.policy определяет фундаментальный механизм освобождения диска:
delete(по умолчанию): Сообщения удаляются по истечении заданного времени (retention.ms) либо при превышении лимита размера на диске (retention.bytes). Используется для потоков событий, логов, метрик и аналитики.compact(компактизация лога): Kafka сохраняет только последнее актуальное состояние (latest value) для каждого уникального ключа (key). Старые версии значений с тем же ключом вычищаются. Используется для таблиц состояний (Changelog), профилей пользователей, кэшей и настроек.compact,delete: Гибридный режим. Для записей выполняется компактизация по ключам, но если запись не обновлялась дольшеretention.ms, она удаляется окончательно.
Базовые параметры конфигурации
# 1. Время жизни данных (в миллисекундах)
retention.ms=604800000 # 7 суток (7 * 24 * 60 * 60 * 1000)
# 2. Максимальный размер лога на одну партицию (в байтах)
retention.bytes=10737418240 # 10 ГБ на КАЖДУЮ партицию (не на весь топик!)
# 3. Политика очистки
cleanup.policy=delete
Управление через Kafka CLI
# Изменение времени удержания существующего топика на 3 дня (259200000 мс)
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name orders \
--alter --add-config retention.ms=259200000
# Проверка установленных параметров retention
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name orders --describe
🟡 Middle Level
Взаимодействие Retention и сегментов лога (Log Segments)
Kafka никогда не удаляет отдельные сообщения из середины файла. Лог партиции на диске разбит на неизменяемые файлы фиксированного размера — сегменты (.log), сопровождаемые индексами смещений (.index) и временных меток (.timeindex).
Директория партиции /var/lib/kafka/data/orders-0/:
├── 00000000000000000000.log (Закрытый сегмент, 1 ГБ, создан 10 дней назад) -> УДАЛЯЕТСЯ
├── 00000000000000000000.index
├── 00000000000000000000.timeindex
├── 00000000000001948520.log (Закрытый сегмент, 1 ГБ, создан 3 дня назад) -> ХРАНИТСЯ
├── 00000000000001948520.index
├── 00000000000003891440.log (АКТИВНЫЙ сегмент, 400 МБ, пишется сейчас) -> НЕ ТРОГАЕТСЯ
Правила удаления сегментов:
- Единица удаления — сегмент целиком: Фоновый поток
kafka-log-retentionпросыпается каждыеlog.retention.check.interval.ms(по умолчанию 5 минут) и анализирует закрытые сегменты. - Сегмент удаляется только тогда, когда ВСЕ сообщения в нем устарели: Брокер берет максимальную временную метку (
maxTimestamp) из.timeindexзакрытого сегмента. Если $(T_{current} - T_{segment_max}) > \text{retention.ms}$, сегмент помечается для удаления (.deleted) и физически удаляется черезlog.segment.delete.delay.ms(по умолчанию 1 минута). - Активный сегмент (Active Segment) НЕ подлежит удалению НИКОГДА: Сегмент, в который в текущий момент поступают записи от продюсеров, не может быть удален по retention, даже если он превысил
retention.msилиretention.bytes.
Опасная ошибка: retention.bytes рассчитывается НА ПАРТИЦИЮ
Критическая ошибка начинающих инженеров — предположение, что retention.bytes ограничивает суммарный размер топика.
\(\text{Максимальный размер топика на диске} = \text{retention.bytes} \times \text{Количество партиций} \times \text{Replication Factor}\)
Пример аварии: Администратор выставил
retention.bytes = 50 GB, полагая, что топик займет 50 ГБ. Однако у топика 32 партиции и фактор репликации 3. Суммарный объем на кластере составит: $50 \times 32 \times 3 = 4800\text{ ГБ (4.8 ТБ)}$! В результате дисковые массивы брокеров переполняются, приводя к отказу всего кластера.
🔴 Senior Level
Внутреннее устройство Log Cleaner: Clean vs Dirty Log
Для топиков с cleanup.policy=compact удаление устаревших версий ключей выполняется пулом фоновых потоков LogCleaner (kafka.log.LogCleaner).
СТРУКТУРА ПАРТИЦИИ ПРИ КОМПАКТИЗАЦИИ
┌───────────────────────────────────────────────┬───────────────────────────────┐
│ CLEAN SECTION │ DIRTY SECTION │
│ (Уже компактизированные закрытые сегменты) │ (Новые закрытые сегменты) │
│ Каждый ключ встречается строго 1 раз │ Могут быть дубликаты ключей │
└───────────────────────────────────────────────┴───────────────────────────────┘
▲ ▲
│ │
firstDirtyOffset activeSegment
(Никогда не компактится)
- SkimpyOffsetMap:
Поток
LogCleanerсканирует «грязную» часть лога (dirty section) и строит в оперативной памяти (off-heap) компактную хэш-таблицуSkimpyOffsetMap(размер задается черезlog.cleaner.dedupe.buffer.size, по умолчанию 128 МБ). Карта связывает 128-битный хэш ключа с его последним известным смещением (offset). - Перезапись сегментов (Compaction Pass):
Очиститель последовательно читает чистую и грязную секции. Если для ключа записи в карте найден более свежий офсет, старая запись отбрасывается. Живые записи копируются во временный сегмент
.clean, который затем атомарно подменяет исходные файлы. - Удаление ключей через Tombstone (Маркер удаления):
Чтобы удалить ключ из compact-топика, продюсер отправляет запись с данным ключом и
value == null(tombstone record).- Tombstone сохраняется в логе, чтобы консьюмеры успели прочитать факт удаления ключа.
- По истечении параметра
delete.retention.ms(по умолчанию 24 часа) процесс Log Cleaner окончательно стирает сам маркер tombstone из сегментов лога.
KIP-405: Tiered Storage (Многоуровневое хранение)
Начиная с Kafka 3.x в архитектуру внедрен механизм Tiered Storage (многоуровневое хранение данных), решивший проблему масштабирования retention:
[Kafka Broker]
│
├── Быстрый локальный диск (NVMe/SSD):
│ local.retention.ms = 86400000 (1 сутки) / local.retention.bytes
│ (Хранит горячие данные для real-time потребителей с нулевой задержкой)
│
└── Фоновая выгрузка через RemoteStorageManager (RSM):
retention.ms = 31536000000 (1 год)
[Объектное хранилище S3 / GCS / Azure Blob / MinIO]
(Хранит холодные закрытые сегменты логов по цене $0.02 за ГБ)
- Локальные диски брокеров очищаются по агрессивному
local.retention.ms(например, 24 часа), освобождая дорогой NVMe-кэш. - Исторические консьюмеры (ETL, аудит, ML-модели) прозрачно вычитывают старые сегменты напрямую из S3 через стандартный Kafka Fetch Protocol без ручного экспорта в HDFS/DataLake.
4 Tricky Questions
1. Что произойдет, если продюсер случайно отправит батч сообщений с ошибочной временной меткой из далекого будущего (например, 2099 год), при message.timestamp.type = CreateTime?
Ответ:
Произойдет полная блокировка удаления сегмента по retention.ms (дисковая утечка):
- При
CreateTimeброкер проставляет в индекс.timeindexвременную метку, указанную клиентом в метаданных записи. - В закрытом сегменте
maxTimestampзафиксируется как «2099 год». - Проверяя условие удаления, брокер вычисляет: \(\text{now}() - \text{segment.maxTimestamp} < \text{retention.ms}\) Поскольку разница отрицательна, условие никогда не выполнится в обозримом будущем.
- Данный сегмент лога и все последующие сегменты партиции никогда не удалятся по тайм-ауту времени! Они будут храниться на диске вечно, пока не сработает ограничение по
retention.bytes(если оно настроено). - Защита: На брокере необходимо настраивать валидацию меток времени:
log.message.timestamp.difference.max.ms(отклонение не более чем на $N$ часов от часов сервера).
2. Консьюмер был выключен на техническое обслуживание в течение 48 часов. На топике настроена cleanup.policy = compact и delete.retention.ms = 86400000 (24 часа). К чему это приведет при включении консьюмера?
Ответ: Произойдет рассинхронизация локального состояния консьюмера (Silent State Inconsistency):
- Во время простоя консьюмера в топик поступил tombstone (
key="user_123", value=null) для удаления пользователя. - Поскольку консьюмер был выключен 48 часов, а
delete.retention.msсоставляет 24 часа, Log Cleaner успел выполнить компактизацию и физически вычистить маркер tombstone из сегментов лога. - Когда консьюмер проснется, в логе уже не будет никаких следов записи
user_123. - Консьюмер прочитает оставшиеся офсеты, но в его локальной базе данных пользователь
user_123так и останется активным навсегда! - Правило: Значение
delete.retention.msв compact-топиках обязано гарантированно превышать максимально допустимое время простоя (Max SLA Downtime) любых потребителей данного топика.
3. Может ли Kafka удалить активный сегмент партиции, если объем партиции превысил retention.bytes в 10 раз?
Ответ:
Нет, активный сегмент (activeSegment) не может быть удален ни при каких обстоятельствах.
Даже если на топике выставлен retention.bytes = 100 MB, а в активный сегмент уже записано 1 ГБ данных, брокер не тронет его, пока сегмент не закроется (не произойдет Roll по достижении segment.bytes или segment.ms). Только после того, как текущий сегмент «запечатывается» и открывается новый активный файл, старый закрытый сегмент проверяется политикой retention и удаляется.
4. В чем разница между log.roll.ms и retention.ms, и как их неправильное соотношение может вызвать переполнение диска?
Ответ:
retention.ms— время, по истечении которого закрытый сегмент может быть удален.log.roll.ms— максимальное время, в течение которого активный сегмент остается открытым для дозаписи, даже если он не заполнился доsegment.bytes.- Ловушка конфигурации на низконагруженных топиках:
Если в топик пишутся редкие события, и
segment.bytes = 1 GB(по умолчанию), аlog.roll.msне настроен (по умолчанию 7 суток), активный сегмент может оставаться открытым месяцами. Даже еслиretention.ms = 86400000(1 день), брокер не сможет удалить ни одного сообщения, пока сегмент не закроется через 7 дней. В результате фактическое время удержания данных составит $\text{log.roll.ms} + \text{retention.ms} = 8$ суток вместо ожидаемых 24 часов.
🎯 Шпаргалка для интервью
30-секундный ответ (Elevator Pitch)
«Retention Policy в Kafka определяет жизненный цикл данных на брокере через параметр
cleanup.policy. По умолчанию используется режимdelete, где устаревшие данные удаляются закрытыми сегментами по достижении лимита времени (retention.ms) или объема диска (retention.bytes). В режимеcompactKafka сохраняет последнее известное значение для каждого ключа, удаляя устаревшие версии через фоновый процессLogCleaner, а удаление ключей оформляется через tombstone-сообщения (value=null). Главный нюанс: retention работает только с закрытыми сегментами и никогда не удаляет активный сегмент. Начиная с Kafka 3.x, доступен Tiered Storage (KIP-405), позволяющий разделять горячие данные на локальных NVMe (local.retention.ms) и холодные архивы в S3».
Ключевые метрики для Production-мониторинга
kafka.log:type=Log,name=Size,topic={topic},partition={id}: физический размер партиции на диске в байтах.kafka.log:type=LogCleanerManager,name=max-dirty-percent: процент «грязных» данных в compact-топиках. Если $> 0.5$, очиститель отстает.disk_used_percentна серверах брокеров: важнейший триггер для предупреждения аварийного перехода брокера в Read-Only.
Красные флаги (Чего говорить НЕЛЬЗЯ)
- ❌ «retention.bytes — это максимальный размер всего топика в кластере» — Это лимит на каждую отдельную партицию. Реальный объем равен $\text{retention.bytes} \times \text{партиции} \times \text{репликация}$.
- ❌ «Kafka удаляет каждое сообщение ровно в ту миллисекунду, когда истек retention.ms» — Сообщения удаляются только пачками в составе целых закрытых сегментов раз в
log.retention.check.interval.ms. - ❌ «Активный сегмент удаляется, если диск заполнен на 100%» — Активный сегмент никогда не удаляется очистителем; переполнение диска приводит к падению процесса брокера.
- ❌ «В compact-топиках ключи удаляются мгновенно» — Удаление требует отправки tombstone (
value=null), который хранится в логе ещеdelete.retention.ms(по умолчанию сутки).