Як видаляються старі повідомлення з топика
В 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) — найважливіший захисний механізм розподіленого читання:
- Захист активних операцій читання консьюмерів (In-flight Reads):
У момент, коли потік очищення вирішив видалити сегмент, певний консьюмер або follower-репліка вже могли відкрити файловий дескриптор і вичитувати дані через системний виклик
sendfile(2). - Перейменування файлу на
*.deletedнегайно виключає сегмент із метаданих і маршрутів пошуку для нових запитівFetchRequest, але залишає файл на диску. - Клієнти, які вже почали читання, спокійно завершують передачу порції байтів протягом 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)}\)
- Брокер повертає клієнту помилку
OFFSET_OUT_OF_RANGE. - Поведінка консьюмера визначається конфігурацією
auto.offset.reset:earliest: консьюмер автоматично скине свій покажчик на новийLogStartOffset(150 000) і продовжить читання. Повідомлення 120 000–149 999 втрачені безповоротно.latest: консьюмер перестрибне в самий кінець партиції (LogEndOffset).none: брокер викине фатальний винятокOffsetOutOfRangeExceptionбезпосередньо в прикладний код Java.
Роловер сегментів (Segment Rolling) та захист від I/O Burst
Щоб застарілі дані могли видалятися регулярно, нові сегменти мають створюватися своєчасно (Segment Rollover). Сегмент закривається і «запечатується» при настанні будь-якої з подій:
- Розмір активного сегмента досяг
segment.bytes(за замовчуванням 1 ГБ). - Час із моменту першого запису в активний сегмент перевищив
segment.ms(абоlog.roll.ms, за замовчуванням 7 діб). - Індекс зміщень
.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 днів?
Відповідь: Дані не будуть видалені.
- Активний сегмент (
activeSegment) — це єдиний сегмент партиції, відкритий для дозапису нових даних. - В архітектурі Kafka діє фундаментальне залізне правило: активний сегмент захищений від видалення за будь-яких обставин, навіть якщо абсолютно всі повідомлення в ньому прострочені за часом, а розмір перевищує всі ліміти.
- Сегмент зможе видалитися лише після того, як у топік надійде нове повідомлення, яке спровокує Rollover (закриття поточного сегмента і створення нового активного), або спрацює примусовий таймер
log.roll.ms.
2. Як брокер вирішує, чи застарів сегмент: за часом створення файлу в Linux, часом першого запису чи часом останнього запису?
Відповідь:
Брокер спирається на максимальну часову мітку серед усіх записів сегмента (maxTimestamp), яка зберігається в заголовку індексного файлу .timeindex:
- Брокер не використовує системний час створення файлу в ОС (
ctime/mtime), оскільки при копіюванні чи відновленні файлів ці мітки можуть збиватися. - При кожному записі повідомлень у сегмент брокер оновлює максимальний timestamp у
.timeindex. - Сегмент визнається застарілим лише тоді, коли навіть найсвіжіший запис у цьому сегменті старший, ніж
now - retention.ms. Якщо в закритий сегмент потрапило хоча б одне свіже повідомлення, весь гігабайтний сегмент утримуватиметься на диску.
3. Як метод AdminClient.deleteRecords() взаємодіє з механікою сегментів, якщо запитано видалення до зміщення, яке знаходиться всередині активного сегмента?
Відповідь:
- Виклик
deleteRecords(beforeOffset = 500)негайно оновлює метадані партиції на брокері, виставляючиLogStartOffset = 500. - Якщо офсет 500 знаходиться всередині активного сегмента, фізичного видалення файлу не відбувається.
- Проте для будь-яких консьюмерів повідомлення зі зміщеннями від 0 до 499 стають логічно невидимими. Спроба вичитати офсет 300 призведе до
OffsetOutOfRangeException. - Фізичне видалення відбудеться лише тоді, коли активний сегмент закриється, і наступний за ним файл зробить старий сегмент закритим. Після цього брокер побачить, що максимальний офсет старого сегмента $< 500$, і видалить файл цілком.
4. Чи може видалення сегмента на лідер-брокері призвести до аварії реплікації на follower-репліках (Under-Replicated Partitions)?
Відповідь: Так, якщо репліка відстає занадто сильно.
- Кожен брокер виконує видалення сегментів за локальним таймером
kafka-log-retention. - Якщо follower-брокер відключився через мережеву ізоляцію на 8 днів при
retention.ms = 7 днів, лідер встигне видалити сегменти з даними за період простою. - Коли follower відновить зв’язок і запитає дані (
FetchRequest) зі свого останнього відомого офсету, лідер поверне помилкуOFFSET_OUT_OF_RANGE, оскільки цих сегментів на лідері фізично вже немає. - 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.