A Web Streams API oferece interfaces padronizadas como ReadableStream, WritableStream e TransformStream. No Node.js, essas APIs aproximam o backend do modelo usado por Fetch, Service Workers, Deno e navegadores, facilitando código compartilhado e integração com respostas HTTP modernas.
O Node.js também possui streams clássicas em node:stream. As duas famílias podem coexistir e são convertidas por métodos como Readable.fromWeb e Readable.toWeb.
ReadableStream básica
const stream = new ReadableStream({
start(controller) {
controller.enqueue(new TextEncoder().encode('Olá '));
controller.enqueue(new TextEncoder().encode('mundo'));
controller.close();
},
});O controller envia chunks, fecha a stream ou reporta erro. Para texto, os chunks normalmente são Uint8Array.
Lendo com reader
const reader = stream.getReader();
try {
while (true) {
const { value, done } = await reader.read();
if (done) break;
console.log(new TextDecoder().decode(value));
}
} finally {
reader.releaseLock();
}Uma stream pode ficar bloqueada por um reader. Libere o lock quando terminar.
Iteração assíncrona
for await (const chunk of stream) {
process.stdout.write(new TextDecoder().decode(chunk));
}O suporte pode depender do ambiente e da forma como a stream foi criada. Use APIs compatíveis com a versão mínima do projeto.
Consumindo resposta do fetch
const response = await fetch('https://example.com/data.ndjson');
if (!response.ok || !response.body) {
throw new Error(`HTTP ${response.status}`);
}
const reader = response.body.getReader();response.body é uma Web ReadableStream. Isso permite processar downloads sem chamar arrayBuffer() ou text(), que acumulam o corpo inteiro.
TransformStream
const uppercase = new TransformStream({
transform(chunk, controller) {
controller.enqueue(String(chunk).toUpperCase());
},
});
const transformed = source.pipeThrough(uppercase);O transform recebe chunks e pode gerar zero, um ou vários chunks. Trate erros com controller.error(error) ou deixe a exceção rejeitar o pipeline.
TextDecoderStream
const textStream = response.body.pipeThrough(
new TextDecoderStream('utf-8'),
);
for await (const text of textStream) {
console.log(text);
}O decoder preserva caracteres divididos entre chunks, algo que um new TextDecoder().decode(chunk) isolado pode tratar incorretamente sem modo streaming.
TextEncoderStream
const bytes = textSource.pipeThrough(new TextEncoderStream());Essa transformação é útil ao gerar resposta ou enviar dados para uma API que espera bytes.
WritableStream
const writable = new WritableStream({
async write(chunk) {
await persistChunk(chunk);
},
close() {
console.log('Concluído');
},
abort(reason) {
console.error('Cancelado', reason);
},
});
await source.pipeTo(writable);pipeTo aplica controle de fluxo e retorna uma Promise que resolve ao concluir.
Backpressure
Web Streams usam estratégia de fila e desiredSize. Um produtor pode aguardar a solicitação do consumidor por meio de pull:
const stream = new ReadableStream({
async pull(controller) {
const chunk = await nextChunk();
if (chunk === null) {
controller.close();
return;
}
controller.enqueue(chunk);
},
cancel(reason) {
closeSource(reason);
},
});Não gere dados infinitamente no start. Deixe o consumidor solicitar quando possível.
QueuingStrategy
const stream = new ReadableStream(source, {
highWaterMark: 16,
size(chunk) {
return chunk.byteLength;
},
});A unidade depende de size. Ajustes incorretos podem aumentar memória ou reduzir throughput. Meça.
Convertendo Web Stream para Node Stream
import { Readable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
import { createWriteStream } from 'node:fs';
const response = await fetch(url);
if (!response.body) throw new Error('Resposta sem corpo');
await pipeline(
Readable.fromWeb(response.body),
createWriteStream('download.bin'),
);A conversão permite usar o ecossistema clássico de pipeline, zlib e arquivos.
Convertendo Node Stream para Web Stream
import { Readable } from 'node:stream';
import { createReadStream } from 'node:fs';
const webStream = Readable.toWeb(createReadStream('video.mp4'));Isso é útil em frameworks e APIs que esperam uma Web Stream.
Response com stream
const body = new ReadableStream({
start(controller) {
controller.enqueue(new TextEncoder().encode('primeiro\n'));
controller.enqueue(new TextEncoder().encode('segundo\n'));
controller.close();
},
});
const response = new Response(body, {
headers: { 'content-type': 'text/plain; charset=utf-8' },
});Frameworks baseados em padrões web podem devolver essa Response diretamente.
NDJSON
Dados delimitados por linha são adequados para streaming:
function ndjsonStream(items) {
const encoder = new TextEncoder();
return new ReadableStream({
async start(controller) {
try {
for await (const item of items) {
controller.enqueue(encoder.encode(`${JSON.stringify(item)}\n`));
}
controller.close();
} catch (error) {
controller.error(error);
}
},
});
}Para fontes muito grandes, prefira pull ou controle explícito, pois o loop em start pode produzir mais rápido que o consumidor.
Parser de linhas
class LineTransform {
constructor() {
let buffer = '';
return new TransformStream({
transform(chunk, controller) {
buffer += chunk;
const lines = buffer.split('\n');
buffer = lines.pop() || '';
for (const line of lines) controller.enqueue(line);
},
flush(controller) {
if (buffer) controller.enqueue(buffer);
},
});
}
}Combine com TextDecoderStream. Limite o tamanho da linha para impedir uso excessivo de memória.
Tee
const [forStorage, forHash] = source.tee();tee cria duas ramificações. Se uma consumir lentamente, a implementação pode acumular dados. Não use como duplicação gratuita de streams enormes sem medir.
Cancelamento
const controller = new AbortController();
await source.pipeTo(destination, {
signal: controller.signal,
});Ao cancelar, implemente cancel e abort para fechar arquivo, conexão ou cursor.
Fetch com timeout
const signal = AbortSignal.timeout(10_000);
const response = await fetch(url, { signal });Depois de abortar, trate a exceção e descarte recursos parciais. Verifique a versão mínima do Node.js adotada para utilitários de AbortSignal.
BYOB readers
Streams de bytes podem usar readers “bring your own buffer” para reduzir alocações em cenários avançados. A implementação exige controle cuidadoso de views e tamanhos. Só adote após profiling demonstrar pressão de alocação.
Integração com compressão
CompressionStream e DecompressionStream podem estar disponíveis conforme o runtime:
const compressed = source.pipeThrough(new CompressionStream('gzip'));Para compatibilidade ampla e opções avançadas, node:zlib continua relevante. Confirme formatos suportados.
Erros
Uma falha em qualquer etapa rejeita pipeTo. Use try/catch:
try {
await source
.pipeThrough(transform)
.pipeTo(destination);
} catch (error) {
logger.error({ error }, 'Stream falhou');
}Decida se o cancelamento da origem e o abort do destino devem ser prevenidos por opções, mas evite manter recursos abertos.
Observabilidade
Meça bytes, duração, taxa de transferência, cancelamentos, tamanho de fila, erros e tempo de processamento por chunk. Não gere métricas por chunk individual em streams muito intensas; agregue.
Testes
Teste chunks divididos, consumidor lento, erro de transformação, cancelamento, linhas grandes, Unicode e corpo vazio. Simule atrasos em write para confirmar backpressure.
Quando usar cada modelo
Use Web Streams quando trabalha com Fetch, padrões web, código compartilhado ou frameworks baseados em Request/Response. Use streams clássicas quando depende do ecossistema Node, EventEmitter, arquivos e módulos existentes. Converta na fronteira em vez de misturar APIs em todas as camadas.
Erros comuns
- acumular o corpo com text ou arrayBuffer;
- decodificar UTF-8 por chunk sem estado;
- produzir tudo no start;
- não implementar cancel;
- usar tee com consumidor lento;
- ignorar erros de pipeTo;
- confundir Web Streams e Node Streams;
- ajustar highWaterMark sem benchmark;
- não limitar linhas e downloads.
Fluxo recomendado
Escolha uma família de streams como padrão da camada, converta somente nas fronteiras, respeite backpressure, propague AbortSignal e teste falhas. Veja também Streams e Backpressure no Node.js, child_process e medições com perf_hooks.
Consulte a documentação oficial de Web Streams no Node.js e a referência da Streams API na MDN.



