Зачем Kafka: журнал вместо отчёта

Топик — это журнал: сегменты и срок хранения

22 мин
Чему научишься
  • читать таблицу сегментов партиции: где пишут, где закрыто, где начало журнала
  • объяснять, почему «шесть часов» означает «не меньше шести»
  • предсказывать, что увидит отставший консьюмер: молчаливую дыру или отказ
  • входить в журнал по времени и знать, почему это почти так же дёшево, как чтение с конца
  • подбирать размер сегмента и срок так, чтобы начало журнала легло куда надо

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 и родня — «брокер очередей»
потребитель эфиравнешняя система, которой уходит результат: Читальный зал и его соседи. Консьюмер — программа, потребитель эфира — тот, кто пользуется результатом

Английские имена — во врезке про собеседование в конце урока.

Воспроизведём оба симптома. Ячейка заводит два топика с одним сроком хранения — шесть часов — и печатает три вещи. Первая: таблицу сегментов партиции 0 плотного signal_raw и что с ними сделала уборка. Вторая: что увидит консьюмер, который просит эфир с 11:00, — при двух разных настройках. Третья: редкий signal_ack под тем же сроком. Смотри на колонку «время» первого сегмента.
python · kafka

Почему «шесть часов» значит «не меньше шести»

Журнал партиции — не один файл, а сегменты подряд: каждый сегмент — отдельный файл на диске. Пишут всегда только в последний, активный сегмент. Когда он дорастает до предела по размеру или по времени, его закрывают и открывают новый.

Удалить из середины файла одно сообщение нельзя. Поэтому срок хранения работает целыми сегментами, по трём правилам:

  • удаляют только целый сегмент — ни одного сообщения по отдельности;
  • решают по самому свежему сообщению в сегменте: сегмент уходит, когда даже оно старше срока;
  • активный сегмент не трогают, пока в нём есть свежее сообщение.

Отсюда оба симптома. В 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 (минута и час в миллисекундах) уже заданы.
python · kafka
Ни один путь не читает журнал подряд: сегмент находят по границам времени в памяти брокера, место в нём — по маленьким файлам-указателям рядом с сегментом. Поэтому перечитать вчерашнее — штатная операция, а не авария.

Что в этом уроке настоящее, а что учебное

В учебной 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.

Проверь себя
Топик с небольшим потоком: сегмент закрывается по времени, раз в 3 часа, срок хранения — 6 часов, пишут в него постоянно. Каким может оказаться возраст самого старого доступного сообщения сразу после уборки?
Главное из урока
топиксегментысрокудаленосамому старому
signal_raw, партиция 0по размеру, 100 с6 ч252 сегмента, начало журнала 0 → 504 0006 ч 00 мин
signal_ackпо умолчанию, один активный6 чничего13 ч 00 мин

Срок хранения работает целыми сегментами и решает по самому свежему сообщению в сегменте; активный сегмент не трогают, пока в нём есть свежее сообщение. Поэтому «шесть часов» значит «не меньше шести». Что увидит отставший консьюмер — дыру или отказ — решает его настройка. А вход в журнал по времени дешёвый: рядом с каждым сегментом лежат указатели.

КВЕРИ: Ты резал журнал по пятнадцать минут. Эфир порезал его по сто секунд — сам.