Семантика 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); жизненным циклом управляют хуки транзакции.