booblik
DocsFeatures

Записать в лог и прочитать по оффсету

1. Суть

Производитель отдаёт брокеру батч записей и получает оффсет первой из них. Потребитель называет оффсет и получает всё, что за ним лежит, до заданного предела в байтах. Оффсет — целое число, монотонное и без пропусков внутри партиции; помнит его потребитель, брокер за него не отвечает.

Это весь брокер. Всё остальное в этом проекте — сегменты, индекс, zero-copy, актор — существует только ради того, чтобы эти две операции были быстрыми.

Что реализовано на сегодня: всё. Хранилище (сегмент, разреженный индекс, партиция с роллингом и retention, восстановление после рестарта, актор записи с батчами и групповым коммитом), сеть (свой селектор, бинарный протокол, FETCH мимо кучи) и клиент (конвейерное соединение, Producer с накоплением, Consumer с собственной позицией). Нет конфигурации, метрик и исполняемого артефакта — это M5.

2. Бизнес-ограничения

  • Оффсеты внутри партиции монотонны и без пропусков: n-я принятая запись получает baseOffset + n. Проверяемо: после k записей nextOffset == baseOffset + k.
  • Запись, однажды получившая оффсет, не меняется и не переезжает. Байты по позиции p неизменны навсегда.
  • Оффсет становится видимым только после того, как байты записаны. Читатель не может увидеть оффсет, за которым ничего нет.
  • Запись, не помещающаяся в сегмент, отвергается, а не обрезается.
  • Запись не долговечна, пока не вызван force(). Это относится к обоим путям записи — через маппинг данные так же теряются при отключении питания, просто переживают падение процесса. Формулировка «ОС сама сбросит» гарантией не является.
  • Чтение за пределами записанного возвращает «нет данных», никогда — мусор из хвоста предразмеченного файла.

3. Как это работает

Сегодня (в коде):

  1. LogSegment.append(payload)SegmentWriter.append кладёт [int32 length][payload] и возвращает позицию → SparseOffsetIndex.append решает, ставить ли отметку → nextOffset увеличивается последним.
  2. LogSegment.read(offset)positionOf ищет в индексе ближайшую отметку не дальше цели → идёт вперёд по префиксам длин → читает тело.
  3. LogSegment.transferTo(from, maxBytes, target) отдаёт сырые байты в канал мимо heap.

По сети (M3–M4): сессия разбирает кадр → PartitionWriter принимает батч и отвечает baseOffset; FETCH находит сегмент по оффсету → пишет заголовок ответа → transferTo тела мимо кучи. Контракт — protocol-wire.

У клиента (M4): Producer копит записи и отправляет батчами, Consumer держит позицию и двигает её только за целые записи.

4. Якоря кода

МодульКод
booblik-coresrc/main/kotlin/ru/workinprogress/booblik/storage/LogSegment.kt — сборка всего вместе
booblik-core.../storage/SegmentWriter.kt — контракт записи и правило про долговечность
booblik-core.../storage/SparseOffsetIndex.kt — поиск позиции по оффсету
booblik-coresrc/test/kotlin/.../storage/LogSegmentTest.kt — сценарии ниже
booblik-netsrc/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).

On this page

1. Суть2. Бизнес-ограничения3. Как это работает4. Якоря кода5. Сценарии (BDD)Сценарий: оффсеты монотонны и без пропусковСценарий: прочитанное равно записанномуСценарий: разреженный индекс не промахивается мимо записиСценарий: чтение за концом лога — «нет данных», а не мусорСценарий: переполненный сегмент отказывает, а не затираетСценарий: байты уезжают в канал в точности как на дискеСценарий: батч получает одно подтверждение и оффсет первой записиСценарий: параллельные производители не сталкиваются оффсетамиСценарий: групповой коммит разменивает барьеры на ожиданиеСценарий: лог переживает рестарт и продолжается с того же местаСценарий: оборванная запись не возвращается из мёртвыхСценарий: retention удаляет сегменты целиком и никогда активныйСценарий: FETCH за границей логаСценарий: FETCH обрывается на середине записиСценарий: запись крупнее maxBytes не приезжает никогдаСценарий: несколько запросов в полёте не путают ответыСценарий: партиции не делят оффсетыСценарий: брокер, которого просят о неизвестном, отвечает, а не догадывается6. Что не входит в скоуп7. Известные особенности (Quirks)