booblik — сценарии использования и граница «в брокере / рядом»
Документ отвечает на один вопрос: что должно стать фичей booblik, а что живёт рядом — в клиентской библиотеке, в образце или вообще за пределами проекта.
Повод конкретный. Задумывался образец из паблишера и трёх сервисов, которые «по очереди разбирают» записи, и выяснилось, что за этой фразой стоит не деталь, а развилка архитектуры. Прежде чем писать код, надо разложить, какие вообще бывают сценарии и чего каждый требует от брокера.
Разбор ведётся по ресёрчу архитектуры: booblik — лог, а не очередь, и большая часть ответов вытекает из этого одного факта.
1. Проверенные факты
| Факт | Где проверено |
|---|---|
| Kafka перечисляет свои сценарии сама, и их семь: Messaging, Website Activity Tracking, Metrics, Log Aggregation, Stream Processing, Event Sourcing, Commit Log | apache/kafka, docs/getting-started/uses.md (trunk) |
| Партиция в группе потребителей читается ровно одним потребителем: «each of which is consumed by exactly one consumer within each subscribing consumer group at any given time» | apache/kafka, docs/design/design.md, раздел Consumer Position |
| Kafka прямо объясняет, чего стоит подтверждение на сообщение: «the broker must keep multiple states about every single message (first to lock it so it is not given out a second time, and then to mark it as permanently consumed so that it can be removed)» | там же |
Групп потребителей для очередей не хватило самой Kafka: заведены Share Groups — отдельный тип группы, где «partitions may be assigned to multiple consumers», «the number of consumers in a share group can exceed the number of partitions», «records are acknowledged individually», а у записи есть блокировка на 30 с по умолчанию (share.record.lock.duration.ms) с подсчётом попыток доставки | apache/kafka, docs/design/design.md, раздел The Share Consumer |
| Компактификация определена через ключ: «Log compaction ensures that Kafka will always retain at least the last known value for each message key»; удаление — это «a message with a key and a null payload» | apache/kafka, docs/design/design.md, раздел Log Compaction |
В booblik ключа на проводе нет: PRODUCE несёт partitionId, ackPolicy, recordCount и записи вида [int32 length][payload] | protocol-wire §3 |
| Ключ используется только на стороне клиента — для выбора партиции | booblik-client/.../client/Publishing.kt, TopicHandle.partitionFor |
OffsetStore объявлен и не реализован; единственные реализации — NoopStore в пробе и заглушка в тесте | Subscription.kt:232, поиск по репозиторию |
| Брокер не хранит позиций потребителей и не знает о них | Subscription.kt:37, Metrics.kt:25 |
Оговорка. Факты про Kafka взяты из исходников её документации в apache/kafka — это первоисточник,
но он описывает Kafka, а не индустрию. Утверждения вида «так делают все» ниже не встречаются: там,
где нужен общий обзор, он назван ориентиром, а не свидетельством.
2. Ось, вдоль которой раскладывается всё остальное
Разница между очередью и логом не в API, а в том, кто владеет состоянием доставки.
Очередь (RabbitMQ, SQS): состояние принадлежит брокеру. Он знает, что сообщение выдано, кому, до какого момента и подтверждено ли. Потребитель при этом почти без состояния — упал, и его работу переразметят. Цена записана у Kafka в разделе Consumer Position: брокер обязан держать несколько состояний на каждое сообщение — заблокировать, чтобы не выдать дважды, и пометить съеденным, чтобы удалить. Плюс отдельный вопрос, что делать с выданным и никогда не подтверждённым.
Лог (Kafka, booblik): состояние принадлежит потребителю и представляет собой одно число — оффсет. Брокер не знает ни кто читает, ни докуда дочитал. Отсюда всё остальное: дешёвое «подтверждение», перемотка назад, независимые читатели, zero-copy на отдаче (трогать байты незачем, потому что и состояния на них нет).
Из этого следует то, что сначала выглядит контринтуитивно: очередь задач с раздачей по одной записи — не «фича, которой пока нет», а другая модель хранения состояния. Это подтверждается со стороны Kafka: групп потребителей для таких нагрузок ей не хватило, и появились Share Groups с блокировкой на запись, индивидуальным подтверждением и счётчиком попыток — то есть ровно то состояние на каждое сообщение, от которого лог уходил.
И второе следствие, которое стоит держать в голове при слове «группы»: группа потребителей в Kafka не раздаёт записи — она раздаёт партиции, по одной на потребителя. Параллелизм упирается в число партиций, а медленная запись задерживает свою партицию целиком.
3. Сценарии
Список — семь из uses.md Kafka плюс восьмой, ради которого документ и затевался.
| Сценарий | Что требует от брокера | booblik сегодня |
|---|---|---|
| Messaging (развязка, буферизация) | долговечность, независимые читатели | есть |
| Website Activity Tracking | высокий поток записи, много читателей одного потока | есть |
| Metrics | то же плюс дешёвый append | есть |
| Log Aggregation | поток append, батчевые читатели | есть |
| Stream Processing (топик → топик) | чтение и запись из одного процесса, at-least-once | есть |
| Event Sourcing | долгое хранение, перемотка в начало | есть, если retention выключен |
| Commit Log / changelog | компактификация по ключу | невозможно — ключа нет на проводе |
| Очередь задач (competing consumers) | состояние на запись: блокировка, подтверждение, перераздача | нет, и это другая модель |
3.1 Что booblik уже умеет, и это не мало
Шесть сценариев из семи закрываются тем, что есть: follow() и replay(), позиция на стороне
потребителя, StartPosition.Earliest/Latest/At, долгий FETCH вместо опроса, батчи на записи.
Ни один из шести не требует ни координатора, ни хранения позиций на брокере: у каждого потребителя своя позиция и своя судьба. Именно это лог и продаёт.
3.2 Commit Log отпадает навсегда, и это следствие старого решения
Компактификация определена через ключ сообщения. У booblik ключ не доходит до брокера: он
существует только в клиенте, чтобы выбрать партицию (TopicHandle.partitionFor), а на провод уходит
[int32 length][payload].
Значит компактификацию нельзя добавить, не поменяв формат кадра и не заставив брокер разбирать записи — то есть не разрушив zero-copy-путь, ради которого писался свой сетевой слой (решение Р3). Это не «пока не сделано», а закрытая дверь, и лучше, чтобы она была закрыта явно.
Побочно это объясняет, почему в booblik нет и не будет темы «хранить позиции потребителей в самом
брокере, как Kafka в __consumer_offsets»: тот топик компактифицируется по ключу
(группа, топик, партиция), а компактификации у нас нет — некомпактифицируемый топик позиций рос бы
вечно.
3.3 Очередь задач
Требуется три вещи, и все три — состояние на запись: не выдать дважды, узнать, что обработана, переразать, если обработчик умер.
Три способа их получить:
(а) Сделать фичей брокера. Это Share Groups: блокировка записи, индивидуальное подтверждение, счётчик попыток. Возможно, но означает отменить решение, на котором стоит проект: брокер начинает держать состояние на каждое сообщение, а вместе с ним появляются запись при чтении, обход заблокированных и вопрос «что делать с выданным и брошенным». Отдача перестаёт быть чтением куска файла в сокет. Веха такого размера в проекте на 10 тысяч строк — это не веха, а второй проект.
(б) Координатор сбоку. Отдельный сервис раздаёт задачи по HTTP с арендой. Работает, объясняется за минуту, но в получившемся образце booblik — транспорт под координатором, а сам механизм очереди к брокеру отношения не имеет.
(в) Протокол поверх лога. Заявки пишутся в отдельную партицию, и побеждает та, что оказалась в логе первой. Партиция — тотальный порядок, одинаково видимый всем читателям, поэтому арбитра не нужно: его роль играет порядок. Аренда — поле со сроком в заявке; истёк срок без записи «готово» — задача снова свободна.
Что (в) честно не даёт:
- at-least-once, не exactly-once. Зависший дольше аренды воркер очнётся и доделает задачу, которую уже забрал другой. Ограждения (fencing) в схеме нет — для него нужна проверка на стороне, которая примет результат только от текущего арендатора;
- лог заявок растёт и живёт только на retention;
- каждый воркер читает все заявки — O(воркеры × задачи) трафика. Для образца безразлично, для тысячи воркеров — нет;
- задержка на взятие задачи — это круг «записал заявку → дочитал лог до неё». Очередь на таком протоколе не будет быстрой; она будет понятной.
4. Решения
Р11. Очередь задач не становится фичей брокера
Состояние на запись — это ровно то, от чего уходит лог, и Kafka подтверждает цену прямым текстом.
Отклонённая альтернатива: реализовать аналог Share Groups. Отклонена не потому, что сложно, а потому
что меняет предмет: получился бы другой брокер, и все замеры пришлось бы снимать заново — отдача
перестала бы быть transferTo из page cache в сокет.
Условие пересмотра: появление потребителя, которому нужна именно поштучная раздача, и готовность считать это отдельным проектом со своим ресёрчем и своими замерами.
Р12. Компактификация невозможна и записана как закрытая дверь
Не «не сделано», а исключено конструкцией: компактификация определена через ключ, ключа на проводе нет, а завести его — значит заставить брокер разбирать записи и потерять zero-copy. Из этого же следует отсутствие серверного хранения позиций.
Сценарий Commit Log / changelog объявляется вне области проекта. Это единственный из семи, который booblik не закрывает.
Р13. Очередь задач живёт в dev/ как протокол поверх лога
Способ (в): топик заявок, арбитр — порядок в логе, аренда по времени. Образец обязан называть свои гарантии: at-least-once, без ограждения, лог заявок растёт, все читают всё.
Отклонённая альтернатива — координатор сбоку (б): проще, но демонстрирует координатор, а не брокер.
Побочная польза: образцу нужна настоящая реализация OffsetStore — файл на томе. Интерфейс объявлен
и до сих пор не реализован ни разу за пределами заглушек, то есть заодно проверится, что его форма
вообще пригодна.
Р14. Совместимости с протоколом Kafka не будет; вместо неё — ретранслятор
Заговорить по-кафковски означает хранить по-кафковски. Заголовок RecordBatch у Kafka
(docs/implementation/message-format.md) начинается с baseOffset и batchLength, содержит
crc, который покрывает всё от attributes до конца батча, и дальше producerId,
baseSequence и прочее. То есть оффсет и контрольная сумма лежат внутри хранимых байтов —
именно поэтому Kafka может отдавать их sendfile.
У booblik ровно та же связка, но со своим форматом: тело ответа FETCH — «сырые записи сегмента, побайтово как на диске» (protocol-wire §4). Совместимость по протоколу поэтому означала бы одно из двух: принять формат хранения Kafka целиком — это замена слоя хранения вместе со всеми привязанными к нему замерами, — или собирать кафковские батчи в куче на каждом ответе, потеряв 2,5–2,8× процессора, измеренные в замере 13.
Отклонённая альтернатива — прокси-транслятор в отдельном процессе. Он оставил бы брокер чистым,
но даже producer-only требует ApiVersions, Metadata и Produce с разбором RecordBatch v2,
а consumer тянет координацию групп — то, чего нет по решению Р11. И «почти Kafka» ломается
не отказом, а странностями: клиент договаривается через ApiVersions и дальше пользуется тем,
что ему объявили.
Принято: ретранслятор. Отдельный сервис читает Kafka обычным kafka-clients и пишет в booblik
клиентом booblik, или наоборот. Никакой реализации протокола, две библиотеки по краям, оба формата
целы, а трансляция платится один раз в пользовательском коде. Это настоящий сценарий —
зеркалирование и миграция — и он ставит booblik рядом с экосистемой, а не притворяется ею.
Чего ретранслятор не переносит и перенести не может: ключ Kafka. На проводе booblik поля для ключа нет (Р12), поэтому по дороге в booblik ключ ещё выбирает партицию — порядок по ключу сохраняется, — но не сохраняется сам, и обратно его не восстановить. Заворачивать ключ в конверт означало бы заставить любого другого потребителя booblik разбирать конверт, которого он не просил.
5. Что из этого строим
Порядок и приёмка — в BACKLOG.md, веха M10.
Образец делается в два слоя, и первый — не «упрощённая версия» второго, а самостоятельная вещь:
- Слой, который показывает booblik таким, какой он есть. Паблишер и потребители,
follow()иreplay(), позиции вOffsetStoreна томе, перезапуск с продолжением. Это шесть сценариев из семи, и здесь нет ни одной строки, которая обходила бы отсутствующую фичу. - Слой поверх: очередь задач как протокол. Заявки, аренда, переразача. С документом, который объясняет, почему это не фича брокера.
Зависимость — от опубликованного артефакта (io.github.youndie.booblik:booblik-client) и образа
из GHCR, а не от project(...): образец тогда заодно постоянная версия той проверки, которая поймала
неработоспособный 0.1.1.
6. Открытые вопросы
- Задержка взятия задачи в схеме (в) не измерена. Круг «заявка → дочитать до неё» упирается в долгий FETCH и в групповой коммит писателя. Число нужно снять на образце — иначе «очередь на логе» останется утверждением о форме, а не о применимости.
- Поведение при равном старте не проверено. Три воркера, увидев свободную задачу одновременно, напишут три заявки; двое проиграют и пойдут дальше. Сколько работы уходит в проигранные заявки при трёх воркерах и при тридцати — вопрос к образцу, а не к рассуждению.
- Retention лога заявок. Заявки нужны, пока задача в работе; сколько держать — вопрос настройки, и ошибка здесь тихая: слишком короткий retention удалит заявку раньше, чем истечёт аренда, и задача уйдёт второму исполнителю при живом первом.