Опции и методы Topic Reader/Writer
Reader
topic:string | TopicReaderSource | TopicReaderSource[]TopicReaderSource:{ path: string; partitionIds?: bigint[]; maxLag?: number | string | Duration; readFrom?: Date | Timestamp }
consumer:stringcodecMap?:Map<Codec | number, CompressionCodec>— дополнительные кодеки распаковки (встроенный ZSTD требует Node.js 22.15+ / 23.8+)maxBufferBytes?:bigint— лимит внутреннего буфера (по умолчанию 8 МиБ)updateTokenIntervalMs?:number— период обновления токена (по умолчанию 60000)gracefulShutdownTimeoutMs?:number— дедлайн принудительного закрытия для gracefulclose(), после него ожидающие коммиты отбрасываются (по умолчанию 30000)recoveryWindowMs?:number— окно терминального реконнекта; по умолчанию неограниченно (реконнект вечно), передайте конечное значение в мс, чтобы ограничитьretryOnSchemeError?:boolean— ретраить SCHEME_ERROR (например, топик ещё не создан); по умолчанию выключено, включите, чтобы дождаться создания топикаonPartitionSessionStart?:(session, committedOffset, { start, end }) => Promise<void | { readOffset?, commitOffset? }>onPartitionSessionStop?:(session, committedOffset) => Promise<void>onCommittedOffset?:(session, committedOffset) => void
Методы и поведение:
read({ limit?, batchWindowMs?, signal? }):AsyncIterable<TopicMessage[]>- Возвращает последовательность батчей сообщений.
limitограничивает общее число сообщений, извлекаемых за один «пробег» итератора, чтобы контролировать задержку и память.batchWindowMsзадаёт максимальное время накопления батча до отдачи; на простаивающем топике итератор вернёт пустой батч[], что позволяет неблокирующую интеграцию в event loop.signalпозволяет прервать ожидание/чтение. - Почему так: длительные блокировки чтения мешают кооперативной многозадачности; «пустые» отдачи по таймауту упрощают планирование работы без busy‑wait.
- Возвращает последовательность батчей сообщений.
commit(messages | message):Promise<void>- Подтверждает обработку до соответствующего оффсета в каждой затронутой партиции (идемпотентно). Коммит гарантирует, что последующее чтение начнётся после подтверждённого оффсета. Можно вызывать на массиве сообщений (одного батча) или одном сообщении.
- Зачем: это реализация как минимум один раз (at‑least‑once). Коммит отделяет «прочитано» от «обработано» и позволяет безопасное восстановление.
- Перфоманс: ожидание
await commit()в горячем пути снижает пропускную способность. Допустима стратегия fire‑and‑forget (void reader.commit(batch)) c наблюдением черезonCommittedOffset.
close():Promise<void>- Завершает чтение «мягко»: перестаёт принимать новые данные, дожидается завершения ожидающих коммитов (с защитным таймаутом) и корректно останавливает фоновые задачи.
destroy(reason?):void- Немедленно останавливает все операции, отклоняет ожидающие коммиты, освобождает ресурсы.
Writer
topic:stringtx?:TX— запись внутри транзакцииproducer?:stringcodec?:CompressionCodecmaxBufferBytes?:bigint— по умолчанию 256 МБmaxInflightCount?:number— по умолчанию 1000flushIntervalMs?:number— по умолчанию 1000 мсupdateTokenIntervalMs?:number— по умолчанию 60000gracefulShutdownTimeoutMs?:number— по умолчанию 30000recoveryWindowMs?:number— окно терминального реконнекта; по умолчанию неограниченно (реконнект вечно), передайте конечное значение в мс, чтобы ограничитьretryOnSchemeError?:boolean— ретраить SCHEME_ERROR (например, топик ещё не создан); по умолчанию выключено, включите, чтобы дождаться создания топикаpartitionId?/messageGroupId?— привязка/маршрутизация записи (взаимоисключающие)onAck?(seqNo, status):(seqNo: bigint, status: 'skipped' | 'written' | 'writtenInTx') => void
Методы и поведение:
write(payload: Uint8Array, extra?):void- Кладёт сообщение в буфер. Опционально можно задать
seqNo(ручной режим),createdAt,metadataItems. Запись не блокирует; фактическая отправка выполняется приflush()или периодическим флашером. ИтоговыйseqNoможно получить черезflush()илиonAck. - Почему
seqNo: на продюсереproducer + seqNoобеспечивает идемпотентность и детерминизм подтверждений (и упорядоченность в партиции).
- Кладёт сообщение в буфер. Опционально можно задать
flush():Promise<bigint>- Выгружает накопленные сообщения в сеть, дожидается подтверждений «в полёте» и возвращает последний подтверждённый
seqNo. Используйте в контрольных точках (например, при остановке сервиса).
- Выгружает накопленные сообщения в сеть, дожидается подтверждений «в полёте» и возвращает последний подтверждённый
close():Promise<void>- «Мягко» завершает работу: прекращает приём новых сообщений, дожидается флаша, освобождает ресурсы. Отклоняется, если корректно доставить буфер не удалось.
destroy():void- Немедленное прекращение без гарантии доставки.
Подтверждения:
onAck(seqNo, status): уведомляет о судьбе сообщения.status:written— записано вне транзакции;writtenInTx— записано в транзакции (станет видимым после коммита);skipped— пропущено (например, из‑за конфликтаseqNo).
Повторные попытки и устойчивость:
- Подключение к TopicService — потоковое; при обрывах прозрачно переподнимается (экспоненциальная задержка + jitter). По умолчанию реконнект неограничен (ждёт сервер/топик); передайте
recoveryWindowMs, чтобы задать таймаут, иretryOnSchemeError, чтобы дождаться ещё не созданного топика. Сообщения «в полёте» отправляются повторно; прозрачный реконнект не приводит к ошибкам ожидающих записей.
Транзакционные варианты:
createTopicTxReader(tx, ...)иcreateTopicTxWriter(tx, ...)привязаны к транзакции Query.- TxReader отслеживает прочитанные оффсеты и отправляет
updateOffsetsInTransactionнаtx.onCommit. - TxWriter инициирует
flushнаtx.onCommitи корректно сворачивается наtx.onRollback/onClose. - Оба реализуют
AsyncDisposable/Disposable, ноusingне обязателен: жизненным циклом управляет транзакция.
- TxReader отслеживает прочитанные оффсеты и отправляет