Skip to content

Потоковые данные

Потоковые этапы (SSE / AsyncIterable)

Этап, функция stream которого возвращает AsyncIterable<T>. Оркестратор собирает все полученные чанки в массив (результат этапа). onChunk вызывается для каждого чанка в реальном времени.

js
const orchestrator = createPipeline([
  { key: 'auth', request: async () => getToken() },
  {
    key: 'liveData',
    stream: async function* ({ prev }) {
      const source = new EventSource(`/api/stream?token=${prev}`)
      yield* eventSourceToAsyncIterable(source)
    },
    onChunk: (chunk, sharedData) => {
      sharedData.partial = (sharedData.partial ?? '') + chunk
      updateUI(sharedData.partial)
    },
  },
  {
    key: 'finalize',
    // allResults.liveData.data — полный массив чанков
    request: async ({ allResults }) => allResults.liveData.data.join(''),
  },
])
  • Учитывает abort() — проверяет сигнал прерывания между каждым чанком.
  • Поддерживает continueOnError — провалившиеся потоковые этапы можно пропускать, как и любой другой шаг.
  • Генерирует стандартные события шага: step:start, step:success, step:error.

Пагинация

paginate() обходит страницы постраничного API как AsyncGenerator<T[]>, скрывая разницу между курсорными и offset/limit-based API:

ts
import { paginate, paginateAll, flattenPages } from 'rest-pipeline-js'

// На основе курсора (стратегия по умолчанию)
for await (const page of paginate({
  fetchPage: (cursor) => client.get('/items', { params: { cursor } }).then((r) => r.data),
})) {
  console.log(page.length, 'items')
}

// На основе offset/limit
for await (const page of paginate({
  strategy: 'offset',
  limit: 50,
  fetchPage: (offset, limit) =>
    client.get('/items', { params: { offset, limit } }).then((r) => r.data),
})) {
  console.log(page.length, 'items')
}

fetchPage возвращает { items, nextCursor } (курсорная стратегия — остановка, когда nextCursor равен null/undefined) либо { items, total? } (offset-стратегия — останавливается, когда страница короче limit, либо когда offset достигает total, если API его сообщает).

  • paginateAll(options) — собирает все страницы в один плоский массив; самый простой вариант, когда весь набор данных достаточно мал, чтобы поместиться в памяти целиком.
  • flattenPages(pages) — превращает поток страниц в поток отдельных элементов; полезно как источник StreamStageConfig.stream, когда onChunk должен срабатывать на каждый элемент, а не на каждую страницу (см. examples/pagination-stream.ts).
  • И paginate(), и её коллбэк fetchPage принимают опциональный signal для поддержки abort().

Если вы заранее знаете количество страниц и хотите получать их параллельно, а не последовательно, см. вместо этого examples/pagination-fanout.ts.

Этапы WebSocket

Этап, работающий поверх постоянного WebSocket-соединения вместо одиночного запроса/ответа — чаты/фиды присутствия, живые биржевые стаканы, события совместного редактирования и т. д. Сообщения, возвращённые onMessage, собираются в массив data этапа (тот же паттерн, что и чанки потоковых этапов); onChunk срабатывает на каждое сообщение в реальном времени.

js
const orchestrator = pipe()
  .step({ key: 'auth', request: async () => getToken() })
  .websocket({
    key: 'chatFeed',
    url: ({ prev }) => `wss://chat.example.com/rooms/general?token=${prev}`,
    onOpen: () => console.log('connected'),
    onMessage: (data) => JSON.parse(data),
    onChunk: (message, sharedData) => updateUI(message),
    closeOn: (message) => message.text === '__end__',
    onClose: ({ wasClean }) => console.log('closed, clean:', wasClean),
    onError: (error) => console.error(error),
    timeoutMs: 5 * 60_000,
  })
  .build()
  • url — строка или функция от { prev, allResults, sharedData, signal }, те же параметры, что получает request.
  • createWebSocket — фабрика для базовой реализации. По умолчанию globalThis.WebSocket (браузеры, Deno, Node ≥22). Для Node <22 передайте реализацию на базе пакета ws: createWebSocket: (url, protocols) => new WS(url, protocols).
  • onMessage (обязателен) — вызывается на каждое сообщение с event.data; может быть async. Значение, отличное от undefined, собирается в массив результатов этапа и передаётся в onChunk.
  • Успех/ошибка определяются событием закрытия, а не событием ошибки: большинство реализаций WebSocket генерируют error непосредственно перед close, так что сам по себе onError не проваливает этап — чистое закрытие (wasClean: true) успешно завершает этап со всем, что было собрано к тому моменту; нечистое закрытие отклоняет его, проходя через continueOnError, как и любой другой этап.
  • closeOn(data) — верните true, чтобы закрыть соединение и успешно завершить этап, как только вы увидели нужное, вместо ожидания закрытия сервером.
  • timeoutMs — общий таймаут соединения (не сбрасывается сообщениями); закрывает сокет и проваливает этап, если срабатывает до того, как произойдёт чистое закрытие само по себе.
  • Учитывает abort() — закрывает базовое соединение и отклоняет этап.

Полную аннотированную версию см. в examples/websocket-stage.ts.