Топик — это журнал: сегменты и срок хранения
Чему научишься
- читать таблицу сегментов партиции: где пишут, где закрыто, где начало журнала
- объяснять, почему «шесть часов» означает «не меньше шести»
- предсказывать, что увидит отставший консьюмер: молчаливую дыру или отказ
- входить в журнал по времени и знать, почему это почти так же дёшево, как чтение с конца
- подбирать размер сегмента и срок так, чтобы начало журнала легло куда надо
17:40. Один топик стёр всё старше шести часов, другой — ничего
Утром, расконсервировав уровень, ты в 04:40 включил двенадцать антенн на приём и поставил на топик signal_raw срок хранения шесть часов: места на уровне связи мало. Ещё ты велел резать журнал почаще. На диске журнал партиции лежит не одним файлом, а несколькими подряд — сегментами, и новый сегмент у тебя начинается каждые пятнадцать минут или каждые 900 мегабайт, что наступит раньше, — чтобы старое удалялось аккуратнее. Топик поднял наспех, на трёх партициях.
Сейчас 17:40. Читальный зал просит перечитать эфир с 11:00: у них в отчётах дыра между одиннадцатью и половиной первого. Ты ставишь консьюмера на нужный оффсет в партиции 0 — с настройкой «не перескакивать молча» — и получаешь не данные, а отказ: оффсета 456 000 больше нет, самый ранний доступный — 504 000, это 11:40. Одиннадцать утра было шесть часов сорок минут назад: эти кадры перешли шестичасовую границу сорок минут назад и уже стёрты.
Ты идёшь в соседний топик signal_ack — туда диспетчерская пишет ответы на заявки, и срок там такой же, шесть часов. И находишь сообщения тринадцатичасовой давности.
Один топик удалил всё старше шести часов — и заявку на 11:00 уже не исполнить. Соседний не удалил ничего. Ни один не соврал.
КВЕРИ: Шесть часов значит «не раньше». Про «ровно» речи не было.

Опорное числоТемп эфира
48 кадров в секунду — так эфир идёт в норме: двенадцать антенн, по 4 кадра в секунду с каждой.
Кадр весит в среднем 450 КБ, поэтому в секунду набегает 21,6 МБ, за минуту — 2 880 кадров, за час — 172 800, за сутки — 4 147 200. Эти числа ещё встретятся в курсе.
В этом уроке смотрим на партицию 0. Оффсеты у каждой партиции свои, а ключ source_id кладёт в партицию 0 пять антенн из двенадцати — поэтому в ней 20 кадров в секунду, 9 МБ/с. Почему ключ раскладывает антенны неровно, разберём в следующем уроке.
ТерминыЕщё три слова
| слово | что это |
|---|---|
| сообщение и кадр | одно и то же: в механике говорим «сообщение», а в сценах смены сообщение эфира зовётся кадром |
| брокер | машина кластера Kafka, которая держит партиции на своём диске; RabbitMQ и родня — «брокер очередей» |
| потребитель эфира | внешняя система, которой уходит результат: Читальный зал и его соседи. Консьюмер — программа, потребитель эфира — тот, кто пользуется результатом |
Английские имена — во врезке про собеседование в конце урока.
signal_raw и что с ними сделала уборка. Вторая: что увидит консьюмер, который просит эфир с 11:00, — при двух разных настройках. Третья: редкий signal_ack под тем же сроком. Смотри на колонку «время» первого сегмента.Почему «шесть часов» значит «не меньше шести»
Журнал партиции — не один файл, а сегменты подряд: каждый сегмент — отдельный файл на диске. Пишут всегда только в последний, активный сегмент. Когда он дорастает до предела по размеру или по времени, его закрывают и открывают новый.
Удалить из середины файла одно сообщение нельзя. Поэтому срок хранения работает целыми сегментами, по трём правилам:
- удаляют только целый сегмент — ни одного сообщения по отдельности;
- решают по самому свежему сообщению в сегменте: сегмент уходит, когда даже оно старше срока;
- активный сегмент не трогают, пока в нём есть свежее сообщение.
Отсюда оба симптома. В signal_raw сегмент закрылся не через пятнадцать минут, а через 100 секунд: 900 МБ при 9 МБ/с набираются быстрее. На плотном потоке сегменты режет размер. Сегменты мелкие, поэтому срез ложится вплотную к шести часам: удалено 252 сегмента, самое старое уцелевшее сообщение ровно шестичасовое.
В signal_ack всего четыре ответа в час, а сегмент по умолчанию закрывается раз в семь суток — за тринадцать часов он не закрылся ни разу. Он всё ещё активный, и в нём есть свежее сообщение, поэтому его не трогают: отсюда тринадцать часов вместо шести. А если бы писать в него перестали и он устарел целиком, брокер открыл бы новый пустой сегмент и удалил старый — это нижний ряд на диаграмме ниже.
Оффсет самого старого сообщения, которое ещё можно прочитать, — начало журнала партиции (log start offset). Оно своё у каждой партиции, как и сами оффсеты: в партиции 0 после уборки оно уехало вперёд с 0 на 504 000.
Что видит отставший консьюмер
Консьюмер попросил 11:00 — а этих сообщений уже нет. Что он увидит, решает одна настройка консьюмера, а не брокер:
- с
auto_offset_reset='earliest'консьюмер молча перепрыгивает на начало журнала партиции и получает дыру 11:00–11:40 без единой ошибки: данные потеряны, а никто об этом не знает; - с
'latest'— а так по умолчанию настроены и kafka-python, и клиент на Java — он прыгает в конец и пропускает всё, что лежало между. Тоже молча; - с
'none'он получает отказ на первом же чтении: оффсета нет, и об этом тебе говорят, а не молчат.
Как перечитать окно из прошлого и когда эта настройка вообще срабатывает — разбор в главе 3.
Как войти в журнал по времени
Границы времени каждого сегмента брокер держит в памяти, поэтому нужный сегмент он находит, не открывая журнал. А рядом с каждым сегментом лежат два маленьких файла-указателя: один по времени находит оффсет, другой по оффсету — место в файле. Поэтому попросить «с 11:00» почти так же дёшево, как читать с конца: журнал подряд не перебирают ни там, ни там, и перечитывание вчерашнего — штатная операция, а не авария.
SEGMENT_MS (и при желании RETENTION_MS) для signal_ack так, чтобы после уборки в 17:40 самое старое доступное сообщение было не старше семи часов и не свежее шести. Ответы пишутся раз в пятнадцать минут, с 04:40 до 17:25, в одну партицию.
Функция oldest_age(segment_ms, retention_ms) уже подключена к редактору, её код скрыт: она заводит топик, пишет 52 ответа, запускает уборку в 17:40 и возвращает возраст самого старого уцелевшего ответа в миллисекундах — или None, если не уцелел ни один. Нижние строки заготовки вызывают её и печатают возраст — их вывод виден, когда жмёшь «Проверить». Константы MIN и HOUR (минута и час в миллисекундах) уже заданы.Что в этом уроке настоящее, а что учебное
В учебной Kafka сегмент — счётчик и пара границ в памяти вкладки, а не файл на диске, и его размер считается по учебному весу кадра в 450 КБ. Сами кадры ячейка не хранит: она их считает и создаёт только те, что ты попросил прочитать. На настоящем кластере за эти тринадцать часов эфир записал бы около терабайта.
Уборку в песочнице запускаешь ты сам, вызовом cluster.apply_retention(...). Настоящий брокер делает это сам, в фоне и по своему расписанию. Поэтому в граница сдвигается не в тот момент, когда ты посмотрел, а когда до неё дошла очередь, — ещё одна причина, по которой «шесть часов» значит «не меньше шести».
Настоящее — правила, и они переносятся на боевой кластер без оговорок: оффсет живёт внутри партиции, удаляют целым сегментом, решают по самому свежему сообщению, активный сегмент не трогают, начало журнала партиции двигается, а что увидит отставший консьюмер — дыру или отказ, — решает его настройка.
Вопрос с собеседования
Как это спрашивают на собеседовании
«Поставили .ms на 6 часов, а данные лежат 13. Почему?» Ждут три слова: сегмент (log segment), активный сегмент (active segment) и то, что решение принимают по самому свежему сообщению сегмента. Бонус — сказать, что уборщик ходит по расписанию, а не в момент истечения срока.
«Что будет с консьюмером, если закоммиченный оффсет его группы указывает на уже удалённое?» Правильный ответ начинается со слов «зависит от auto.offset.reset»: earliest — молчаливая дыра, latest — молчаливый пропуск, none — ошибка. И отдельно: умолчание клиентов — latest, это стоит знать наизусть.
«Как Kafka ищет сообщение по времени?» Через два индекса рядом с каждым сегментом: времени (time index) даёт оффсет, индекс оффсетов (offset index) — место в файле. Поэтому offsetsForTimes не сканирует журнал. Имена, которые прозвучат: log, retention.ms, segment.ms, segment.bytes, log start offset.
Главное из урока
| топик | сегменты | срок | удалено | самому старому |
|---|---|---|---|---|
signal_raw, партиция 0 | по размеру, 100 с | 6 ч | 252 сегмента, начало журнала 0 → 504 000 | 6 ч 00 мин |
signal_ack | по умолчанию, один активный | 6 ч | ничего | 13 ч 00 мин |
Срок хранения работает целыми сегментами и решает по самому свежему сообщению в сегменте; активный сегмент не трогают, пока в нём есть свежее сообщение. Поэтому «шесть часов» значит «не меньше шести». Что увидит отставший консьюмер — дыру или отказ — решает его настройка. А вход в журнал по времени дешёвый: рядом с каждым сегментом лежат указатели.
КВЕРИ: Ты резал журнал по пятнадцать минут. Эфир порезал его по сто секунд — сам.