booblik
DocsFeatures

Подписаться на топик и публиковать в него

Веха M7 закрыта целиком: METADATA (M-70), долгий FETCH (M-75), подписка и контрольные точки (M-71, M-72), публикация (M-73). Код — Subscription.kt и Publishing.kt. Consumer.poll() остаётся как низкоуровневый примитив под всем этим. Не сделан замер M-74 — до него числа подписки неизвестны.

1. Суть

Читателю нужен поток записей, а не цикл вокруг poll(). Писателю нужно назвать топик один раз, а не таскать пару «топик, партиция» через весь код. Обе задачи упираются в одно и то же ограничение, и оно определяет форму сильнее, чем вкус.

2. Три ограничения, из которых всё следует

Брокер не хранит позиции потребителей. Это то самое решение, которое убрало координатор групп (ресёрч, раздел 2). Подписочный API соблазняет притвориться, что система помнит, где вы остановились; она не помнит. Отсюда правило: позиция видна в каждой выдаче, а сохраняет её вызывающий. И отсюда же запрет на слово commit в именах — оно означало бы, что кто-то на той стороне принял позицию к сведению.

Клиент не знает, сколько у топика партиций. Набор фиксируется при старте брокера, запроса метаданных в протоколе нет. Поэтому subscribe(topic) без перечисления партиций — не вопрос удобства обёртки, а новый apiKey. Он заводится (M-70).

Горячий путь не аллоцирует и не берёт локи. Второе проверяется байткод-гейтом (M-62). Значит DSL не строит промежуточные структуры на публикацию, а слияние потоков делается корутинами и каналом, а не разделяемым состоянием под замком.

3. METADATA (сделано, M-70)

ApiKey.METADATA(3). Тело запроса — список имён топиков; пустой список означает «все».

Ответ на партицию несёт три числа, и третье здесь не для украшения:

ПолеЗачем
partitionIdсобственно то, ради чего запрос
logStartOffsetначало живого лога после retention; без него «читать сначала» пришлось бы понимать как Offset.ZERO, а это OFFSET_OUT_OF_RANGE на любом топике, который что-то удалял
highWatermark«читать только новое» и вычисление отставания без пробного FETCH

Точный формат кадра — protocol-wire, раздел METADATA (пишется в M-70).

4. Подписка (сделано, M-71)

Единица выдачи — батч, а не запись. Единица чтения — то, ради чего сделан весь проект (замер 2.2: батчи против записи по одной дали 54 раза), и API не имеет права её прятать. Плоский поток записей остаётся отдельной операцией для тех, кому он удобнее.

data class RecordBatch(
    val topic: TopicName,
    val partition: PartitionId,
    val baseOffset: Offset,
    val records: List<ByteArray>,
    val highWatermark: Offset,
) {
    /** Единственное число, которое нужно перезапустившемуся потребителю. */
    val nextOffset: Offset get() = baseOffset + records.size
    val lag: Long get() = highWatermark.value - nextOffset.value
}

Два режима, и путать их дорого:

ОперацияКогда догнали хвостДля чего
follow()ждёт и продолжаетработа
replay()поток кончаетсяпересчёты, миграции, отладка

replay() фиксирует highWatermark на момент старта. Иначе на живом топике он не кончится никогда, а «прочитать всё, что было» — ровно тот запрос, ради которого его зовут.

Старт задаётся не числом, а намерением:

sealed interface StartPosition {
    data object Earliest : StartPosition       // logStartOffset, а не ноль
    data object Latest : StartPosition         // highWatermark
    data class At(val offset: Offset) : StartPosition
}

Порядок гарантирован внутри партиции и только. Подписка на несколько партиций читает каждую своей корутиной и сливает результаты; межпартиционного порядка не существует, потому что его не существует и в логе.

Обратное давление даётся даром: коллектор приостанавливается — цикл выборки перестаёт просить. Это единственная причина, по которой здесь Flow, а не колбэк.

Цена измерена (замер 17): машинерия Flow — 5,5 нс на выдачу, канал callbackFlow — ещё 54 нс, при сквозном пути 712 нс на запись. Платится за выдачу, а не за запись, поэтому при батче 100 это 0,008 % и 0,08 % пути. Это же — арифметическая причина отдавать батч, а не запись: при выдаче по одной канал стоил бы уже 8 %.

4а. Долгий FETCH (сделано, M-75)

Опрос с интервалом был бы дешёвым в реализации и дорогим навсегда: сто подписчиков на догнанном хвосте — это сто холостых запросов каждые 50 мс, а поменять это потом означает сломать и follow(), и обработчик FETCH, и формат кадра. Поэтому форма фиксируется здесь, а не после M-71.

Сигнал — StateFlow, а не SharedFlow. Разница не стилистическая. У SharedFlow между «посмотрели текущий highWatermark» и «начали слушать» есть окно, в которое обновление проваливается, и потерянное пробуждение выглядит как зависший раз в сутки потребитель — то есть как самый дорогой класс багов в этом проекте. StateFlow хранит текущее значение, поэтому

withTimeoutOrNull(maxWait) { handle.writer.highWatermark.first { it > request.fetchOffset } }

не может пропустить то, что уже произошло. Конфляция здесь безвредна: watermark монотонен, и интересно только последнее значение.

Это запись на горячем пути, и её надо померить. Публикация watermark происходит на каждом закоммиченном батче — то есть ровно там, где проект запрещает лишнюю работу. Смягчение — публиковать только когда есть ожидающие; но и атомарное чтение счётчика не бесплатно, поэтому M-75 не считается закрытой без числа (правило вехи).

Ожидание блокирует конвейер на этом соединении, и это самое опасное место. Session обслуживает запросы строго по одному — так устроена гарантия порядка, на которой держится сопоставление по correlationId. Значит удержанный FETCH задерживает всё, что послано следом по тому же соединению, включая PRODUCE. Сегодня этого не видно, потому что FETCH отвечает сразу.

Решение на M-75 — контракт, а не переделка сессии: follow() владеет своим соединением. Разбирать сессию на конкурентную обработку с очередью записи означало бы менять и клиентский контракт тоже; если это когда-нибудь понадобится, это отдельная веха, а не деталь этой.

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

Тихая смерть TCP — реальная плата за удержание. Соединение без трафика минуту роняют NAT и файрволы. Отсюда умолчание maxWaitMillis = 30 с (ниже обычных таймаутов простоя) и SO_KEEPALIVE на принятых сокетах.

Метрика удерживаемых запросов обязательна. Иначе здоровый брокер со ста подписчиками выглядит как мёртвый: соединений сто, produce 0/s, fetch 0/s. После M-64 мы знаем цену состояния, которое неотличимо от поломки.

4б. Контракт долгого FETCH в сценариях (частично закрыт тестами)

Сценарии — спецификация, а не иллюстрация: каждый получает тест, и пока теста нет, строка помечена целевой. Исполняемого Gherkin (Cucumber) в проекте нет и заводить его отдельным решением — здесь сценарии живут в документе, а рядом стоит имя теста.

Сценарий: догнавший потребитель ждёт, а не опрашивает
  Дано потребитель на highWatermark и maxWaitMillis = 5000
  Когда он посылает FETCH
  Тогда брокер не отвечает, пока не появится запись или не истечёт ожидание
  И запрос считается в метрике удерживаемых

Сценарий: запись во время ожидания будит немедленно
  Дано удерживаемый FETCH с maxWaitMillis = 5000
  Когда продюсер коммитит батч
  Тогда ответ приходит сразу, а не по истечении ожидания
  И содержит записи этого батча

Сценарий: запись, пришедшая между проверкой и ожиданием, не теряется
  Дано потребитель на highWatermark
  Когда батч коммитится ровно в момент между чтением watermark и началом ожидания
  Тогда FETCH возвращается с данными, а не виснет до таймаута

Сценарий: истёкшее ожидание — пустой ответ
  Дано удерживаемый FETCH с maxWaitMillis = 100 и пустой хвост
  Когда ожидание истекает
  Тогда приходит ответ с нулём записей и текущим highWatermark
  И это не ошибка, и соединение остаётся открытым

Сценарий: minBytes больше maxBytes отвергается кадром
  Когда клиент посылает FETCH с minBytes 2048 и maxBytes 1024
  Тогда ответом будет CORRUPT_REQUEST
  И соединение остаётся открытым

Сценарий: retention во время ожидания
  Дано удерживаемый FETCH с оффсета, который ещё жив
  Когда retention удаляет сегмент, содержащий этот оффсет
  Тогда ответом будет OFFSET_OUT_OF_RANGE, а не пустой ответ и не зависание

Сценарий: остановка брокера во время удержания
  Дано удерживаемый FETCH
  Когда брокер останавливается
  Тогда он пытается ответить пустым ответом до закрытия сокета
  И клиент, увидевший обрыв вместо ответа, повторяет запрос с того же оффсета
  И ни одна запись не потеряна и не продублирована

Сценарий: обрыв сети во время удержания
  Дано удерживаемый FETCH и соединение, оборванное посередине
  Когда клиент переподключается и повторяет FETCH с того же оффсета
  Тогда он получает ровно те записи, которых ещё не видел

Сценарий: удержание не задерживает чужие запросы
  Дано follow() на собственном соединении и продюсер на другом
  Когда FETCH удерживается пять секунд
  Тогда PRODUCE отвечает без задержки
  # Обратное — PRODUCE в конвейере за удержанным FETCH на ОДНОМ соединении — задержится,
  # и это свойство сессии, а не дефект: см. раздел 4а.

5. Контрольные точки (сделано вместе с M-71)

interface OffsetStore {
    suspend fun load(topic: TopicName, partition: PartitionId): Offset?
    suspend fun save(topic: TopicName, partition: PartitionId, offset: Offset)
}

Объявлен интерфейс — и всё. Реализации библиотека не поставляет намеренно: файл, БД и транзакционная семантика вокруг сохранения — это решение о чужой системе, ровно как экспортёр метрик, которого здесь тоже нет (см. Metrics).

checkpointing(store) сохраняет позицию после того, как коллектор обработал батч. Это at-least-once: падение между обработкой и сохранением повторит батч. Сохранять раньше значило бы at-most-once, то есть терять; выбор здесь не про вкус, а про то, какая ошибка допустима, и он называется в имени и в KDoc.

6. Публикация (сделано, M-73)

Два разных инструмента, и разница между ними — гарантия, а не сахар.

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

val orders = producer.topic(TopicName("orders"))     // партиции придут из METADATA
orders.send(record, key = userId)

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

Ключ не уходит на провод: в формате записи его нет, брокер про ключи не знает. Отсюда следует и то, чего не будет никогда: компакции по ключу.

Батч одним запросом даёт то, чего в API сегодня нет:

val offsets = producer.batch(orders, partition = 0) {
    add(a)
    add(b)
}

Записи одного запроса получают смежные оффсеты. Не атомарные: при обрыве восстановление останавливается на первой несошедшейся записи и может сохранить префикс батча. Написано прямо, потому что batch { } иначе прочитают как транзакцию, а транзакций здесь нет.

Аккумулятор при этом остаётся один. Второй, поверх существующего, дал бы два места, где запись ждёт отправки, и два объяснения любой задержке.

7. BDD (целевое — код ещё не написан)

Сценарий: подписка на топик без перечисления партиций
  Дано топик orders с тремя партициями
  Когда клиент подписывается на orders с позиции Earliest
  Тогда он получает записи всех трёх партиций
  И порядок сохраняется внутри каждой партиции

Сценарий: Earliest — это начало живого лога, а не ноль
  Дано топик, у которого retention удалил первые сегменты
  Когда клиент подписывается с позиции Earliest
  Тогда чтение начинается с logStartOffset
  И ошибки OFFSET_OUT_OF_RANGE не происходит

Сценарий: replay кончается, follow — нет
  Дано топик, в который продолжают писать
  Когда клиент вызывает replay()
  Тогда поток завершается на highWatermark, снятом при старте
  Когда клиент вызывает follow()
  Тогда поток не завершается и отдаёт новые записи по мере появления

Сценарий: контрольная точка сохраняется после обработки
  Дано подписка с OffsetStore
  Когда обработка батча бросает исключение
  Тогда позиция в хранилище остаётся прежней
  И следующий запуск читает этот батч заново

8. Чего здесь нет и не будет

  • Групп потребителей и ребалансировки. Партиции называет вызывающий или он же берёт их все из METADATA. Распределить их между несколькими процессами — задача того, кто эти процессы запускает; брокер в этом не участвует.
  • Хранения позиций на брокере. См. раздел 2.
  • Ключей на проводе и компакции. См. раздел 6.
  • Транзакций. batch { } даёт смежность, а не атомарность.

On this page