📨 Розділ 15 · Питання #28

Як видаляються старі повідомлення з топика

В Apache Kafka повідомлення ніколи не видаляються поштучно.


🟢 Junior Level

Механіка фізичного видалення повідомлень

В Apache Kafka повідомлення ніколи не видаляються поштучно.

Оскільки лог партиції спроєктований як високопродуктивний послідовний файл із прямим доступом (Append-only Log), видалення окремого запису з середини файлу вимагало б повного перезапису гігабайтів даних на диску та перерахунку зміщень. Це катастрофічно зруйнувало б продуктивність і throughput брокера.

Замість цього видалення відбувається цілими блоками файлів — закритими сегментами логу (LogSegment):

Файлова система ОС (/var/lib/kafka/data/events-0/):
[Сегмент 0: 00000.log]   [Сегмент 1: 00500.log]   [Сегмент 2: 01000.log (АКТИВНИЙ)]
      (1 ГБ)                   (1 ГБ)                     (400 МБ)
   Offsets: 0 - 499         Offsets: 500 - 999         Offsets: 1000 - ...
         │
         ▼ (Сплив retention.ms)
  [ВИДАЛЯЄТЬСЯ ЦІЛКОМ]

Структура сегмента на диску

Кожен сегмент партиції складається з набору узгоджених файлів із однаковою базовою назвою (стартовим зміщенням сегмента):

  • 00000000000000000000.log — фізичні бінарні дані пакетів (батчів) повідомлень.
  • 00000000000000000000.index — розріджений індекс зміщень (offset -> physical position in bytes).
  • 00000000000000000000.timeindex — часовий індекс (timestamp -> offset).
  • 00000000000000000000.txnindex — індекс транзакцій (якщо топік транзакційний).

Коли брокер приймає рішення про видалення сегмента, усі ці файли видаляються з операційної системи атомарно як єдине ціле.


🟡 Middle Level

Життєвий цикл видалення закритого сегмента

Видалення застарілих сегментів керується внутрішнім планувальником брокера та виконується за такими кроками:

1. Фоновий потік `kafka-log-retention` прокидається кожні `log.retention.check.interval.ms` (5 хв).
2. Виконується перевірка кожного закритого сегмента:
   - За часом: maxTimestamp з .timeindex < (now - retention.ms)
   - За розміром: сумарний обсяг партиції > retention.bytes (видаляються найстаріші сегменти).
3. Перейменування файлів:
   00000.log       -> 00000.log.deleted
   00000.index     -> 00000.index.deleted
   00000.timeindex -> 00000.timeindex.deleted
4. Запуск таймера відстрочки `log.segment.delete.delay.ms` (за замовчуванням 60 секунд).
5. Після закінчення 60 секунд фоновий потік викликає системне видалення файлів із файлової системи.
6. Зміщення початку логу (`LogStartOffset`) зміщується на базовий офсет наступного живого сегмента.

Навіщо потрібен суфікс .deleted і затримка 60 секунд

Затримка видалення (log.segment.delete.delay.ms = 60000) — найважливіший захисний механізм розподіленого читання:

  1. Захист активних операцій читання консьюмерів (In-flight Reads): У момент, коли потік очищення вирішив видалити сегмент, певний консьюмер або follower-репліка вже могли відкрити файловий дескриптор і вичитувати дані через системний виклик sendfile(2).
  2. Перейменування файлу на *.deleted негайно виключає сегмент із метаданих і маршрутів пошуку для нових запитів FetchRequest, але залишає файл на диску.
  3. Клієнти, які вже почали читання, спокійно завершують передачу порції байтів протягом 60-секундного вікна без аварійних збоїв введення-виведення (IOException).

Програмне видалення даних: AdminClient.deleteRecords()

Починаючи з версії Kafka 0.11, розробники можуть програмно ініціювати видалення старих повідомлень із Java-коду, не чекаючи спрацьовування автоматичних політик retention:

Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

try (AdminClient adminClient = AdminClient.create(props)) {
    TopicPartition partition = new TopicPartition("orders", 0);
    
    // Задаємо зміщення, ДО якого всі записи мають бути видалені (до offset 150000)
    Map<TopicPartition, RecordsToDelete> recordsToDelete = Map.of(
        partition, RecordsToDelete.beforeOffset(150000L)
    );

    // Виконуємо асинхронний запит до брокера
    DeleteRecordsResult result = adminClient.deleteRecords(recordsToDelete);
    result.all().get();
    
    System.out.println("Зміщення початку логу orders-0 успішно зсунуто на 150000");
}

Важливо розуміти: Виклик deleteRecords(beforeOffset) не розрізає файли сегментів навпіл. Він зсуває покажчик LogStartOffset партиції на значення 150000. Будь-які закриті сегменти, чий кінцевий офсет виявився строго меншим за 150000, негайно позначаються як .deleted і знищуються.


🔴 Senior Level

Низькорівнева взаємодія з Linux VFS та Inode

На рівні ядра операційної системи Linux видалення сегмента логу спирається на стандартну семантику файлових дескрипторів:

[Потік LogRetention] ──unlink("00000.log.deleted")──> Видалення імені зі структури директорії (dentry)
                                                                 │
                                                                 ▼
                                                  Лічильник жорстких посилань Inode = 0
                                                                 │
       ┌─────────────────────────────────────────────────────────┴───────────────────────────────────────┐
       ▼                                                                                                 ▼
Якщо немає відкритих дескрипторів (FD count = 0):                               Якщо консьюмер читає файл (FD count > 0):
- ОС негайно звільняє дискові блоки.                                             - ОС НЕ звільняє дискові блоки.
- Дисковий простір стає доступним у `df -h`.                                     - Консьюмер продовжує читати файл через Zero-Copy sendfile(2).
                                                                                 - Лише після close(fd) ядро ОС фізично звільняє блоки!

Пастка переповнення диска при завислих дескрипторах: Якщо в застосунку або процесі брокера відбувається витік дескрипторів (дескриптор відкрито стороннім процесом бекапу або підвислим сокетом), команда rm або перейменування в .deleted видалить запис про файл, але дисковий простір не звільниться. Утиліта ls покаже порожню директорію, а df -h рапортуватиме про 100% зайнятого місця. Проблема діагностується командою lsof | grep deleted.

Зміщення LogStartOffset та виняток OffsetOutOfRangeException

При фізичному видаленні сегментів мінімальний доступний офсет партиції монотонно зростає:

Початковий стан:
LogStartOffset = 0 ─────────────────────────────> LogEndOffset = 250000

Після видалення трьох сегментів:
                     LogStartOffset = 150000 ───> LogEndOffset = 250000
                     (Офсети 0..149999 фізично видалені з диска)

Якщо повільний консьюмер спробує прочитати офсет 120 000, брокер виявить, що: \(\text{Requested Offset (120 000)} < \text{LogStartOffset (150 000)}\)

  1. Брокер повертає клієнту помилку OFFSET_OUT_OF_RANGE.
  2. Поведінка консьюмера визначається конфігурацією auto.offset.reset:
    • earliest: консьюмер автоматично скине свій покажчик на новий LogStartOffset (150 000) і продовжить читання. Повідомлення 120 000–149 999 втрачені безповоротно.
    • latest: консьюмер перестрибне в самий кінець партиції (LogEndOffset).
    • none: брокер викине фатальний виняток OffsetOutOfRangeException безпосередньо в прикладний код Java.

Роловер сегментів (Segment Rolling) та захист від I/O Burst

Щоб застарілі дані могли видалятися регулярно, нові сегменти мають створюватися своєчасно (Segment Rollover). Сегмент закривається і «запечатується» при настанні будь-якої з подій:

  1. Розмір активного сегмента досяг segment.bytes (за замовчуванням 1 ГБ).
  2. Час із моменту першого запису в активний сегмент перевищив segment.ms (або log.roll.ms, за замовчуванням 7 діб).
  3. Індекс зміщень .index заповнився до межі segment.index.bytes (за замовчуванням 10 МБ).

Проблема Thundering Herd при роловері: Якщо в топіку 64 партиції, і всі вони були створені одночасно в понеділок о 10:00, то рівно через 7 діб у понеділок о 10:00 усі 64 партиції одночасно закриють активні сегменти, синхронно створять 64 нових файли та розпочнуть скидання fsync на диск. Це викличе різкий сплеск затримок дискової підсистеми (I/O Spike). Рішення: Параметр log.roll.jitter.ms (рандомізований розкид затримки) вносить випадкове зміщення в момент роловера, плавно розподіляючи навантаження в часі.


4 Tricky Questions

1. Що станеться, якщо в партиції залишився лише один сегмент, і він є активним (activeSegment), але всі повідомлення в ньому були записані 30 днів тому при retention.ms = 7 днів?

Відповідь: Дані не будуть видалені.

  1. Активний сегмент (activeSegment) — це єдиний сегмент партиції, відкритий для дозапису нових даних.
  2. В архітектурі Kafka діє фундаментальне залізне правило: активний сегмент захищений від видалення за будь-яких обставин, навіть якщо абсолютно всі повідомлення в ньому прострочені за часом, а розмір перевищує всі ліміти.
  3. Сегмент зможе видалитися лише після того, як у топік надійде нове повідомлення, яке спровокує Rollover (закриття поточного сегмента і створення нового активного), або спрацює примусовий таймер log.roll.ms.

2. Як брокер вирішує, чи застарів сегмент: за часом створення файлу в Linux, часом першого запису чи часом останнього запису?

Відповідь: Брокер спирається на максимальну часову мітку серед усіх записів сегмента (maxTimestamp), яка зберігається в заголовку індексного файлу .timeindex:

  1. Брокер не використовує системний час створення файлу в ОС (ctime / mtime), оскільки при копіюванні чи відновленні файлів ці мітки можуть збиватися.
  2. При кожному записі повідомлень у сегмент брокер оновлює максимальний timestamp у .timeindex.
  3. Сегмент визнається застарілим лише тоді, коли навіть найсвіжіший запис у цьому сегменті старший, ніж now - retention.ms. Якщо в закритий сегмент потрапило хоча б одне свіже повідомлення, весь гігабайтний сегмент утримуватиметься на диску.

3. Як метод AdminClient.deleteRecords() взаємодіє з механікою сегментів, якщо запитано видалення до зміщення, яке знаходиться всередині активного сегмента?

Відповідь:

  1. Виклик deleteRecords(beforeOffset = 500) негайно оновлює метадані партиції на брокері, виставляючи LogStartOffset = 500.
  2. Якщо офсет 500 знаходиться всередині активного сегмента, фізичного видалення файлу не відбувається.
  3. Проте для будь-яких консьюмерів повідомлення зі зміщеннями від 0 до 499 стають логічно невидимими. Спроба вичитати офсет 300 призведе до OffsetOutOfRangeException.
  4. Фізичне видалення відбудеться лише тоді, коли активний сегмент закриється, і наступний за ним файл зробить старий сегмент закритим. Після цього брокер побачить, що максимальний офсет старого сегмента $< 500$, і видалить файл цілком.

4. Чи може видалення сегмента на лідер-брокері призвести до аварії реплікації на follower-репліках (Under-Replicated Partitions)?

Відповідь: Так, якщо репліка відстає занадто сильно.

  1. Кожен брокер виконує видалення сегментів за локальним таймером kafka-log-retention.
  2. Якщо follower-брокер відключився через мережеву ізоляцію на 8 днів при retention.ms = 7 днів, лідер встигне видалити сегменти з даними за період простою.
  3. Коли follower відновить зв’язок і запитає дані (FetchRequest) зі свого останнього відомого офсету, лідер поверне помилку OFFSET_OUT_OF_RANGE, оскільки цих сегментів на лідері фізично вже немає.
  4. Follower не зможе наздогнати лідера стандартним шляхом, випадає з ISR і переходить в аварійний стан. Адміністраторам доведеться або перевизначати репліку з нуля, або використовувати Tiered Storage / дзеркалювання.

🎯 Шпаргалка для інтерв’ю

30-секундна відповідь (Elevator Pitch)

«Kafka видаляє старі повідомлення не поштучно, а цілими закритими сегментами логу (LogSegment), що забезпечує максимальну продуктивність введення-виведення без фрагментації диска. Фоновий потік kafka-log-retention кожні 5 хвилин інспектує закриті сегменти. Якщо максимальна часова мітка сегмента (maxTimestamp у .timeindex) перевищила retention.ms або сумарний обсяг партиції перевищив retention.bytes, файли сегмента перейменовуються в *.deleted. Через захисний інтервал log.segment.delete.delay.ms (60 секунд), необхідний для завершення активних операцій читання консьюмерами через sendfile(2), файли фізично видаляються з ОС, а LogStartOffset партиції зсувається вперед. Активний сегмент захищений від видалення завжди».

Ключові метрики для Production-моніторингу

  • kafka.log:type=Log,name=LogStartOffset,topic={topic},partition={id}: мінімальний фізично доступний офсет у лозі партиції.
  • kafka.log:type=Log,name=NumLogSegments: загальна кількість файлів сегментів на диску для даної партиції.
  • kafka.log:type=LogManager,name=OfflineLogDirectoryCount: кількість пошкоджених точок монтування дисків брокера.

Червоні прапорці (Чого говорити НЕ МОЖНА)

  • ❌ «Kafka знаходить старе повідомлення в середині логу і стирає його байти» — Лог Kafka незмінний (Append-only); видалення відбувається лише цілими файлами сегментів.
  • ❌ «При видаленні файлу в Linux консьюмер, який читав цей файл, негайно падає з помилкою FileNotFound» — Ядро Linux утримує дискові inode відкритими до закриття всіх файлових дескрипторів процесу.
  • ❌ «deleteRecords() може розрізати сегмент логу на дві частини» — Метод лише зсуває LogStartOffset, а фізичне видалення сегментів відбувається стандартним блочним способом.
  • ❌ «Активний сегмент видаляється, якщо в ньому немає нових даних» — Активний сегмент ніколи не видаляється, доки не відбудеться Rollover.

Пов’язані теми