booblik
DocsResearch

booblik — сценарии использования и граница «в брокере / рядом»

Документ отвечает на один вопрос: что должно стать фичей booblik, а что живёт рядом — в клиентской библиотеке, в образце или вообще за пределами проекта.

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

Разбор ведётся по ресёрчу архитектуры: booblik — лог, а не очередь, и большая часть ответов вытекает из этого одного факта.


1. Проверенные факты

ФактГде проверено
Kafka перечисляет свои сценарии сама, и их семь: Messaging, Website Activity Tracking, Metrics, Log Aggregation, Stream Processing, Event Sourcing, Commit Logapache/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.

Образец делается в два слоя, и первый — не «упрощённая версия» второго, а самостоятельная вещь:

  1. Слой, который показывает booblik таким, какой он есть. Паблишер и потребители, follow() и replay(), позиции в OffsetStore на томе, перезапуск с продолжением. Это шесть сценариев из семи, и здесь нет ни одной строки, которая обходила бы отсутствующую фичу.
  2. Слой поверх: очередь задач как протокол. Заявки, аренда, переразача. С документом, который объясняет, почему это не фича брокера.

Зависимость — от опубликованного артефакта (io.github.youndie.booblik:booblik-client) и образа из GHCR, а не от project(...): образец тогда заодно постоянная версия той проверки, которая поймала неработоспособный 0.1.1.


6. Открытые вопросы

  1. Задержка взятия задачи в схеме (в) не измерена. Круг «заявка → дочитать до неё» упирается в долгий FETCH и в групповой коммит писателя. Число нужно снять на образце — иначе «очередь на логе» останется утверждением о форме, а не о применимости.
  2. Поведение при равном старте не проверено. Три воркера, увидев свободную задачу одновременно, напишут три заявки; двое проиграют и пойдут дальше. Сколько работы уходит в проигранные заявки при трёх воркерах и при тридцати — вопрос к образцу, а не к рассуждению.
  3. Retention лога заявок. Заявки нужны, пока задача в работе; сколько держать — вопрос настройки, и ошибка здесь тихая: слишком короткий retention удалит заявку раньше, чем истечёт аренда, и задача уйдёт второму исполнителю при живом первом.

On this page