booblik
DocsServices

booblik-core

1. Зона ответственности

Хранение. Модуль владеет тем, как запись ложится на диск и как её оттуда достают: доменные идентификаторы, файл сегмента, индекс оффсетов, отдача байтов в произвольный канал.

Чем не занимается — эта половина важнее:

  • не знает про сеть, сокеты и протокол на проводе;
  • не знает про топики и партиции как про сущности с конфигурацией — оперирует одним сегментом;
  • не занимается репликацией, группами потребителей и хранением их позиций;
  • не решает, когда звать force(). Умеет звать; политику задаёт вызывающий (решение Р4).

Главный инвариант: файл сегмента append-only. Байты, однажды записанные по позиции p, не меняются никогда. На этом стоит всё остальное — читатель без блокировки, transferTo мимо heap, публикация nextOffset последним действием.

2. Контракт

Публичный Kotlin-API, HTTP тут нет. Формат записи на диске — [int32 payloadSize][int32 crc32c][payload], big-endian; он же уезжает в сокет как есть на zero-copy-пути, см. protocol-wire.

2а. Ключевые файлы (якоря кода)

ФайлЧто там
src/main/kotlin/ru/workinprogress/booblik/Ids.ktOffset, Position, TopicName, PartitionId
.../storage/Log.ktто, что видит писатель: append / force / nextOffset
.../storage/SegmentWriter.ktинтерфейс записи; здесь же правило про долговечность
.../storage/FileChannelSegmentWriter.ktзапись через FileChannel, gathering write
.../storage/MappedSegmentWriter.ktзапись через FFM-маппинг
.../storage/SparseOffsetIndex.ktразреженный индекс в LongArray
.../storage/LogSegment.ktсегмент целиком: запись, индекс, чтение, transferTo, восстановление
.../storage/PartitionLog.ktпартиция: роллинг, поиск сегмента по оффсету, retention
.../log/PartitionWriter.ktактор записи: батчи, групповой коммит
.../log/AckPolicy.ktтри обещания: NONE, WRITTEN, FORCED
.../log/FlushPolicy.ktфоновый сброс: ограничивает окно потери, но не даёт долговечности

3. Как устроено

Писатель и читатель — разные объекты над одним файлом, и они не синхронизируются. Читатель работает на своём FileChannel, писатель — на своём. Гонки нет по двум причинам: файл append-only, поэтому читатель может увидеть только префикс; и верхняя граница, до которой читателю разрешено смотреть, публикуется писателем через SegmentWriter.size. LogSegment.nextOffset присваивается после записи байтов — оффсет становится виден только тогда, когда за ним уже что-то есть.

Две записи за одним интерфейсом. FILE_CHANNEL и MAPPED (SegmentMode) — не настройка для пользователя, а предмет замера (решение Р1). Умолчание — MAPPED с M-45; FILE_CHANNEL остаётся рабочим режимом и путём отката, и данные читаются в обе стороны (ModeMigrationTest). Ничто выше SegmentWriter не должно уметь их различать, и тесты LogSegmentTest гоняют оба режима через один и тот же набор утверждений именно для того, чтобы это оставалось правдой.

Индекс разреженный. Поиск даёт запись не дальше 4 КиБ до цели, остальное — проход вперёд по префиксам длин. Почему не ConcurrentSkipListMap — ресёрч §1.6 и решение Р2.

Когерентность mmap и read(). Записанное через маппинг видно последующему read() на другом дескрипторе — это свойство единого буферного кэша Linux и macOS, а не JDK. Именно оно позволяет писать через маппинг, а читать плоским FileChannel. На системе без единого кэша (нам такие не нужны, но упомянуть стоит) режим MAPPED пришлось бы читать тоже через маппинг.

4. Зависимости

ТипИмяДля чего
Librarykotlinx-coroutines-coreактор записи и групповой коммит
JDKjava.lang.foreignмаппинг сегмента без потолка 2 ГБ и с детерминированным освобождением

5. Конфигурация

Конфигурационных файлов нет — всё параметрами конструкторов.

ПараметрДефолтСмысл
LogSegment.DEFAULT_CAPACITY512 МиБпотолок сегмента; жёсткая граница — Int.MAX_VALUE
SparseOffsetIndex.DEFAULT_INTERVAL_BYTES4 КиБсколько лога приходится на одну запись индекса
SparseOffsetIndex.DEFAULT_MAX_ENTRIES128 Кипотолок; массив растёт до него, а не выделяется сразу
SegmentModeMAPPEDкакой из двух путей записи; см. Р1 и M-45

6. Инфраструктура и деплой

Нечего деплоить: это библиотека внутри проекта. Публикации в Maven тоже нет и не планируется до того, как появится сеть.

7. Локальный запуск

./gradlew :booblik-core:test

8. Сознательные ограничения / грабли

  • Умолчание MAPPED означает, что сегмент занимает свою ёмкость на бумаге с первой записи. Свежая партиция — файл на 512 МиБ, в котором одна запись. Он разрежённый, место не занято, но ls и сигнализация по месту читают видимый размер. И -Xmx при этом не ограничивает маппинг вовсе, то есть флаги кучи ничего не говорят про rss процесса.
  • Записи нулевой длины запрещены, и это решение хранилища, а не валидация входа. Восстановление идёт по префиксам длин, а в предразмеченном mapped-файле за концом лога лежат нули — поэтому нулевая длина и есть маркер конца. Легальная пустая запись была бы от него неотличима.
  • В mapped-пути длина пишется ПОСЛЕ тела и суммы. Длина — то, что восстановление читает первым, поэтому она обязана появляться последней: если длина видна, всё остальное уже записано. Сама по себе эта очередность не доказательство (переупорядочение страниц никто не запрещал) — она решает лишь, где остановится восстановление; отличить целую запись от рваной умеет контрольная сумма.
  • Заголовок записи — [int32 длина][int32 crc32c]. Сумма считается при записи и проверяется при восстановлении и при обычном чтении. На zero-copy-пути она не проверяется и не может: он до байтов не дотрагивается. Там проверяет клиент.
  • Восстановление останавливается на первой несошедшейся записи и сохраняет всё до неё. За дырой оффсеты уже не значат того, что значили, поэтому «пропустить и продолжить» было бы хуже.
  • Индекс растёт по мере надобности, а не выделяется сразу. Раньше на каждый сегмент выделялся мегабайтный LongArray независимо от его размера, и сотня мелких сегментов убивала брокер на 64 МиБ. Считать размер от ёмкости сегмента — соблазнительно и неверно: записи бывают больше интервала индекса, поэтому ёмкость / интервал не верхняя граница, и индекс начинает заполняться раньше сегмента.
  • PositionInt, и это потолок сегмента, а не небрежность. И маппинг, и один вызов transferTo индексируются int (ресёрч §1.4, §1.5). Long здесь обещал бы диапазон, который хранилище всё равно не обслужит.
  • Запись int в маппинг идёт JAVA_INT_UNALIGNED. С JAVA_INT FFM падает на невыровненном смещении — а записи переменной длины почти всегда невыровнены (ресёрч §1.7). Выглядит как случайная константа; это не она.
  • force() в mapped-пути делает msync И fsync. Не перестраховка: замерено (ресёрч §1.9), что на APFS после msync последующий fsync стоит 79 % от полного — то есть там msync барьером не является; на ext4 он им является, и второй вызов стоит 1 % (замер 12.1 на настоящем ext4; на WSL2 было 2 %). Портируемый вариант — звать оба. Убрать второй вызов значит тихо ослабить FORCED на одном из двух путей записи.
  • Позиция FileChannel обязана совпадать с written. У gathering write нет позиционной перегрузки, он всегда пишет туда, где канал стоит сейчас. Забытый channel.position(...) при открытии восстановленного сегмента — это не ошибка, а тихая порча: первая же запись после рестарта затирала начало лога. Поймано RecoveryTest.
  • Маппинг не владеет каналом, из которого сделан. MappedSegmentWriter закрывает его сам — иначе брокер терял бы дескриптор на каждый сегмент, и заметно это стало бы только на роллинге.
  • Retention удаляет файл сразу, а дескрипторы — потом. Читатель, уже начавший сегмент, дочитывает его: на Linux и macOS отвязанный файл остаётся доступен через открытые дескрипторы. Считает читателей LogSegment.acquire/release.
  • retainNewerThan смотрит на mtime файла. В mapped-режиме время изменения обновляет writeback ядра, а не наши записи, поэтому возраст сегмента там менее предсказуем, чем в FILE_CHANNEL. Retention по размеру таких оговорок не имеет.

On this page