Skip to content

Семантика Topic Reader/Writer

Reader

  • Потоковая модель: сервер шлёт пачки сообщений по активным partition‑сессиям.
  • read({ limit, batchWindowMs, signal }) — асинхронный итератор батчей;
    • limit — максимум сообщений за вызов итерации (по умолчанию нет лимита).
    • batchWindowMs — максимальное время накопления батча до отдачи; на простаивающем топике возвращается пустой батч.
    • signal — отмена ожидания/чтения.
  • Коммиты: commit(batch|msg) — подтверждает обработку на стороне сервера.
  • Hooks:
    • onPartitionSessionStart(session, committedOffset, { start, end }) — можно сдвинуть readOffset/commitOffset.
    • onPartitionSessionStop(session, committedOffset) — финализация/коммиты.
    • onCommittedOffset(session, committedOffset) — уведомления об ack.

Буферизация:

  • Внутренний буфер ограничивается maxBufferBytes.
  • codecMap управляет распаковкой полезной нагрузки.

Устойчивость:

  • Прозрачный реконнект на уровне stream (экспоненциальная задержка + jitter); ожидающие коммиты повторно отправляются на новой partition‑сессии, поэтому цикл read() + commit() переживает реконнекты.
  • Работающий reader, у которого удалили топик, простаивает, пока сервер не закроет устаревший stream (~1 мин), затем прозрачно переподключается и автоматически возобновляет чтение, если топик снова существует.

Writer

  • write(payload) — кладёт сообщение в буфер; возвращает void. Итоговый seqNo приходит через flush() или onAck.
  • flush() — отправляет буфер на сервер; возвращает последний подтверждённый seqNo.
  • onAck(seqNo, status) — уведомления о подтверждении записи.
  • Ограничения: maxBufferBytes, maxInflightCount, flushIntervalMs.
  • Прозрачный реконнект на уровне stream (экспоненциальная задержка + jitter); по умолчанию неограничен (ретраит вечно, ждёт сервер/топик), ограничивается recoveryWindowMs, если задан.

Транзакции

  • TxReader: отслеживает прочитанные оффсеты и отправляет updateOffsetsInTransaction на коммите транзакции.
  • TxWriter: дожидается flush перед коммитом (через tx.onCommit).
  • using не обязателен (оба реализуют AsyncDisposable/Disposable); жизненным циклом управляют хуки транзакции.