Pipelines
History
pipeTo(source, ...transforms?, writer, options?): Promise
AsyncIterable | IterableObjectwrite(chunk) method.ObjectAbortSignalpreventFail is true.booleantrue, do not call writer.end() when
the source ends. Default: false.booleantrue, do not call writer.fail() on
error. Default: false.PromisePipe a source through transforms into a writer. If the writer has a
writev(chunks) method, entire batches are passed in a single call (enabling
scatter/gather I/O).
If the writer implements the optional *Sync methods (writeSync, writevSync,
endSync), pipeTo() will attempt to use the synchronous methods
first as a fast path, and fall back to the async versions only when the sync
methods indicate they cannot complete (e.g., backpressure or waiting for the
next tick). fail() is always called synchronously.
import { from, pipeTo } from 'node:stream/iter'; import { compressGzip } from 'node:zlib/iter'; import { open } from 'node:fs/promises'; const fh = await open('output.gz', 'w'); const totalBytes = await pipeTo( from('Hello, world!'), compressGzip(), fh.writer({ autoClose: true }), );
const { from, pipeTo } = require('node:stream/iter'); const { compressGzip } = require('node:zlib/iter'); const { open } = require('node:fs/promises'); async function run() { const fh = await open('output.gz', 'w'); const totalBytes = await pipeTo( from('Hello, world!'), compressGzip(), fh.writer({ autoClose: true }), ); } run().catch(console.error);
pipeToSync(source, ...transforms?, writer, options?): number
Synchronous version of pipeTo(). The source, all transforms, and the
writer must be synchronous. Cannot accept async iterables or promises.
The writer must have the *Sync methods (writeSync, writevSync,
endSync) and fail() for this to work.
pull(source, ...transforms?, options?): AsyncIterable
AsyncIterable | IterableObjectAbortSignalAsyncIterableUint8Array[]Create a lazy async pipeline. Source conversion and streamable protocol
dispatch occur when pull() is called, but data is not read from source
until the returned iterable is consumed. A signal that is already aborted is
thrown synchronously after source conversion. Transforms are applied in order.
import { from, pull, text } from 'node:stream/iter'; const asciiUpper = (chunks) => { if (chunks === null) return null; return chunks.map((c) => { for (let i = 0; i < c.length; i++) { c[i] -= (c[i] >= 97 && c[i] <= 122) * 32; } return c; }); }; const result = pull(from('hello'), asciiUpper); console.log(await text(result)); // 'HELLO'
const { from, pull, text } = require('node:stream/iter'); const asciiUpper = (chunks) => { if (chunks === null) return null; return chunks.map((c) => { for (let i = 0; i < c.length; i++) { c[i] -= (c[i] >= 97 && c[i] <= 122) * 32; } return c; }); }; async function run() { const result = pull(from('hello'), asciiUpper); console.log(await text(result)); // 'HELLO' } run().catch(console.error);
Using an AbortSignal:
import { pull } from 'node:stream/iter'; const ac = new AbortController(); const result = pull(source, transform, { signal: ac.signal }); ac.abort(); // Pipeline throws AbortError on next iteration
const { pull } = require('node:stream/iter'); const ac = new AbortController(); const result = pull(source, transform, { signal: ac.signal }); ac.abort(); // Pipeline throws AbortError on next iteration
pullSync(source, ...transforms?): Iterable
IterableIterableUint8Array[]Synchronous version of pull(). Source conversion and streamable protocol
dispatch occur when pullSync() is called. All transforms must be synchronous.