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

Компакция: по ключу остаётся последнее

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

19:10. Справочник на одну строку

2 ноября 2184 года, диспетчерская «Хранилища-9». Читальному залу нужен справочник антенн: двенадцать строк, s01…s12, у каждой — сектор неба и статус. Такой справочник на уровне есть — топик sources. Его вёл прошлый сменщик, и утром ты расконсервировал его вместе с уровнем. За годы в нём накопилось 2,1 миллиона сообщений на эти двенадцать ключей: антенны перенастраивают, статусы моргают, и каждое изменение уходит сообщением с ключом source_id. Срок хранения у топика — «вечно».

Новый консьюмер, которому нужен справочник, читает топик с начала: два с лишним миллиона сообщений ради двенадцати строк, три минуты только на старте. И ты делаешь ровно то, чему научился в 17:40: ставишь топику срок хранения тридцать суток и запускаешь уборку по сроку.

Через несколько минут Читальный зал видит в справочнике одну антенну — s07. Остальные одиннадцать никуда не делись: они просто полгода не менялись, и последнее сообщение о каждой оказалось старше тридцати суток. Срок хранения удалил их целыми сегментами — как ему и положено.

Ты не удалил ни одного сообщения по ошибке. Ты сделал всё по правилам — и потерял одиннадцать строк из двенадцати.

КВЕРИ: Срок хранения умеет спросить только одно: сколько сообщению лет. Справочник интересует другое.

Голографическая справочная доска: двенадцать карточек в два ряда, одиннадцать погасли, светится одна. На доске сидит КВЕРИ, кот-интерфейс архива, и смотрит на неё сверху. За доской стеллажи катушек ленты. Дежурный со спины стоит перед доской с планшетом в руке.
Антенны никуда не делись — погасли только их карточки в справочнике. А на катушках за доской горит по одному витку: справочнику от каждой антенны нужно ровно столько.
Воспроизведём аварию своими руками. Ячейка работает на уменьшенной копии sources: 2 002 сообщения на те же двенадцать ключей вместо 2,1 миллиона — иначе таблица не влезет на экран, а вывод от масштаба не зависит. Сегмент — 500 сообщений. Один и тот же поток льётся в два топика, разница только в политике: sources получает срок хранения 30 суток, а sources_demo заводится с другой политикой — она видна в коде. Третья печать удаляет из sources_demo один ключ. Что делает вторая политика, разберём сразу после ячейки.
python · kafka

Два вопроса к топику

Первая печать — авария из сцены в миниатюре: удалены два сегмента, ключей было 12, осталось 1. Срок хранения сработал правильно. Ошибка в другом: топику задали не тот вопрос.

К топику бывает два разных вопроса:

  • «что происходило?» — это история. Каждое сообщение ценно само по себе, а старое можно выбросить, когда оно никому не нужно. Эфир signal_raw — такой, и ему нужен срок хранения;
  • «как обстоит дело сейчас?» — это состояние. Из всех сообщений об антенне важно одно, последнее, и его возраст ничего не значит: антенна, которая полгода не менялась, от этого не исчезла. Справочник sources — такой.

Справочнику не нужна вся история — ему нужно последнее значение каждого ключа. Эта фраза почти дословно и есть лекарство. Компакция (log compaction): из всех сообщений с одним ключом оставляют последнее, а предыдущие выбрасывают, и журнал превращается в справочник.

Компакцию включают топику одной настройкой, и такой топик зовут компактным. (То же правило — «последняя версия на ключ» — держат таблицы ReplacingMergeTree в ClickHouse.)

Компакцию делаешь не ты. По закрытым сегментам ходит фоновая уборка брокера — чистильщик. От каждого ключа он оставляет одно сообщение — последнее во всех закрытых сегментах, а не в каждом по отдельности. В активный сегмент он не заходит.

Что показала вторая печать

  • Масштаб. Сообщений было 2 002, осталось 14, ключей как было 12: последние значения одиннадцати антенн в закрытых сегментах, а у s07 — её последнее значение среди закрытых сегментов (оффсет 1999) и ещё два сообщения в активном сегменте. Два миллиона сообщений настоящего топика тоже живут ради двенадцати строк.
  • Оффсеты не переназначают. Уцелевшие стоят на прежних оффсетах — 988, 989, 991… и сразу 1999, а всё, что было между, исчезло навсегда. Это дыры в оффсетах. Начало журнала партиции осталось 0, но первым консьюмер получит оффсет 988. Код, который считает, что за оффсетом n всегда идёт n + 1, сломается на первом же компактном топике. Дыры перешагивает сам poll(): позиция консьюмера проходит сквозь них и только растёт.
  • Активный сегмент не тронут. Оффсеты 1999, 2000 и 2001 — три сообщения ключа s07, и все на месте: чистильщик не смотрит в активный сегмент и не знает, что сообщение 1999 уже устарело. Поэтому «в компактном топике одно сообщение на ключ» неверно — верно «со временем».
Шестнадцать сообщений, четыре ключа. Чистильщик прошёл по закрытым сегментам, оффсеты 0–11: от каждого ключа осталось последнее — D на оффсете 3, C на 6, B на 10, A на 11, а оффсеты 0–2, 4, 5 и 7–9 исчезли навсегда. Активный сегмент, оффсеты 12–15, не тронут: повторы ключей в нём на месте. Ключей как было четыре.

Как удалить ключ, если остаётся последнее

Правило компакции хранит последнее сообщение ключа сколько угодно долго. Как тогда сказать «такой антенны больше нет»? Отправить tombstone — сообщение с ключом и без тела: p.send('sources_demo', None, key='s11'). Для чистильщика оно и есть последнее сообщение ключа, поэтому все прежние значения s11 он выбросит. А консьюмер, который собирает справочник, увидев пустое тело, убирает ключ у себя.

Третья печать сделана на копии — настоящую s11 никто не списывал. В ней видно два шага:

  • в 19:10 tombstone лёг в активный сегмент, и чистильщику до него пока нет дела. В 20:10 s07 прислала новое значение, и сегмент с tombstone закрылся. Только тогда компакция его увидела и убрала три сообщения — прежнее значение s11 и два старых значения s07;
  • s11 в справочнике больше нет, ключей 11. Но сам tombstone в журнале лежит, и у него есть срок: 24 часа.

Этот срок — окно жизни удаления: столько tombstone ещё лежит в журнале, чтобы отставшие консьюмеры успели его увидеть. Отсчёт идёт от компакции, которая впервые застала его в закрытом сегменте. Когда окно закроется, следующая компакция выбросит и его — о ключе в журнале не останется ни слова.

Отсюда цена. Консьюмер, у которого s11 уже лежит в собранном справочнике и который ушёл в простой дольше окна, про удаление не узнает никогда: вернётся — а tombstone уже нет, и s11 останется у него навсегда.

Практика: напиши код
Напиши build_directory(consumer): консьюмер проходит компактный топик sources_demo с самого начала до конца, а функция возвращает справочник — словарь антенна → тело. Топик тот же, что после третьей печати ячейки: оффсеты идут с дырами, а ключ s11 удалён через tombstone. Удалённого ключа в справочнике быть не должно — ни с каким значением, даже с None. Консьюмера даёт make_consumer(): он в новой группе и начнёт с начала журнала партиции — в учебной Kafka это умолчание, настоящему консьюмеру для этого пишут auto_offset_reset='earliest'. Проверка соберёт справочник дважды, двумя новыми консьюмерами, и ждёт один и тот же ответ.
python · kafka
Окно в 24 часа отсчитывается не от записи tombstone, а от компакции, которая застала его в закрытом сегменте. Журнал у всех консьюмеров один, а справочники разошлись — и по топику этого не видно.
Проверь себя
Продюсер по ошибке отправляет в компактный топик sources_demo сообщение без ключа: p.send('sources_demo', {'sector': 2, 'status': 'в работе'}). Что будет?

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

Компакция в песочнице — явный вызов cluster.compact(...). Настоящий чистильщик приходит сам и когда сочтёт нужным: берётся за журнал, только когда часть, ещё не прошедшая компакцию, перевалила порог, и может не приходить часами. Отсюда оговорка «со временем»: в любой момент в топике может лежать хвост, которого компакция ещё не касалась.

Учебные числа ячейки. Порог этой части занижен — 0,1 вместо умолчания 0,5, иначе на маленькой копии вторая компакция не случилась бы. Сегмент — ровно 500 сообщений, потому что все они одного веса. А окно жизни удаления в 24 часа — не учебное: это умолчание настоящей Kafka.

Сегмент по времени песочница закрывает по часам кластера, а настоящий брокер — по отметкам сообщений. Поэтому в ячейке сообщения 10:00 и 18:40 лежат в одном активном сегменте, а на боевом кластере история со старыми отметками легла бы иначе; правило компакции от этого не меняется.

Правило из вопроса про сообщение без ключа ячейка не печатает, но песочница держит его так же, как брокер: ответь на вопрос, а потом проверь ответ сам одной строкой.

Настоящие в уроке правила: последнее сообщение на ключ, активный сегмент компакция не трогает, оффсеты не переназначают, tombstone живёт ограниченное время — и правило из вопроса про сообщение без ключа.

Вопрос с собеседования

Как это спрашивают на собеседовании

Первый вопрос звучит почти всегда одинаково и почти всегда по-английски: «чем log compaction отличается от по времени и когда какую берут?» Отвечай не настройками, а вопросом к данным. Если топик отвечает на вопрос «что происходило» — это журнал событий, ему нужен срок хранения (retention.ms). Если на вопрос «как обстоит дело сейчас» — это состояние, ему нужна компакция (cleanup.policy=compact).

В вопросе зашита ловушка: многие отвечают «она убирает дубликаты». Нет: она оставляет последнее сообщение на ключ. Два одинаковых сообщения с разными ключами — не дубликаты, и никуда они не денутся.

Второй вопрос — «как удалить ключ из компактного топика?» Ждут не только слова tombstone — сообщение с ключом и пустым телом, — но и оговорку: оно живёт ограниченное время (delete.retention.ms, по умолчанию сутки), и консьюмер, отставший дольше, удаление пропустит.

Третий, самый зрелый: «может ли компактный топик заменить справочник в базе?» Да — и хороший кандидат называет цену: историю изменений вы потеряли. Если она нужна, топиков должно быть два. Ещё имена, которые прозвучат: log cleaner, min.cleanable.dirty.ratio, active segment и non-contiguous offsets — оффсеты с дырами.

Закрепление: реши задачи
Решено 0 из 1
Главное из урока
топикполитикасообщенийключей
sourcesсрок хранения 30 суток2 002 → 1 00212 → 1
sources_demoкомпакция2 002 → 1412 → 12
sources_demo, удалена s11компакция13, вместе с tombstone11

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

Настоящий sources ты переводишь на компакцию и заливаешь двенадцать строк заново: вернуть стёртое сроком не умеет ни одна политика.

КВЕРИ: Журнал помнит всё. Справочнику хватает последнего слова каждой антенны.