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.kt | Offset, 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. Зависимости
| Тип | Имя | Для чего |
|---|---|---|
| Library | kotlinx-coroutines-core | актор записи и групповой коммит |
| JDK | java.lang.foreign | маппинг сегмента без потолка 2 ГБ и с детерминированным освобождением |
5. Конфигурация
Конфигурационных файлов нет — всё параметрами конструкторов.
| Параметр | Дефолт | Смысл |
|---|---|---|
LogSegment.DEFAULT_CAPACITY | 512 МиБ | потолок сегмента; жёсткая граница — Int.MAX_VALUE |
SparseOffsetIndex.DEFAULT_INTERVAL_BYTES | 4 КиБ | сколько лога приходится на одну запись индекса |
SparseOffsetIndex.DEFAULT_MAX_ENTRIES | 128 Ки | потолок; массив растёт до него, а не выделяется сразу |
SegmentMode | MAPPED | какой из двух путей записи; см. Р1 и M-45 |
6. Инфраструктура и деплой
Нечего деплоить: это библиотека внутри проекта. Публикации в Maven тоже нет и не планируется до того, как появится сеть.
7. Локальный запуск
./gradlew :booblik-core:test8. Сознательные ограничения / грабли
- Умолчание
MAPPEDозначает, что сегмент занимает свою ёмкость на бумаге с первой записи. Свежая партиция — файл на 512 МиБ, в котором одна запись. Он разрежённый, место не занято, ноlsи сигнализация по месту читают видимый размер. И-Xmxпри этом не ограничивает маппинг вовсе, то есть флаги кучи ничего не говорят про rss процесса. - Записи нулевой длины запрещены, и это решение хранилища, а не валидация входа. Восстановление идёт по префиксам длин, а в предразмеченном mapped-файле за концом лога лежат нули — поэтому нулевая длина и есть маркер конца. Легальная пустая запись была бы от него неотличима.
- В mapped-пути длина пишется ПОСЛЕ тела и суммы. Длина — то, что восстановление читает первым, поэтому она обязана появляться последней: если длина видна, всё остальное уже записано. Сама по себе эта очередность не доказательство (переупорядочение страниц никто не запрещал) — она решает лишь, где остановится восстановление; отличить целую запись от рваной умеет контрольная сумма.
- Заголовок записи —
[int32 длина][int32 crc32c]. Сумма считается при записи и проверяется при восстановлении и при обычном чтении. На zero-copy-пути она не проверяется и не может: он до байтов не дотрагивается. Там проверяет клиент. - Восстановление останавливается на первой несошедшейся записи и сохраняет всё до неё. За дырой оффсеты уже не значат того, что значили, поэтому «пропустить и продолжить» было бы хуже.
- Индекс растёт по мере надобности, а не выделяется сразу. Раньше на каждый сегмент
выделялся мегабайтный
LongArrayнезависимо от его размера, и сотня мелких сегментов убивала брокер на 64 МиБ. Считать размер от ёмкости сегмента — соблазнительно и неверно: записи бывают больше интервала индекса, поэтомуёмкость / интервалне верхняя граница, и индекс начинает заполняться раньше сегмента. Position—Int, и это потолок сегмента, а не небрежность. И маппинг, и один вызовtransferToиндексируютсяint(ресёрч §1.4, §1.5).Longздесь обещал бы диапазон, который хранилище всё равно не обслужит.- Запись
intв маппинг идётJAVA_INT_UNALIGNED. СJAVA_INTFFM падает на невыровненном смещении — а записи переменной длины почти всегда невыровнены (ресёрч §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 по размеру таких оговорок не имеет.