Streams permitem processar dados em partes, sem carregar o conteúdo inteiro na memória. No Node.js, elas aparecem em arquivos, HTTP, compressão, criptografia, processos filhos, sockets e bibliotecas de banco. O conceito central para manter estabilidade é backpressure: o produtor precisa desacelerar quando o consumidor não acompanha.
Ignorar backpressure pode aumentar memória, pressionar o garbage collector e derrubar o processo durante arquivos grandes ou picos de tráfego.
Tipos de stream
- Readable: produz dados;
- Writable: recebe dados;
- Duplex: lê e escreve;
- Transform: recebe, transforma e produz.
Exemplos incluem fs.createReadStream, resposta HTTP, sockets TCP e zlib.createGzip.
Copiando um arquivo
import { createReadStream, createWriteStream } from 'node:fs';
import { pipeline } from 'node:stream/promises';
await pipeline(
createReadStream('input.bin'),
createWriteStream('output.bin'),
);pipeline conecta as streams, propaga erros e fecha recursos. Prefira essa API a encadear pipe sem tratamento completo.
O que é backpressure
Uma Writable possui um buffer interno. Quando write() retorna false, o produtor deve parar e aguardar drain.
async function writeChunks(writable, chunks) {
for (const chunk of chunks) {
if (!writable.write(chunk)) {
await new Promise((resolve) => writable.once('drain', resolve));
}
}
writable.end();
}Continuar escrevendo após false não descarta dados; apenas aumenta a fila em memória.
Pipeline aplica backpressure
Ao conectar Readable e Writable com pipeline, o mecanismo pausa e retoma a leitura conforme o destino. Isso é uma das razões para usar abstrações de stream em exportações e uploads.
Compressão
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';
import { pipeline } from 'node:stream/promises';
await pipeline(
createReadStream('access.log'),
createGzip(),
createWriteStream('access.log.gz'),
);Cada etapa respeita o ritmo da próxima. Se o disco de destino for lento, leitura e compressão desaceleram.
Transform personalizada
import { Transform } from 'node:stream';
class UppercaseTransform extends Transform {
_transform(chunk, encoding, callback) {
try {
callback(null, chunk.toString().toUpperCase());
} catch (error) {
callback(error);
}
}
}Chame o callback uma única vez. Para processamento assíncrono, espere a operação antes de chamá-lo; não dispare várias tarefas ilimitadas.
Object mode
import { Transform } from 'node:stream';
const normalize = new Transform({
objectMode: true,
transform(user, encoding, callback) {
callback(null, {
id: user.id,
name: user.name.trim(),
});
},
});Em object mode, highWaterMark representa quantidade de objetos, não bytes. Objetos podem ter tamanhos muito diferentes, então monitore memória.
highWaterMark
O highWaterMark é um limiar do buffer, não um limite absoluto. Aumentá-lo pode melhorar throughput em alguns cenários, mas também aumenta memória e latência até que um chunk seja processado.
createReadStream('large.csv', {
highWaterMark: 64 * 1024,
});Teste tamanhos diferentes com carga real; não copie valores sem medir.
Iteração assíncrona
const stream = createReadStream('data.txt', { encoding: 'utf8' });
for await (const chunk of stream) {
await processChunk(chunk);
}O loop aguarda cada processamento e aplica controle natural. Cuidado: chunks não correspondem necessariamente a linhas ou registros completos.
Processando linhas
import { createInterface } from 'node:readline';
const rl = createInterface({
input: createReadStream('events.ndjson'),
crlfDelay: Infinity,
});
for await (const line of rl) {
if (!line.trim()) continue;
const event = JSON.parse(line);
await handleEvent(event);
}Trate linhas inválidas, limite tamanho e decida se um erro interrompe ou é enviado para dead letter.
Concorrência limitada
Processar uma linha de cada vez pode ser lento. Por outro lado, iniciar milhares de promises destrói o controle de fluxo. Use um limite:
const inFlight = new Set();
const limit = 8;
for await (const item of source) {
const task = processItem(item).finally(() => inFlight.delete(task));
inFlight.add(task);
if (inFlight.size >= limit) {
await Promise.race(inFlight);
}
}
await Promise.all(inFlight);Preservar ordem exige estratégia adicional. Confirme semântica de falha e retry.
HTTP download
import { Readable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
const response = await fetch(url);
if (!response.ok || !response.body) {
throw new Error(`HTTP ${response.status}`);
}
await pipeline(
Readable.fromWeb(response.body),
createWriteStream('download.bin'),
);Defina timeout, limite tamanho, valide URL e remova arquivo parcial em caso de erro.
HTTP upload e resposta
app.get('/export', async (req, res, next) => {
try {
res.setHeader('content-type', 'text/csv; charset=utf-8');
await pipeline(createCsvStream(), res);
} catch (error) {
if (!res.headersSent) next(error);
else res.destroy(error);
}
});Depois que headers e parte do corpo foram enviados, não é possível retornar uma resposta JSON de erro normal.
Cancelamento
Quando o cliente desconecta, interrompa banco, arquivo ou transformação:
const controller = new AbortController();
req.on('close', () => controller.abort());
await pipeline(source, transform, res, {
signal: controller.signal,
});Confirme que as fontes respeitam o sinal e liberam recursos.
Erros
Streams emitem error. pipeline centraliza tratamento:
try {
await pipeline(source, transform, destination);
} catch (error) {
logger.error({ error }, 'Pipeline falhou');
throw error;
}Evite combinar múltiplos listeners e promises de forma que o mesmo erro seja tratado duas vezes.
final e flush
Writable pode implementar _final; Transform pode implementar _flush para emitir dados restantes:
const transform = new Transform({
transform(chunk, encoding, callback) {
// acumular fragmentos
callback();
},
flush(callback) {
this.push(remainingData);
callback();
},
});Chunks e Unicode
Um caractere UTF-8 pode ser dividido entre chunks. Use setEncoding('utf8') ou StringDecoder quando transformar texto. Não faça chunk.toString() isoladamente em protocolos que podem dividir caracteres.
Buffers e cópias
Evite concatenar repetidamente buffers grandes. Cada concatenação pode copiar dados. Prefira pipeline, armazenamento em arquivo ou uma lista com tamanho total controlado.
Segurança
Defina limites para uploads, linhas, arquivos compactados e resposta externa. Streams evitam carregar tudo de uma vez, mas não impedem consumo infinito de disco, CPU ou tempo.
Observabilidade
Meça:
- bytes processados;
- taxa de transferência;
- duração;
- tempo aguardando drain;
- buffer e fila;
- erros por etapa;
- cancelamentos;
- memória e event loop delay.
Uma taxa baixa pode estar no disco, rede, compressão ou consumidor; crie métricas por etapa.
Testes
Teste chunks pequenos, fronteiras aleatórias, erro no meio, destino lento e cancelamento. Um teste com um único chunk não revela bugs de parsing.
import { Readable, Writable } from 'node:stream';
const source = Readable.from(['a', 'b', 'c']);
const output = [];
const destination = new Writable({
write(chunk, encoding, callback) {
output.push(chunk.toString());
callback();
},
});
await pipeline(source, destination);Erros comuns
- ignorar retorno de write;
- usar pipe sem tratar erros completos;
- acumular toda saída em memória;
- assumir que chunk é linha;
- iniciar concorrência ilimitada;
- aumentar highWaterMark sem medir;
- não cancelar após desconexão;
- ignorar Unicode dividido;
- não limitar upload e descompressão;
- não testar consumidor lento.
Fluxo recomendado
Use pipeline, respeite backpressure, limite concorrência e tamanho, propague cancelamento e monitore cada etapa. Para processos externos, combine com child_process; para medir gargalos, use perf_hooks; para carga, teste com Autocannon.
Consulte a documentação oficial de Streams e o guia oficial de backpressure no Node.js.



