Записать в лог и прочитать по оффсету
1. Суть
Производитель отдаёт брокеру батч записей и получает оффсет первой из них. Потребитель называет оффсет и получает всё, что за ним лежит, до заданного предела в байтах. Оффсет — целое число, монотонное и без пропусков внутри партиции; помнит его потребитель, брокер за него не отвечает.
Это весь брокер. Всё остальное в этом проекте — сегменты, индекс, zero-copy, актор — существует только ради того, чтобы эти две операции были быстрыми.
Что реализовано на сегодня: всё. Хранилище (сегмент, разреженный индекс, партиция с роллингом
и retention, восстановление после рестарта, актор записи с батчами и групповым коммитом), сеть
(свой селектор, бинарный протокол, FETCH мимо кучи) и клиент (конвейерное соединение, Producer
с накоплением, Consumer с собственной позицией). Нет конфигурации, метрик и исполняемого
артефакта — это M5.
2. Бизнес-ограничения
- Оффсеты внутри партиции монотонны и без пропусков:
n-я принятая запись получаетbaseOffset + n. Проверяемо: послеkзаписейnextOffset == baseOffset + k. - Запись, однажды получившая оффсет, не меняется и не переезжает. Байты по позиции
pнеизменны навсегда. - Оффсет становится видимым только после того, как байты записаны. Читатель не может увидеть оффсет, за которым ничего нет.
- Запись, не помещающаяся в сегмент, отвергается, а не обрезается.
- Запись не долговечна, пока не вызван
force(). Это относится к обоим путям записи — через маппинг данные так же теряются при отключении питания, просто переживают падение процесса. Формулировка «ОС сама сбросит» гарантией не является. - Чтение за пределами записанного возвращает «нет данных», никогда — мусор из хвоста предразмеченного файла.
3. Как это работает
Сегодня (в коде):
LogSegment.append(payload)→SegmentWriter.appendкладёт[int32 length][payload]и возвращает позицию →SparseOffsetIndex.appendрешает, ставить ли отметку →nextOffsetувеличивается последним.LogSegment.read(offset)→positionOfищет в индексе ближайшую отметку не дальше цели → идёт вперёд по префиксам длин → читает тело.LogSegment.transferTo(from, maxBytes, target)отдаёт сырые байты в канал мимо heap.
По сети (M3–M4): сессия разбирает кадр → PartitionWriter принимает батч и отвечает
baseOffset; FETCH находит сегмент по оффсету → пишет заголовок ответа → transferTo тела мимо
кучи. Контракт — protocol-wire.
У клиента (M4): Producer копит записи и отправляет батчами, Consumer держит позицию и двигает
её только за целые записи.
4. Якоря кода
| Модуль | Код |
|---|---|
| booblik-core | src/main/kotlin/ru/workinprogress/booblik/storage/LogSegment.kt — сборка всего вместе |
| booblik-core | .../storage/SegmentWriter.kt — контракт записи и правило про долговечность |
| booblik-core | .../storage/SparseOffsetIndex.kt — поиск позиции по оффсету |
| booblik-core | src/test/kotlin/.../storage/LogSegmentTest.kt — сценарии ниже |
| booblik-net | src/main/kotlin/.../net/Session.kt — PRODUCE и FETCH на проводе |
| booblik-net | .../net/client/Producer.kt, .../net/client/Consumer.kt — как этим пользуются |
5. Сценарии (BDD)
Сценарий: оффсеты монотонны и без пропусков
- Дано: пустой сегмент с
baseOffset = 0. - Когда: записано 100 записей подряд.
- Тогда: выданные оффсеты — ровно
0…99,nextOffset == 100. - Автоматизирован:
LogSegmentTest.append assigns gap-free offsets from the base offset(оба режима записи).
Сценарий: прочитанное равно записанному
- Дано: 500 записей разной длины.
- Когда: каждая читается по своему оффсету.
- Тогда: байты совпадают с записанными.
- Автоматизирован:
LogSegmentTest.record read back equals record written(оба режима).
Сценарий: разреженный индекс не промахивается мимо записи
- Дано: 5 000 записей по 64 байта — заведомо больше одного интервала индекса, то есть поиск обязан идти вперёд от отметки, а не от нуля.
- Когда: запрашивается позиция записи 4321 и тело записи 4999.
- Тогда: позиция равна
4321 × (4 + 64), тело совпадает. - Автоматизирован:
LogSegmentTest.sparse index still lands on the exact record after a forward scan.
Сценарий: чтение за концом лога — «нет данных», а не мусор
- Дано: сегмент с одной записью.
- Когда: запрашивается оффсет 1 и оффсет 7.
- Тогда:
readвозвращаетnull,positionOfвозвращаетnull. Существенно дляMAPPED: файл там предразмечен на всю ёмкость и полон нулей, так что «прочитать хвост» технически возможно — и запрещено. - Автоматизирован:
LogSegmentTest.reading past the last written offset yields null rather than garbage.
Сценарий: переполненный сегмент отказывает, а не затирает
- Дано: сегмент ёмкостью 1024 байта, записи по 200 байт.
- Когда: пишем, пока
hasRoomForне скажет «нет». - Тогда: принято ровно 5 записей (
1024 / (4 + 200)), шестая не принимается. - Автоматизирован:
LogSegmentTest.a full segment refuses the append instead of overwriting.
Сценарий: байты уезжают в канал в точности как на диске
- Дано: сегмент с одной записью.
- Когда: содержимое перекачивается
transferToв цикле по готовности. - Тогда: на выходе
[int32 length][payload]побайтово, суммарно4 + length. - Автоматизирован:
LogSegmentTest.transferTo hands over the framed bytes verbatim.
Сценарий: батч получает одно подтверждение и оффсет первой записи
- Дано: партиция с
logEndOffset = N. - Когда: производитель отдаёт батч из
kзаписей. - Тогда: возвращается
N, записи получилиN…N+k-1,nextOffsetсталN+k. - Автоматизирован:
PartitionWriterTest.a batch gets one acknowledgement and consecutive offsets.
Сценарий: параллельные производители не сталкиваются оффсетами
- Дано: 64 производителя по 25 батчей из 4 записей.
- Когда: все пишут одновременно.
- Тогда: выданные базовые оффсеты различны и без пропусков; итог — ровно 6400 записей.
- Автоматизирован:
PartitionWriterTest.concurrent producers get unique, gap-free offsets.
Сценарий: групповой коммит разменивает барьеры на ожидание
- Дано: 200 производителей с политикой
FORCED. - Когда: все просят долговечности одновременно.
- Тогда: барьеров меньше, чем производителей, на порядок — каждый ждёт один барьер, а не свой.
- Автоматизирован:
PartitionWriterTest.group commit amortises one barrier across many producers.
Сценарий: лог переживает рестарт и продолжается с того же места
- Дано: партиция из трёх сегментов, 12 записей.
- Когда: партиция закрыта и открыта заново.
- Тогда: видны все три сегмента,
nextOffsetтот же, следующая запись получает12. - Автоматизирован:
RecoveryTest.reopening a partition restores every segment and keeps their order.
Сценарий: оборванная запись не возвращается из мёртвых
- Дано: сегмент, в котором последняя запись не дописана.
- Когда: сегмент открывают заново.
- Тогда: её нет,
nextOffsetуказывает на последнюю целую границу, и следующая запись переиспользует этот оффсет. - Автоматизирован:
RecoveryTest.a half written record is discarded — …(по тесту на каждый путь записи: их обрывает по-разному, см. Quirks).
Сценарий: retention удаляет сегменты целиком и никогда активный
- Дано: три сегмента.
- Когда: запрошено удержание в пределах одного байта.
- Тогда: остаётся один — активный; файлы удалённых отвязаны;
logStartOffsetвырос. - Автоматизирован:
PartitionLogTest.retention by size drops whole segments and never the active one.
Сценарий: FETCH за границей лога
- Дано:
highWatermark = N. - Когда: приходит FETCH с
fetchOffset > N. - Тогда: ответ с
errorCode = 2(OFFSET_OUT_OF_RANGE), тело пустое. РовноN— не ошибка: так выглядит догнавший потребитель, и он получает пустой батч. - Автоматизирован:
ServerTest.fetching past the end is an error, fetching exactly at the end is not.
Сценарий: FETCH обрывается на середине записи
- Дано:
maxBytesменьше, чем суммарный размер записей послеfetchOffset. - Когда: приходит FETCH.
- Тогда: тело обрывается по границе байтов, не по границе записи; клиент отбрасывает незавершённый хвост и сдвигает позицию только за целые записи, так что следующий запрос начинается с той же обрезанной записи. Брокер батч не разбирает — именно этого разбора избегает zero-copy.
- Автоматизирован:
ServerTest.maxBytes cuts on a byte boundary and the client drops the partial tail,ClientTest.a consumer advances only past whole records when a response is cut.
Сценарий: запись крупнее maxBytes не приезжает никогда
- Дано: первая же запись после
fetchOffsetбольше, чемmaxBytes. - Когда: приходит FETCH.
- Тогда: в ответе нет ни одной целой записи и есть обрезанный хвост. Отбросить хвост
и повторить — правильно во всех остальных случаях — здесь означает делать тот же запрос вечно,
поэтому клиент обязан сообщить:
RecordExceedsMaxBytesExceptionс размером записи и текущимmaxBytes. Молчание неотличимо от «догнал» (M-139). - Автоматизирован:
ClientTest.a record too large for maxBytes is reported rather than looking like the end of the log,ClientTest.a subscription reports the record it cannot fit instead of following nothing, плюс проверка комплекта соответствия во всех шести клиентах.
Сценарий: несколько запросов в полёте не путают ответы
- Дано: одно соединение и пятьдесят одновременных вызовов.
- Когда: все ждут ответа.
- Тогда: каждый получает оффсет своего батча; идентификаторы уникальны.
- Автоматизирован:
ClientTest.many requests in flight get their own answers back.
Сценарий: партиции не делят оффсеты
- Дано: топик из трёх партиций и второй топик.
- Когда: в каждую пишут по батчу.
- Тогда: все начинают с нуля и ничего друг о друге не знают.
- Автоматизирован:
BrokerTest.partitions have independent offsets.
Сценарий: брокер, которого просят о неизвестном, отвечает, а не догадывается
- Дано: брокер с топиком
ordersиз трёх партиций. - Когда: приходит запрос к партиции 3, к чужому топику, с неизвестным
apiKeyили версией. - Тогда:
UNKNOWN_TOPIC_OR_PARTITIONв первых двух случаях,UNSUPPORTED_VERSIONв остальных — и всегда сcorrelationIdзапроса, если заголовок разобрался. - Автоматизирован:
BrokerTest,ErrorCodeTest.
6. Что не входит в скоуп
- Репликация и вообще кластер. Один процесс. Всё, что делает Kafka эксплуатируемой в кластере, — это и есть та сложность, от которой проект отказался.
- Создание топиков. Набор партиций задаётся при старте. Создавать их на лету — значит договариваться о состоянии, а это первый шаг к координатору.
- Группы потребителей и хранение их позиций. Позицию помнит клиент.
- TLS и сжатие. Несовместимы с zero-copy (ресёрч §1.2), а не «не успели».
- Удаление по retention. Появится не раньше, чем появится второй сегмент (M2).
- Транзакции и exactly-once. Никогда.
7. Известные особенности (Quirks)
MAPPEDсоздаёт файл на всю ёмкость сразу. Свежий сегмент на 512 МиБ — это файл 512 МиБ, где ничего нет. На APFS и ext4 он разрежённый, ноls -lпоказывает полный размер, а любая сигнализация по месту на диске читает именно его.nextOffsetприсваивается после записи байтов, и это не стилистика. Поменяв порядок, получим окно, в котором читатель видит оффсет, за которым ещё ничего нет — гонку, которая воспроизводится раз в миллион записей и выглядит как повреждённый лог.- Два пути записи рвутся по-разному, и восстановление ловит их разными признаками. У
FILE_CHANNELграницу задаёт длина файла: заголовок, обещающий больше, чем в файле есть, — доказуемо недописанная запись. УMAPPEDдлина файла не значит ничего (он предразмечен), и единственный маркер — нулевой префикс; поэтому там тело пишется раньше префикса, чтобы крах между двумя записями оставлял ноль, а не правдоподобный заголовок. - Пустых записей не бывает. Нулевая длина занята под маркер конца лога.
FORCEDвMAPPEDстоит примерно столько же, сколько вFILE_CHANNEL. Раньше он выглядел в 63 раза дешевле, но это былmsync, который барьером долговечности не является (ресёрч §1.9). Теперь зовутся оба вызова. Если кто-то захочет «вернуть быстро» — он вернёт не скорость, а более слабое обещание под тем же именем.- Средняя пропускная способность
MAPPEDвыше, а худшая секунда — хуже. На минутной дистанции маппинг гуляет в 15 раз (128 тыс…1,97 млн зап/с),FILE_CHANNELдержится в пределах 20 % (204…255 тыс). Планировать очереди и таймауты приходится по худшей секунде, и по ней маппинг проигрывает (ресёрч §1.11).