stream/promises.pipeline
stream/promises.pipeline conecta Readable → Transform → Writable con la propagación correcta de errores, la destrucción del stream en caso de fallo y una API de Promesas adecuada para manejadores async/await.
Busca en todas las páginas de la documentación
stream/promises.pipeline conecta Readable → Transform → Writable con la propagación correcta de errores, la destrucción del stream en caso de fallo y una API de Promesas adecuada para manejadores async/await.
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';
await pipeline(
createReadStream('input.log'),
createGzip(),
createWriteStream('input.log.gz'),
);Cuándo usarlo:
.on('error') en pipes manualesimport { pipeline } from 'node:stream/promises';
import { createReadStream } from 'node:fs';
import { createServer } from 'node:http';
import { Transform } from 'node:stream';
import { createGzip } from 'node:zlib';
const lineCount = new Transform({
transform(chunk, _enc, cb) {
const lines = String(chunk).split('\n').length - 1;
(this as Transform & { lines?: number }).lines =
((this as Transform & { lines?: number }).lines ?? 0) + lines;
cb(null, chunk);
},
});
const server = createServer(async (req, res) => {
if (req.url !== '/export') {
res.writeHead(404).end();
return;
}
try {
res.writeHead(200, {
'content-type': 'application/gzip',
'content-disposition': 'attachment; filename="export.log.gz"',
});
await pipeline(
createReadStream('app.log'),
lineCount,
createGzip(),
res,
);
console.log('lines processed', (lineCount as Transform & { lines?: number }).lines);
} catch (err) {
if (!res.headersSent) res.writeHead(500);
res.end();
console.error('pipeline failed', err);
}
});
server.listen(3000);Lo que esto demuestra:
pipeline acepta el res Writable de HTTP como sumidero finalheadersSent requieren finalizar la respuesta sin una segunda línea de estadoawait se integra con try/catch como cualquier E/S asíncronapipeline destruye los streams participantes con el error.end.node:stream; la variante de Promesa es preferida en código asíncrono.| Característica | pipeline | pipe manual |
|---|---|---|
| Reenvío de errores | Sí | Manual |
| Destruir en caso de fallo | Sí | Manual |
| API de Promesas | Sí | No |
| AbortSignal | Sí | Manual |
import { pipeline } from 'node:stream/promises';
import type { Readable, Writable } from 'node:stream';
export async function safePump(
source: Readable,
sink: Writable,
signal?: AbortSignal,
): Promise<void> {
await pipeline(source, sink, { signal });
}writeHead antes de pipeline para enviar la respuesta. Solución: establece las cabeceras primero.req.on('aborted') opcionalmente. Solución: pasa AbortSignal vinculado a la solicitud._flush para los datos almacenados en búfer restantes.| Alternativa | Usar cuándo | No usar cuándo |
|---|---|---|
finished + destrucción manual | Código de streams heredado | Nuevo desarrollo |
pump (npm) | Versiones antiguas de Node | Node 24 tiene pipeline nativo |
| Almacenar todo el payload en búfer | Archivos pequeños < 1 MB | Descargas grandes |
Web pipeThrough | fetch Web Streams | fuentes node:fs |
Sí, en caso de éxito o fallo, los streams participantes se destruyen/finalizan adecuadamente.
Sí, argumentos variádicos: pipeline(a, b, c, d, sink).
Pasa la opción { signal: abortController.signal } en Node 24 pipeline.
ENOENT en archivo faltante, ECONNRESET en desconexión del cliente, errores de zlib en gzip corrupto.
API de callback heredada; la versión de Promesa es preferida en manejadores asíncronos.
Sí, los datos fluyen de izquierda a derecha a través de cada Transform.
Todos los streams en la cadena deben estar de acuerdo en el modo de objeto vs. bytes (con un Transform de conversión).
Un stream PassThrough intermedio que cuenta chunk.length.
El mismo patrón: await pipeline(source, res) dentro de una ruta asíncrona con un envoltorio de errores.
Usa reply.send(stream) o pipeline en la respuesta cruda según la documentación de Fastify 5 para un control preciso.
Usa stream/consumers.buffer o un Writable que recolecte fragmentos en memoria para las aserciones.
Heredado de la semántica de los streams; pipeline no elimina la necesidad de un highWaterMark sensato.
Versiones de la pila: Esta página fue escrita para Node.js 24.18.0 (LTS activa), npm 10+, TypeScript 5.6+, Express 5, Fastify 5 y NestJS 11.
Revisado por Chris St. John·Última actualización: 19 jul 2026