Компакция: по ключу остаётся последнее
Чему научишься
- различать два вопроса к топику: «что происходило» и «как обстоит дело сейчас»
- применять правило компакции: от каждого ключа остаётся последнее сообщение
- читать журнал с дырами в оффсетах и не падать на них
- удалять ключ сообщением без тела — и помнить, что у такого сообщения есть срок
- собирать справочник из компактного топика одним проходом консьюмера
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 один ключ. Что делает вторая политика, разберём сразу после ячейки.Два вопроса к топику
Первая печать — авария из сцены в миниатюре: удалены два сегмента, ключей было 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 уже устарело. Поэтому «в компактном топике одно сообщение на ключ» неверно — верно «со временем».
Как удалить ключ, если остаётся последнее
Правило компакции хранит последнее сообщение ключа сколько угодно долго. Как тогда сказать «такой антенны больше нет»? Отправить 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'. Проверка соберёт справочник дважды, двумя новыми консьюмерами, и ждёт один и тот же ответ.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 — оффсеты с дырами.
Главное из урока
| топик | политика | сообщений | ключей |
|---|---|---|---|
sources | срок хранения 30 суток | 2 002 → 1 002 | 12 → 1 |
sources_demo | компакция | 2 002 → 14 | 12 → 12 |
sources_demo, удалена s11 | компакция | 13, вместе с tombstone | 11 |
Срок хранения отвечает на вопрос «что происходило», компакция — на вопрос «как обстоит дело сейчас»: от каждого ключа остаётся последнее сообщение. Чистильщик ходит только по закрытым сегментам и оффсеты не переназначает — отсюда дыры в оффсетах. Ключ удаляют, отправив tombstone — сообщение без тела; живёт он ограниченно — окно жизни удаления: консьюмер, опоздавший дольше, удаление пропустит. Справочник из компактного топика собирают одним проходом консьюмера с начала.
Настоящий sources ты переводишь на компакцию и заливаешь двенадцать строк заново: вернуть стёртое сроком не умеет ни одна политика.
КВЕРИ: Журнал помнит всё. Справочнику хватает последнего слова каждой антенны.