Как удаляются старые сообщения из топика
В Apache Kafka сообщения никогда не удаляются поштучно.
🟢 Junior Level
Механика физического удаления сообщений
В Apache Kafka сообщения никогда не удаляются поштучно.
Поскольку лог партиции спроектирован как высокопроизводительный последовательный файл с прямым доступом (Append-only Log), удаление отдельной записи из середины файла потребовало бы полной перезаписи гигабайтов данных на диске и пересчета смещений, что полностью разрушило бы производительность брокера.
Вместо этого удаление происходит целыми блоками файлов — закрытыми сегментами лога (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-кода, не дожидаясь срабатывания автоматических политик:
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.