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

Как удаляются старые сообщения из топика

В 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) — важнейший защитный механизм распределенного чтения:

  1. Защита активных чтений консьюмеров (In-flight Reads): В момент, когда поток очистки решил удалить сегмент, какой-либо консьюмер или follower-реплика мог уже открыть файловый дескриптор и вычитывать данные через системный вызов sendfile(2).
  2. Переименование файла в *.deleted немедленно исключает сегмент из метаданных и маршрутов поиска новых запросов FetchRequest, но оставляет файл на диске.
  3. Клиенты, уже начавшие чтение, спокойно завершают передачу порции байт в течение 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)}\)

  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.

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