Що таке 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]
(Зберігає холодні закриті сегменти логів за низькою ціною зберігання)
- Локальні диски брокерів очищаються за агресивним
local.retention.ms(наприклад, 24 години), звільняючи дорогий NVMe-диск. - Історичні споживачі (ETL, аудит, ML-моделі) прозоро вичитують старі сегменти безпосередньо з S3 через стандартний Kafka Fetch Protocol без ручного експорту в HDFS або Data Lake.
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(за замовчуванням добу).