Бинарный протокол booblik
Реализовано (M3). Кодек —
booblik-net/src/main/kotlin/.../net/wire/, разбор вRequestDecoder, сборка вResponseEncoder, протокольный цикл вSession. Расхождения между этим документом и кодом означают, что неправ документ.Топики и партиции задаются при старте
Broker.open— создавать их протокол не умеет, и это решение, а не пробел.
1. Зачем свой протокол
Совместимости с Kafka не будет. Она означала бы поддержку десятков версий API, групп потребителей и координатора — то есть ровно ту часть Kafka, от которой этот проект отказывается. Свой протокол занимает страницу, и это его главное свойство.
2. Кадр
Всё big-endian — как и запись на диске (booblik-core). Совпадение не косметическое: оно
позволяет отдавать байты сегмента в сокет transferTo, ничего не переворачивая.
Запрос
int32 frameLength // длина всего, что после этого поля
int16 apiKey // 1 = PRODUCE, 2 = FETCH
int16 apiVersion // 1
int32 correlationId // возвращается в ответе как есть
...payloadОтвет
int32 frameLength
int32 correlationId
int16 errorCode // 0 = OK
...payloadcorrelationId нужен, потому что сессия конвейерная: клиент вправе отправить следующий запрос,
не дождавшись предыдущего ответа, и ответы приходят в порядке отправки. Без него конвейеризация
превращается в «запрос-ответ», и все замеры RPS упираются в RTT, а не в брокер.
3. PRODUCE (apiKey = 1)
u16 topicNameLength
bytes topicName // UTF-8, ограничение длины — TopicName.MAX_LENGTH
int32 partitionId
int8 ackPolicy // 0 = не ждать, 1 = ждать записи, 2 = ждать force()
int32 recordCount
records: // recordCount раз
int32 payloadSize
bytes payloadКонтрольной суммы в PRODUCE нет: её считает брокер при записи. Клиентская сумма проверяла бы арифметику клиента поверх того, что уже проверяет TCP, а защитить нужно диск.
Два ограничения, каждое отвечается CORRUPT_REQUEST при разборе (RequestDecoder).
Записаны здесь в M-131: комплект соответствия наткнулся на них первым же прогоном, а документ
о них молчал — то есть автор клиента на другом языке узнал бы о них в проде.
recordCountдолжен быть положительным. Батч без записей — это запрос, который ничего не просит и на который нечего ответить: оффсета не возникает, аbaseOffsetв ответе обязан на что-то указывать;payloadSizeкаждой записи должен быть положительным — пустых записей не бывает. Причина в восстановлении, а не в экономии: оно читает префикс длины первым и останавливается на первой записи, чьи байты не сходятся с её заголовком. Нулевая длина неотличима от незаписанного места в свежем сегменте, поэтому сохранённая пустая запись обрывала бы лог на себе при следующем старте. Клиент обязан донести отказ до вызывающего: проглоченный отказ теряет записи молча, что хуже самого ограничения.
Записи идут батчем, и это не удобство, а суть. Подтверждение на каждое сообщение — это
CompletableDeferred и круг по каналу на запись, которая сама стоит десятки наносекунд
(ресёрч, риск 3). Единица протокола — батч, единица адресации — запись.
Формат тела батча совпадает с форматом на диске, запись в запись. Поэтому PRODUCE дописывается в сегмент без перекладывания, а FETCH отдаётся из сегмента без сборки.
Ответ
int64 baseOffset // оффсет первой записи батча; остальные идут подряд
int64 logEndOffsetackPolicy = 0 не отвечает вообще — ответа нет, а не «ответ с обещанием». Так вышло не из
экономии: оффсета не существует, пока писатель не дошёл до батча, поэтому вернуть в мгновенном
ответе было бы нечего, кроме числа, которое клиенту не с чем сверить. Реализация устроена так же
(PartitionWriter.append с AckPolicy.NONE возвращает null). Это единственный режим, в котором
брокер вправе потерять принятое, и цена честности здесь в том, что клиент не узнает и о перегрузке.
ackPolicy = 2 не означает «один барьер на запрос». Писатель собирает всех, кто уже стоит в
очереди, и делает один force на всю группу — иначе потолок был бы 250 запросов в секунду
независимо от числа производителей (ресёрч §1.9, замер группового коммита в
benchmarking).
4. FETCH (apiKey = 2)
Версия 1 — то, что есть; версия 2 — целевая (M-75) и добавляет два поля в конец запроса. Оба поля пишутся здесь заранее, до реализации, по одной причине: долгий FETCH меняет не удобство, а форму взаимодействия, и клиентский API нельзя замораживать поверх опроса, чтобы потом ломать и его, и обработчик, и формат.
u16 topicNameLength
bytes topicName
int32 partitionId
int64 fetchOffset
int32 maxBytes
int32 maxWaitMillis // v2: сколько держать запрос, если данных нет; 0 = как в v1
int32 minBytes // v2: не отвечать, пока не накопится столько; 0 или 1 = отдать что естьОтвет формата не меняет. Истёкшее ожидание — это пустой ответ, а не ошибка: догнавший потребитель — нормальное состояние, и оно уже описано так же в v1.
Три правила, каждое закрывает конкретный способ повесить брокер или клиента:
minBytesбольшеmaxBytes—CORRUPT_REQUESTпри разборе. Такой запрос нельзя удовлетворить никогда: он просит не отвечать, пока не наберётся больше, чем он согласен принять. Проверяется в кодеке, а не в сессии, потому что это свойство кадра;maxWaitMillisограничен сверху брокером (целевое — 60 с). Клиент, попросивший час, получает более ранний пустой ответ; это безобидно, а незажатое значение держит корутину и сокет столько, сколько попросит кто угодно. Ограничение названо в документе, а не спрятано;- ожидание происходит до
openFetch, а не после. Сегмент удерживается от подсчёта длины до последнего отправленного байта; удерживать его ещё и всё время ожидания значило бы блокировать retention на минуту. Следствие: за время ожиданияlogStartOffsetможет уйти вперёд, и тогда ответом будетOFFSET_OUT_OF_RANGE— законный исход, а не гонка.
Ответ
int64 highWatermark // первый оффсет, которого ещё нет
int32 payloadBytes
bytes payload // сырые записи сегмента, побайтово как на диске:
// int32 payloadSize
// int32 crc32c
// bytes payloadКонтрольную сумму в ответе проверяет клиент, а не брокер. Брокер не может: на zero-copy-пути он до этих байтов не дотрагивается — в этом весь его смысл. Сумма считается один раз при записи и защищает диск, а не сеть; за сеть отвечает TCP.
payload уезжает в сокет через LogSegment.transferTo, минуя heap. Отсюда три следствия,
и все три пришлось учесть в Session:
- заголовок ответа пишется обычной записью, и только после него начинается
transferTo. Порядок нарушать нельзя — сокет один; payloadBytesсчитается до передачи и не может зависеть от того, сколькоtransferToотдал за раз: он отдаёт частями и это норма (ресёрч §1.4);- сегмент удерживается от момента, когда посчитали длину, до момента, когда отправили последний
байт (
PartitionLog.openFetch). Иначе retention уведёт его между обещанием и отправкой, и ответ окажется короче собственного префикса длины.
Тело можно отдавать и через heap (FetchMode.HEAP) — это не запасной путь, а контроль
в эксперименте M-35. Ниже гигабайта в секунду разницы между ними нет вовсе
(benchmarking, замер 7).
Ответ может оборваться на середине записи: maxBytes — граница по байтам, а не по записям.
Клиент обязан отбросить незавершённый хвост и запросить его следующим FETCH — ровно так же
устроен Kafka, и по той же причине: иначе брокеру пришлось бы разбирать батч, чтобы найти
границу, а это тот самый разбор, которого zero-copy избегает.
Из этого следует случай, который обязан быть ошибкой, и записан он в M-136. Если целых записей
в ответе нет ни одной, а обрезанный хвост есть, — следующая запись крупнее maxBytes и целиком
не приедет никогда. Клиент, который отбрасывает хвост и повторяет, делает вечно тот же самый
запрос: он продолжает работать, не сообщает ничего и не двигается. Единственное лекарство —
поднять maxBytes, поэтому клиент обязан сказать об этом вызывающему, а не молчать. Отличать
от fetchOffset == highWatermark: там пустой ответ — норма, и она разрешается сама собой,
как только кто-нибудь допишет.
4а. METADATA (apiKey = 3)
Заведён в M-70 ради одной вещи: набор партиций фиксируется при старте брокера и до сих пор нигде на проводе не появлялся, поэтому подписаться на топик можно было, только перечислив партиции руками. Это не слой метаданных — создания топиков нет и кластера, который надо описывать, тоже, — это единственный вопрос, который читателю приходится задать до чтения.
int32 topicCount // 0 = все топики этого брокера
per topic:
u16 topicNameLength
bytes topicNameОтвет
int32 topicCount
per topic:
u16 topicNameLength
bytes topicName
int32 partitionCount
per partition:
int32 partitionId
int64 logStartOffset // начало ЖИВОГО лога, после retention
int64 highWatermark // первый оффсет, которого ещё нетТри числа на партицию, и первые два — весь смысл запроса. logStartOffset — начало живого
лога: без него «читать сначала» пришлось бы понимать как ноль, а это OFFSET_OUT_OF_RANGE на
любом топике, который хоть раз удалял сегмент. highWatermark отвечает на «читать только новое»
и даёт отставание без пробного FETCH.
Названный топик, которого нет, роняет весь запрос с UNKNOWN_TOPIC_OR_PARTITION, а не молча
выпадает из ответа. Иначе «топика нет» и «топик есть и пуст» приходили бы одинаково, и подписчик,
для которого эта разница существенна, читал бы пустоту вечно, ничего об этом не узнав. Запрос без
имён просит всё и на это наткнуться не может.
Порядок ответа детерминирован — топики по имени, партиции по номеру. Порядок обхода хеш-таблицы сделал бы два одинаковых запроса разными на вид, и обнаружил бы это первый, кто сравнит два ответа.
5. Коды ошибок
| Код | Имя | Когда |
|---|---|---|
| 0 | NONE | — |
| 1 | UNKNOWN_TOPIC_OR_PARTITION | нет такого топика или партиции |
| 2 | OFFSET_OUT_OF_RANGE | fetchOffset вне [logStartOffset, highWatermark]. Ровно highWatermark — не ошибка: так выглядит догнавший потребитель |
| 3 | RECORD_TOO_LARGE | запись не помещается в сегмент |
| 4 | UNSUPPORTED_VERSION | неизвестный apiKey или apiVersion. Кадр при этом исправен |
| 5 | CORRUPT_REQUEST | кадр не разбирается |
Ответ на ошибку возвращает correlationId запроса во всех случаях, кроме одного: если кадр
короче заголовка, идентификатора просто нет, и тогда возвращается ноль. Это не педантизм —
клиент сопоставляет ответы по идентификатору, поэтому ответ с чужим номером не теряет информацию,
а разрешает чужой запрос.
Соединение при ошибке не рвётся: кадрирование было исправным, значит собеседник говорит на протоколе. Рвётся оно только на длине кадра вне допустимого диапазона — там уже непонятно, где начинается следующее сообщение.
6. Чего в протоколе нет и не будет в MVP
- TLS. Не «потом добавим»: шифрование и zero-copy физически несовместимы — шифровать можно только то, чего коснулся процессор (ресёрч §1.2). Это отдельный путь чтения, а не флаг.
- Аутентификация. Следствие предыдущего пункта: без канала, которому можно доверять, она даёт ложное чувство защищённости.
- Метаданные, группы потребителей, коммит оффсетов. Позицию помнит клиент. Это половина причины, по которой брокер помещается в голову.
- Сжатие. По той же причине, что и TLS.
7. Два алгоритма, которые обязаны совпадать у всех
Здесь описано то, чего на проводе не видно, но что обязано совпадать у любых двух клиентов — иначе они разойдутся молча. Раздел заведён в M13-1, когда стало ясно, что клиенты появятся не только на JVM.
7.1. Партиционирование по ключу
Ключа в протоколе нет и не будет: в формате записи для него нет места (§3), поэтому брокер не маршрутизирует по ключу, не уплотняет по ключу и не сможет. Партицию выбирает клиент, и на провод уезжает уже её номер.
Отсюда требование, которое до первого неJVM-клиента выглядело выполненным само собой:
Два продюсера обязаны выбрать одну и ту же партицию для одного ключа.
Не выберут — записи под одним ключом лягут в две партиции, и порядок по ключу, ради которого партиции и существуют, пропадёт. Отказ тихий: запись уходит в существующую партицию, просто не в ту, и ни ошибки, ни лога не будет.
До 0.3.0 умолчанием был java.util.Arrays.hashCode. Он задан JDK, чего достаточно, чтобы
сошлись две JVM, и недостаточно ни для чего другого: это алгоритм платформы, а не контракт
на проводе. С 0.3.0 умолчание — FNV-1a, выписанный в байтах:
hash = 0x811C9DC5 // FNV offset basis, 32 бита
для каждого байта b ключа: // b беззнаковый, 0..255
hash = (hash XOR b) * 0x01000193 // FNV prime, умножение по модулю 2³²
partition = hash mod partitionCount // остаток БЕЗЗНАКОВЫЙОба «беззнаковый» здесь несущие, а не уточняющие. Языки расходятся ровно в этих двух местах:
- байт. В Java и Kotlin элемент массива знаковый, поэтому
0x80без маски войдёт в сумму как0xFFFFFF80. В Go, Rust и C байт беззнаковый и маска не нужна. Реализация, ошибшаяся здесь, проходит все латинские ключи и расходится на первом ключе с высоким байтом — то есть на первом же имени не на латинице; - остаток. Знаковый остаток от отрицательного хеша даёт отрицательный номер партиции.
Беззнаковый остаток убирает вопрос целиком, поэтому выбран он, а не
floorMod.
Взят FNV-1a, а не murmur2 из Kafka: четыре строки против двадцати, без беззнакового сдвига вправо и без разбора хвоста. Качество распределения решает здесь меньше, чем кажется — свёртка идёт по числу партиций в единицах.
Ключа нет — партиция берётся круговым перебором, и это намеренно не специфицировано. Счётчик живёт в клиенте и в его памяти; запись без ключа не несёт обещания о порядке, поэтому совпадение между клиентами тут не требуется и требовать его значило бы запретить любую другую раскладку.
7.2. Контрольная сумма
CRC-32C (Castagnoli): полином 0x1EDC6F41, вход и выход отражённые, init и xorout —
0xFFFFFFFF. Табличная реализация со сдвигом вправо использует зеркало полинома, 0x82F63B78;
перепутать их — получить сумму устойчивую, правдоподобную и всюду неверную.
Проверять её обязан клиент, и только он: на zero-copy-пути брокер до этих байтов не дотрагивается (§4). Клиент, который её пропустит, молча выключает единственную в проекте защиту от порчи диска.
Это не CRC-32 из zlib и не System.IO.Hashing.Crc32 из .NET — другой полином, и подмена
тихая: обе функции называются «CRC32», обе возвращают правдоподобное число. Контрольное значение
для строки 123456789 — 0xE3069283; реализация, не дающая его, неверна.
7.3. Вектора
Оба алгоритма закреплены golden-векторами в conformance/vectors/: канонический вид — TSV,
рядом лежит зеркало в JSON. Их считает вторая реализация на Python
(conformance/vectors/generate.py), а не выгрузка из Kotlin-клиента: выгрузка доказывала бы
только то, что клиент согласен сам с собой.
Вектора подобраны под то, что они ловят, а не под покрытие: 0x7f и 0x80 соседние
специально — на них расходится знаковое прочтение, и p2/p16 у них разные.
Тест падает — неверен код, а не вектора. Перегенерация ради зелёного прогона превращает эталон в эхо. Единственный честный повод перегенерировать — добавление векторов.