📨 Раздел 15 · Вопрос #27

Что такое 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 определяет фундаментальный механизм освобождения диска:

  1. delete (по умолчанию): Сообщения удаляются по истечении заданного времени (retention.ms) либо при превышении лимита размера на диске (retention.bytes). Используется для потоков событий, логов, метрик и аналитики.
  2. compact (компактизация лога): Kafka сохраняет только последнее актуальное состояние (latest value) для каждого уникального ключа (key). Старые версии значений с тем же ключом вычищаются. Используется для таблиц состояний (Changelog), профилей пользователей, кэшей и настроек.
  3. 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 МБ, пишется сейчас)     -> НЕ ТРОГАЕТСЯ

Правила удаления сегментов:

  1. Единица удаления — сегмент целиком: Фоновый поток kafka-log-retention просыпается каждые log.retention.check.interval.ms (по умолчанию 5 минут) и анализирует закрытые сегменты.
  2. Сегмент удаляется только тогда, когда ВСЕ сообщения в нем устарели: Брокер берет максимальную временную метку (maxTimestamp) из .timeindex закрытого сегмента. Если $(T_{current} - T_{segment_max}) > \text{retention.ms}$, сегмент помечается для удаления (.deleted) и физически удаляется через log.segment.delete.delay.ms (по умолчанию 1 минута).
  3. Активный сегмент (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
                                                                     (Никогда не компактится)
  1. SkimpyOffsetMap: Поток LogCleaner сканирует «грязную» часть лога (dirty section) и строит в оперативной памяти (off-heap) компактную хэш-таблицу SkimpyOffsetMap (размер задается через log.cleaner.dedupe.buffer.size, по умолчанию 128 МБ). Карта связывает 128-битный хэш ключа с его последним известным смещением (offset).
  2. Перезапись сегментов (Compaction Pass): Очиститель последовательно читает чистую и грязную секции. Если для ключа записи в карте найден более свежий офсет, старая запись отбрасывается. Живые записи копируются во временный сегмент .clean, который затем атомарно подменяет исходные файлы.
  3. Удаление ключей через 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 (дисковая утечка):

  1. При CreateTime брокер проставляет в индекс .timeindex временную метку, указанную клиентом в метаданных записи.
  2. В закрытом сегменте maxTimestamp зафиксируется как «2099 год».
  3. Проверяя условие удаления, брокер вычисляет: \(\text{now}() - \text{segment.maxTimestamp} < \text{retention.ms}\) Поскольку разница отрицательна, условие никогда не выполнится в обозримом будущем.
  4. Данный сегмент лога и все последующие сегменты партиции никогда не удалятся по тайм-ауту времени! Они будут храниться на диске вечно, пока не сработает ограничение по retention.bytes (если оно настроено).
  5. Защита: На брокере необходимо настраивать валидацию меток времени: log.message.timestamp.difference.max.ms (отклонение не более чем на $N$ часов от часов сервера).

2. Консьюмер был выключен на техническое обслуживание в течение 48 часов. На топике настроена cleanup.policy = compact и delete.retention.ms = 86400000 (24 часа). К чему это приведет при включении консьюмера?

Ответ: Произойдет рассинхронизация локального состояния консьюмера (Silent State Inconsistency):

  1. Во время простоя консьюмера в топик поступил tombstone (key="user_123", value=null) для удаления пользователя.
  2. Поскольку консьюмер был выключен 48 часов, а delete.retention.ms составляет 24 часа, Log Cleaner успел выполнить компактизацию и физически вычистить маркер tombstone из сегментов лога.
  3. Когда консьюмер проснется, в логе уже не будет никаких следов записи user_123.
  4. Консьюмер прочитает оставшиеся офсеты, но в его локальной базе данных пользователь user_123 так и останется активным навсегда!
  5. Правило: Значение 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). В режиме compact Kafka сохраняет последнее известное значение для каждого ключа, удаляя устаревшие версии через фоновый процесс 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 (по умолчанию сутки).

Связанные темы