Kodokon kodokon.com

Streams y backpressure: manejar archivos grandes

Procesa archivos de varios gigabytes con memoria constante gracias a los streams, pipeline y backpressure.

9 min · 3 preguntas

Abrir esta lección en Kodokon

readFile carga un archivo entero en memoria: con un log de 4 GB, tu proceso revienta. Los streams procesan los datos en fragmentos con una huella de memoria constante, acotada por highWaterMark (64 KiB por defecto para un stream de archivo). Existen cuatro tipos: Readable, Writable, Duplex y Transform. El problema central que resuelven se llama backpressure: ¿qué haces cuando el productor es más rápido que el consumidor?

JAVASCRIPT
import { createWriteStream } from "node:fs";

const out = createWriteStream("big.txt");
let i = 0;

function writeChunks() {
  let ok = true;
  while (i < 1e6 && ok) {
    ok = out.write(`line ${i}\n`);
    i += 1;
  }
  if (i < 1e6) out.once("drain", writeChunks);
  else out.end();
}

writeChunks();
El contrato write/drain, aplicado a mano.

Cuando write() devuelve false, el búfer interno ha superado highWaterMark: los datos no se pierden, pero seguir escribiendo acumula todo en memoria - exactamente lo que los streams pretendían evitar. El contrato: deja de escribir y espera el evento drain. Eso es lo que pipe() hace por ti, y lo que pipeline() hace aún mejor: además del backpressure, propaga los errores y destruye limpiamente todos los streams, en ambas direcciones.

JAVASCRIPT
import { createReadStream, createWriteStream }
  from "node:fs";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";

const upper = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, chunk.toString().toUpperCase());
  },
});

await pipeline(
  createReadStream("big.txt"),
  upper,
  createWriteStream("big-upper.txt"),
);
console.log("done");
Transformación en streaming, memoria constante.

Los streams Readable son iterables asíncronos: for await consume los fragmentos respetando el backpressure. Cuidado, un fragmento es un trozo de bytes arbitrario: nada garantiza que termine en un límite de línea. Para el procesamiento línea por línea, node:readline se encarga de reensamblar las piezas, incluido un \r\n partido entre dos fragmentos gracias a crlfDelay: Infinity.

JAVASCRIPT
import { createReadStream } from "node:fs";
import { createInterface } from "node:readline";

const rl = createInterface({
  input: createReadStream("big.txt"),
  crlfDelay: Infinity,
});

let count = 0;
for await (const line of rl) {
  if (line.includes("42")) count += 1;
}
console.log(count, "matching lines");
Leer un archivo grande línea por línea.

Prueba de conocimientos

Comprueba que has retenido los puntos clave de esta lección.

  1. ¿Qué significa que writable.write(chunk) devuelva false?
    • El fragmento se perdió y hay que reescribirlo
    • El búfer interno supera highWaterMark: deja de escribir hasta el evento drain
    • El stream fue destruido por un error
  2. ¿Qué ventaja decisiva tiene pipeline() sobre pipe()?
    • pipeline es más rápido porque es multihilo
    • pipe maneja el backpressure, pipeline no
    • pipeline propaga los errores y destruye todos los streams implicados
    • pipeline convierte automáticamente las codificaciones
  3. En objectMode, ¿qué representa highWaterMark?
    • Un número de objetos en cola
    • Un número de bytes, como en modo binario
    • El tamaño máximo de un objeto serializado