booblik
DocsApi

Бинарный протокол 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
...payload

correlationId нужен, потому что сессия конвейерная: клиент вправе отправить следующий запрос, не дождавшись предыдущего ответа, и ответы приходят в порядке отправки. Без него конвейеризация превращается в «запрос-ответ», и все замеры 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  logEndOffset

ackPolicy = 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 больше maxBytesCORRUPT_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. Коды ошибок

КодИмяКогда
0NONE
1UNKNOWN_TOPIC_OR_PARTITIONнет такого топика или партиции
2OFFSET_OUT_OF_RANGEfetchOffset вне [logStartOffset, highWatermark]. Ровно highWatermarkне ошибка: так выглядит догнавший потребитель
3RECORD_TOO_LARGEзапись не помещается в сегмент
4UNSUPPORTED_VERSIONнеизвестный apiKey или apiVersion. Кадр при этом исправен
5CORRUPT_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 и xorout0xFFFFFFFF. Табличная реализация со сдвигом вправо использует зеркало полинома, 0x82F63B78; перепутать их — получить сумму устойчивую, правдоподобную и всюду неверную.

Проверять её обязан клиент, и только он: на zero-copy-пути брокер до этих байтов не дотрагивается (§4). Клиент, который её пропустит, молча выключает единственную в проекте защиту от порчи диска.

Это не CRC-32 из zlib и не System.IO.Hashing.Crc32 из .NET — другой полином, и подмена тихая: обе функции называются «CRC32», обе возвращают правдоподобное число. Контрольное значение для строки 1234567890xE3069283; реализация, не дающая его, неверна.

7.3. Вектора

Оба алгоритма закреплены golden-векторами в conformance/vectors/: канонический вид — TSV, рядом лежит зеркало в JSON. Их считает вторая реализация на Python (conformance/vectors/generate.py), а не выгрузка из Kotlin-клиента: выгрузка доказывала бы только то, что клиент согласен сам с собой.

Вектора подобраны под то, что они ловят, а не под покрытие: 0x7f и 0x80 соседние специально — на них расходится знаковое прочтение, и p2/p16 у них разные.

Тест падает — неверен код, а не вектора. Перегенерация ради зелёного прогона превращает эталон в эхо. Единственный честный повод перегенерировать — добавление векторов.

On this page