Node.js: Streams no Node.js
Última atualização: 2026-08-26
Charlie recebeu uma tarefa urgente: analisar 2 GB de logs de acesso no servidor. Assim que ele os carregou usando o fs.readFile, a memória esgotou-se imediatamente — causando uma falha por falta de memória (OOM). Seu chefe de equipe disse: “Não importa o tamanho do arquivo, não tente processá-lo todo de uma vez — use o Stream para processá-lo aos poucos!” Charlie mudou para o createReadStream para processar os logs linha por linha, e o uso de memória permaneceu estável em 50 MB. A partir daí, ele passou a dizer a todos que encontrava: “O Stream é o canivete suíço para o processamento de big data no Node.js.”
1. Conceitos básicos de fluxos
Stream é uma interface abstrata do Node.js para o processamento de dados em fluxo contínuo, em que os dados são transmitidos em blocos, como um fluxo de água, em vez de serem carregados na memória de uma só vez.
- Conceito-chave: Processar dados em blocos, permitindo que as operações sejam iniciadas sem esperar que todos os dados estejam disponíveis
- Relação de herança: Todos os streams são instâncias de
EventEmittere utilizam eventos para notificar mudanças de estado. - Casos de uso: Leitura e gravação de arquivos grandes, transmissão em rede, compactação e conversão de dados
- Módulos integrados:
fs,zlib,crypto,httpe outros oferecem interfaces de fluxo. - Objetos globais: O módulo
streamfornece as classes baseReadable,Writable,DuplexeTransform
(1) Quatro tipos de fluxos
| Tipo | Descrição | Cenários típicos | Entrada | Saída |
|---|---|---|---|---|
| Legível | Fluxo legível, fonte de dados | fs.createReadStream, process.stdin |
Nenhum | Sim |
| Gravável | Fluxo gravável, ponto de extremidade de dados | fs.createWriteStream, process.stdout |
Sim | Não |
| Duplex | Duplex, leitura/gravação (independente) | net.Socket, tls.Socket |
Sim | Sim |
| Transformar | Fluxo de transformação; a saída é uma transformação da entrada | zlib.createGzip, crypto.createCipheriv |
Sim | Sim |
(2) Dois tipos de fluxo
O fluxo Readable possui dois modos de operação:
| Recurso | Modo contínuo | Modo pausado |
|---|---|---|
| Recuperação de dados | Envio automático; consumo por meio do evento data |
É necessário chamar manualmente read() para fazer a recuperação |
| Método de acionamento | Adicionar ouvinte data / Chamar pipe() / Chamar resume() |
Estado padrão inicial / Chamar pause() |
| Tratamento da contrapressão | pipe() Tratamento automático; a operação manual requer monitoramento drain |
Ritmo controlado pelo consumidor, contrapressão natural |
| Casos de uso | Fontes de dados contínuas e de alto rendimento | Requer controle preciso sobre o tempo de leitura |
2. Uma explicação detalhada sobre eventos de stream
Todos os fluxos são baseados em EventEmitter e utilizam um mecanismo orientado a eventos para transmitir dados e estado.
(1) Eventos gerais
| Evento | Condição de acionamento | Fluxo aplicável | Parâmetros de retorno de chamada |
|---|---|---|---|
data |
Novo bloco recebido | Legível | chunk |
end |
Dados lidos | Legível | Nenhum |
error |
Erro | Todas as transmissões | Error |
close |
Fechar recursos de baixo nível | Todas as transmissões | Nenhuma |
finish |
end() Após todos os dados terem sido gravados |
Gravável | Nenhum |
▶ Exemplo: Monitorando eventos do fluxo “Readable”
const fs = require('fs');
const rs = fs.createReadStream('./access.log', { highWaterMark: 64 * 1024 });
rs.on('data', (chunk) => {
console.log(`Received ${chunk.length} Byte`);
});
rs.on('end', () => {
console.log('Finished reading');
});
rs.on('error', (err) => {
console.error('Output error:', err.message);
});
▶ Exemplo: Monitorando eventos de fluxos graváveis
const ws = fs.createWriteStream('./output.txt');
ws.on('finish', () => {
console.log('All data has been written');
});
ws.on('error', (err) => {
console.error('Write error:', err.message);
});
ws.write('Hello Stream');
ws.end();
dataOs eventos são acionados automaticamente no modo de fluxo- A diferença entre
endefinish:endindica que a leitura foi concluída, efinishindica que a gravação foi concluída. errorÉ necessário monitorar os eventos; caso contrário, as exceções não capturadas farão com que o processo trave.closeNem todos os fluxos acionam o evento; isso depende da implementação subjacente.
3. Encadeamento com pipe()
pipe() é o método mais poderoso do Stream; ele conecta automaticamente a saída de um fluxo de leitura à entrada de um fluxo de gravação e lida automaticamente com a contrapressão.
▶ Exemplo:(1) Uso básico do “pipe”
readable.pipe(writable);
pipe() Retorna o fluxo de destino, para que possa ser encadeado:
readable.pipe(transform1).pipe(transform2).pipe(writable);
▶ Exemplo: Cópia de arquivos
const fs = require('fs');
fs.createReadStream('./source.txt')
.pipe(fs.createWriteStream('./dest.txt'));
▶ Exemplo: Transmissão de respostas HTTP
const http = require('http');
const fs = require('fs');
http.createServer((req, res) => {
res.writeHead(200, { 'Content-Type': 'text/plain' });
fs.createReadStream('./large.txt').pipe(res);
}).listen(3000);
▶ Exemplo:(2) Diagrama de fluxo de dados do tubo de fluxo
flowchart LR
A[Readable<br/>Data Source] -->|chunk| B[Transform<br/>Data Transformation]
B -->|chunk| C[Writable<br/>Data Endpoint]
C -.->|backpressure| A
B -.->|backpressure| A
style A fill:#4CAF50,color:#fff
style B fill:#FF9800,color:#fff
style C fill:#2196F3,color:#fff
pipe()Gerencia automaticamente a taxa de fluxo de dados; o lado de leitura entra em pausa quando o lado de gravação está ocupadopipe()Os erros não são tratados; os eventos devem ser monitorados separadamente para cada fluxoerror- O uso de
stream.pipeline()lida automaticamente com a limpeza e a propagação de erros, tornando-o mais seguro do quepipe()
4. Fluxos do sistema de arquivos e operações com arquivos
O módulo fs fornece fluxos de leitura e gravação relacionados ao sistema de arquivos.
(1) createReadStream / createWriteStream
| Opção | Descrição | Padrão |
|---|---|---|
highWaterMark |
Buffer size (bytes) | Readable: 64KB / Writable: 16KB |
encoding |
Codificação | null (Buffer) |
start |
Posição do byte inicial | 0 |
end |
Posição do byte final (inclusive) | Infinito |
flags |
Sinais de abertura do arquivo | Legível: r / Gravável: w |
▶ Exemplo: Como ler um segmento de arquivo dentro de um intervalo especificado
const fs = require('fs');
const rs = fs.createReadStream('./big.bin', {
start: 100,
end: 199,
highWaterMark: 32
});
rs.on('data', (chunk) => {
console.log(chunk.length);
});
(2) Comparação entre readFile e createReadStream
| Comparação | readFile | createReadStream |
|---|---|---|
| Uso de memória | Todos os arquivos carregados na memória | Utiliza apenas o tamanho do highWaterMark |
| Atraso na inicialização | Aguardar até que todos os arquivos tenham sido lidos antes de chamar a função de retorno | Retornar o fluxo imediatamente e processá-lo à medida que for lido |
| Tamanhos de arquivo adequados | Arquivos pequenos (<10 MB) | Arquivos grandes ou dados contínuos |
| Tratamento de erros | Recuperação de um erro em uma função de retorno | Monitoramento de eventos error |
| Pode ser pausado/retomado | Não é compatível | pause() / resume() |
▶ Exemplo: Processamento de arquivos grandes bloco por bloco
const fs = require('fs');
let totalBytes = 0;
const rs = fs.createReadStream('./2gb.log');
rs.on('data', (chunk) => {
totalBytes += chunk.length;
});
rs.on('end', () => {
console.log(`Total ${totalBytes} Byte`);
});
- Arquivos grandes devem ser processados por meio de fluxos para evitar erros de OOM
highWaterMarkQuanto menor for a memória, maior será a eficiência energética, mas mais frequentes serão as chamadas do sistema.- As opções
start/endpermitem a leitura de arquivos em blocos
5. Mecanismo de contrapressão
A contrapressão é fundamental para o controle de fluxo: quando o lado de gravação não consegue acompanhar a taxa de envio do lado de leitura, uma “contrapressão” é aplicada ao lado de gravação para interromper temporariamente o lado de leitura e evitar o acúmulo de dados na memória.
(1) A função pipe() lida automaticamente com a contrapressão
Ao usar o pipe(), a contrapressão é gerenciada automaticamente de forma interna pelo Node.js, sem a necessidade de intervenção manual.
(2) Controle manual da contrapressão
Quando pipe() não for utilizado, é necessário avaliar manualmente o valor de retorno de write():
const fs = require('fs');
const rs = fs.createReadStream('./source.txt');
const ws = fs.createWriteStream('./dest.txt');
rs.on('data', (chunk) => {
const canContinue = ws.write(chunk);
if (!canContinue) {
rs.pause();
ws.once('drain', () => {
rs.resume();
});
}
});
rs.on('end', () => {
ws.end();
});
▶ Exemplo: Observando o efeito da contrapressão
const rs = fs.createReadStream('./big.log', { highWaterMark: 1024 });
const ws = fs.createWriteStream('./out.log', { highWaterMark: 512 });
let paused = 0;
rs.on('data', (chunk) => {
const ok = ws.write(chunk);
if (!ok) {
paused++;
rs.pause();
ws.once('drain', () => rs.resume());
}
});
rs.on('end', () => {
console.log(`Backpressure triggered ${paused} times`);
ws.end();
});
write()Retornofalseindica que o buffer interno está cheio e que a leitura deve ser interrompidadrainEste evento indica que o buffer foi esvaziado e que a gravação pode ser retomada.- Ignorar a contrapressão fará com que o uso de memória continue a aumentar, levando, eventualmente, a um erro de falta de memória (OOM)
pipeline()é mais recomendado do quepipe(); ele lida automaticamente com a propagação de erros e a limpeza de recursos.
6. Transform Stream
Um fluxo Transform é uma subclasse de Duplex; sua saída é o resultado da transformação da entrada, e ele é comumente usado para compressão de dados, criptografia e conversão de formatos.
(1) A diferença entre “Transform” e “Duplex”
| Item de comparação | Duplex | Transformação |
|---|---|---|
| Relação entre entrada e saída | Independentes; não se influenciam mutuamente | A saída é gerada pela transformação da entrada |
| Métodos obrigatórios | _read() + _write() |
_transform() |
| Usos típicos | Sockets de rede | Compressão, criptografia, conversão de dados |
| Buffer interno | Um para leitura, outro para gravação | Estado intermediário da transformação |
▶ Exemplo: Fluxo de transformação personalizado (conversão de maiúsculas e minúsculas)
const { Transform } = require('stream');
const upper = new Transform({
transform(chunk, encoding, callback) {
callback(null, chunk.toString().toUpperCase());
}
});
process.stdin.pipe(upper).pipe(process.stdout);
▶ Exemplo: arquivos compactados com zlib
const fs = require('fs');
const zlib = require('zlib');
fs.createReadStream('./access.log')
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream('./access.log.gz'));
▶ Exemplo: Descompactação de um arquivo zlib
const fs = require('fs');
const zlib = require('zlib');
fs.createReadStream('./access.log.gz')
.pipe(zlib.createGunzip())
.pipe(fs.createWriteStream('./access_restored.log'));
- Em
_transform(chunk, encoding, callback),callback(null, data)envia os resultados da transformação callback(err)Erro aceitávelzlib.createGzip()/createGunzip()são os fluxos de transformação integrados mais comumente utilizados- Você também pode implementar o processamento
_flush(callback)para os dados restantes ao personalizar o fluxo de transformação.
7. Como pausar e retomar uma transmissão
O fluxo “Readable” está no modo de pausa por padrão e pode ser ativado ou desativado usando pause() e resume().
▶ Exemplo: Controles de pausa e retomada
const fs = require('fs');
const rs = fs.createReadStream('./big.log');
let count = 0;
rs.on('data', (chunk) => {
count++;
if (count % 10 === 0) {
rs.pause();
console.log(`Processed ${count} chunks, pause 1 second`);
setTimeout(() => rs.resume(), 1000);
}
});
▶ Exemplo: Lendo um arquivo linha por linha (readline)
const fs = require('fs');
const readline = require('readline');
const rl = readline.createInterface({
input: fs.createReadStream('./access.log'),
crlfDelay: Infinity
});
rl.on('line', (line) => {
if (line.includes('ERROR')) {
console.log(line);
}
});
rl.on('close', () => {
console.log('The file has been read.');
});
readline.createInterfaceEnvolve um fluxoReadablecomo uma interface para leitura linha por linhacrlfDelay: InfinityLida corretamente com as quebras de linha\r\n- Após
pause(), o eventodatanão é mais acionado até queresume()seja chamado. pipe()mudará a transmissão para o modo ao vivo
8. Exemplo abrangente: Pipeline de processamento de logs
Crie um fluxo completo de processamento de dados: Leia arquivos de log grandes → Analise linha por linha → Filtre as linhas com erros → Comprima a saída.
const fs = require('fs');
const zlib = require('zlib');
const { Transform } = require('stream');
const filterError = new Transform({
transform(chunk, encoding, callback) {
const lines = chunk.toString().split('\n');
const errors = lines.filter(l => l.includes('ERROR')).join('\n');
callback(null, errors ? errors + '\n' : '');
}
});
const src = fs.createReadStream('./app.log');
const dest = fs.createWriteStream('./errors.log.gz');
src
.pipe(filterError)
.pipe(zlib.createGzip())
.pipe(dest);
dest.on('finish', () => {
console.log('The error log has been written in compressed form.');
});
src.on('error', (err) => {
console.error('Read failed:', err.message);
});
- Cada etapa do pipeline é um fluxo; os dados fluem em blocos, o que resulta em um uso de memória extremamente baixo.
pipe()Lida automaticamente com a contrapressão entre fluxos adjacentes- Os erros devem ser monitorados separadamente em cada fluxo ou tratados de maneira uniforme usando
pipeline() - Transformar o fluxo
filterErrorpara reter apenas as linhas que contenhamERROR
❓ Perguntas Frequentes
P: O que é um stream? R: Um stream é uma interface abstrata para o processamento de dados; os dados podem ser lidos ou gravados em blocos, sem a necessidade de serem carregados na memória de uma só vez.
P: Qual é a diferença entre um fluxo legível (Readable) e um fluxo gravável (Writable)? R: Um fluxo legível (Readable) é uma fonte de dados da qual é possível ler dados; um fluxo gravável (Writable) é um destino de dados no qual é possível gravar dados.
P: Quando se deve usar streams? R: Em situações como o manuseio de arquivos grandes, transmissão de rede e processamento de dados em tempo real — quando o volume de dados é grande ou não pode ser carregado na memória de uma só vez.
P: O que o método
pipefaz? R:pipeconecta a saída de um fluxo de leitura à entrada de um fluxo de gravação, gerenciando automaticamente o fluxo de dados e a contrapressão.
P: O que é contrapressão? R: Quando a taxa de saída de dados do fluxo de gravação é menor do que a taxa de entrada de dados do fluxo de leitura, o fluxo de gravação envia um sinal ao fluxo de leitura para interromper a leitura; esse é o mecanismo de contrapressão.
- Quando se deve usar o Stream? Use o Stream quando a quantidade de dados exceder a memória disponível, quando for necessário processar os dados à medida que são lidos ou ao trabalhar com fontes de dados em tempo real; para arquivos pequenos, é mais simples usar o
readFile. - O
pipe()lida automaticamente com a contrapressão? Sim, opipe()verifica internamente o valor de retorno da extremidade de gravaçãowrite()e chama automaticamente opause()/resume()para gerenciar a vazão. - Como faço para ler um arquivo linha por linha? Use
readline.createInterface({ input: createReadStream(path) })e aguarde o eventolinepara recuperar os dados linha por linha. - Qual é a diferença entre Transform e Duplex? No Duplex, a leitura e a gravação são independentes e não estão relacionadas; no Transform, a saída é gerada pela transformação da entrada, exigindo apenas a implementação do método
_transform(). - Como os erros de stream devem ser tratados? Monitore os eventos
errorem cada stream ou usestream.pipeline()para propagar erros automaticamente e liberar recursos, evitando assim falhas no processo causadas por exceções não capturadas. - Qual é o valor adequado para
highWaterMark? O valor padrão de 64 KB é adequado para a maioria dos cenários; para processar grandes blocos de dados binários, é possível aumentá-lo para 256 KB e, quando a memória for limitada, reduzi-lo para 16 KB. - Qual você deve escolher:
pipe()oupipeline()? Recomendamos usarstream.pipeline(), pois ele propaga erros automaticamente e libera recursos;pipe()não lida com erros e não fecha fluxos automaticamente.
📖 Resumo
- Conceitos fundamentais e uso básico de streams
- Uma explicação detalhada dos conceitos fundamentais e do uso dos eventos de stream
- Conceitos básicos e uso de conexões encadeáveis com
pipe() - Conceitos básicos e uso de fluxos do sistema de arquivos e operações com arquivos
- Conceitos fundamentais e utilização do mecanismo de contrapressão
- Conceitos básicos e uso do Transform Stream
- Conceitos-chave e uso da pausa e retomada de streams
- Exemplo abrangente: conceitos básicos e uso do pipeline de processamento de logs
📝 Exercícios
- Conclua todos os exemplos de código desta lição e certifique-se de que cada um deles seja executado corretamente.
- Modifique o exemplo completo e adicione suas próprias extensões
- Analise a documentação oficial, identifique 1 ou 2 APIs que não foram abordadas nesta aula e escreva um código de teste para elas.
- Reflexão: Como você aplicaria o que aprendeu nesta aula a um projeto do mundo real?
- Tente combinar o que você aprendeu nesta aula com o conteúdo das aulas anteriores para criar um pequeno projeto.