Inicio / Artículos / Comprendiendo las secuencias de Node.js: el problema que realmente resuelven

Comprendiendo las secuencias de Node.js: el problema que realmente resuelven

Aprenda por qué existen las transmisiones de Node.js, cómo funciona internamente la canalización y qué significa realmente la contrapresión para manejar grandes volúmenes de datos de manera eficiente.

1261 palabras

La mayoría de los tutoriales sobre flujos comienzan con la interfaz de la API, .pipe(), Readable, Writable, Transform, sin abordar siquiera el problema real que estas herramientas fueron creadas para resolver. Enfocarse de esa manera hace que los flujos parezcan una ceremonia arbitraria que hay que memorizar. Invierte el orden, comienza con el problema, y la API en su mayor parte se explica por sí sola.

El problema: algunos datos son demasiado grandes para almacenarse de una sola vez

Imagina una tarea que requiere convertir una exportación CSV de 4 GB a JSON. El enfoque ingenuo sería el siguiente:

const fs = require("fs");

const data = fs.readFileSync("export.csv", "utf8");
const rows = data.split("\n").map(parseRow);
fs.writeFileSync("export.json", JSON.stringify(rows));

readFileSync no devolverá el control al programa hasta que todo el archivo de 4 GB haya sido cargado en memoria como una única cadena de JavaScript, y esa representación en forma de cadena suele consumir significativamente más memoria de lo que sugeriría el tamaño real del archivo. En una máquina con solo 2 GB de RAM disponible, este código no solo funciona de manera ineficiente, sino que también se cae completamente o deja sin recursos a todos los demás procesos que compiten por el mismo pool de memoria. No hay nada malo en la lógica en sí; todos los pasos de transformación son correctos. El verdadero defecto es una suposición oculta dentro de readFileSync: la idea de que está bien mantener todo un archivo en memoria al mismo tiempo, sin importar cuán grande sea ese archivo.

La idea real detrás de un flujo

En lugar de exigir todo el conjunto de datos desde el principio, la transmisión por flujo funciona según un principio diferente: se solicita una parte, se procesa y luego se solicita la siguiente. En ningún momento el archivo completo permanece en memoria; solo existe una pequeña porción en un momento dado, y esa porción se procesa y libera antes de que aparezca la siguiente.

const fs = require("fs");

const readStream = fs.createReadStream("export.csv", { encoding: "utf8" });

readStream.on("data", (chunk) => {
  console.log(`Received ${chunk.length} characters`);
});

readStream.on("end", () => {
  console.log("Done reading the whole file, piece by piece");
});

Obsérvese que el tamaño total del archivo nunca se tiene en cuenta en este código. Ya sea que el archivo de origen tenga 4 GB o 4 KB, este fragmento funciona de manera idéntica, aunque el consumo máximo de memoria difiere enormemente en los dos casos. Ese es precisamente el objetivo de la transmisión por flujo: se sustituye la necesidad de contar con el conjunto completo de datos al inicio por la posibilidad de comenzar de inmediato y terminar el trabajo sin mantener nunca en memoria más que un pequeño rango de datos.

Tuberización: Conectar una fuente directamente a un destino

Tirar de los datos por partes de forma manual, una a la vez, es un método útil como base, pero en la práctica es mucho más común conectar directamente una secuencia legible a una secuencia escribible, permitiendo que los datos fluyan automáticamente desde el origen hasta el destino en lugar de transmitir cada parte manualmente:

const fs = require("fs");

fs.createReadStream("export.csv")
  .pipe(fs.createWriteStream("export-copy.csv"));

Bastan dos líneas para duplicar un archivo de cualquier tamaño sin cargar todo en memoria. .pipe() no realiza ningún truco especial aquí; simplemente redirige el evento data emitido por la fuente a una llamada write en el destino, además de manejar un detalle adicional que resulta ser más importante que el mecanismo de copiado en sí.

La parte que casi nadie explica bien: la contrapresión

Este es el verdadero problema que aborda .pipe(), más allá de la simple conveniencia. Imagine un escenario en el que la lectura se realiza rápidamente, por ejemplo desde un disco local, mientras que la escritura ocurre lentamente, quizás a través de una conexión de red con ancho de banda limitado.

readStream.on("data", (chunk) => {
  writeStream.write(chunk); // what happens if this can't keep up?
});

Llamar a .write() más rápido de lo que el destino puede absorber realmente los datos no hace que el flujo escribible genere un error ni bloquee la ejecución. En lugar de eso, acumula silenciosamente el exceso en un búfer de memoria interno, guardándolo hasta que tenga la oportunidad de vaciarlo. Si esta discrepancia persiste el tiempo suficiente, con una brecha lo suficientemente grande entre la velocidad de lectura y la de escritura, el problema exacto que se pretendía eliminar con los flujos de datos vuelve a surgir silenciosamente: el uso de memoria crece sin límites, solo se pospone en lugar de ocurrir de inmediato, y suele ser mucho más difícil de detectar antes de que provoque una caída del sistema.

La contrapresión es la medida de protección contra exactamente este modo de fallo: permite a un flujo escribible indicar que ha alcanzado su capacidad y necesita que el productor reduzca la velocidad, y un productor implementado correctamente respeta esa señal en lugar de seguir enviando datos.

readStream.on("data", (chunk) => {
  const canContinue = writeStream.write(chunk);
  if (!canContinue) {
    readStream.pause(); // stop reading until the writable side catches up
  }
});

writeStream.on("drain", () => {
  readStream.resume(); // writable side is ready for more
});

En el momento en que el búfer interno de un flujo escribible supera su límite configurado, .write() devuelve false, lo cual sirve como señal para detener la producción hasta que el flujo emita un evento drain que indique que ha eliminado la acumulación y está listo para aceptar más datos. Este patrón de pausa y reanudación es precisamente lo que .pipe() gestiona automáticamente en segundo plano:

readStream.pipe(writeStream); // handles backpressure for you, silently, correctly

Esto, y no la brevedad, es la verdadera justificación para preferir .pipe() en lugar de conectar manualmente los listeners de data y write. Reenviar los datos por separado ignorando el valor de retorno de .write() vuelve a introducir el mismo problema de memoria ilimitada que los flujos estaban diseñados para evitar desde un principio; es simplemente un paso más lejos del error evidente inherente a readFileSync.

Flujos de Transformación: Procesamiento de Datos en Tránsito

Hay casos en los que no basta con mover los datos sin modificarlos de un lugar a otro, también es necesario darles una nueva forma durante el proceso. Un flujo Transform está diseñado precisamente para esto: se encuentra en medio de una cadena de tuberías, recibe los datos que llegan, aplica alguna operación a cada uno de ellos y reenvía el resultado al siguiente componente:

const { Transform } = require("stream");

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

fs.createReadStream("input.txt")
  .pipe(upperCaseTransform)
  .pipe(fs.createWriteStream("output.txt"));

Cada bloque se convierte a mayúsculas a medida que pasa por la cadena de procesamiento, y en ningún momento es necesario que todo el archivo, ya sea en su forma original o convertida, esté completamente en memoria. Este es precisamente el mecanismo detrás de zlib.createGzip() incorporado en Node. No es más que un flujo Transform que comprime cada bloque a medida que llega, y al igual que cualquier transformación, puede integrarse directamente en una cadena de tuberías de la misma manera que se hizo con el ejemplo de mayúsculas:

const zlib = require("zlib");

fs.createReadStream("export.csv")
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream("export.csv.gz"));

La lectura, compresión y escritura ocurren al mismo tiempo, procesando pequeñas porciones de datos a medida que llegan, sin que ninguna de las tres etapas necesite cargar todo el archivo al mismo tiempo.

Por qué vale la pena entenderlo realmente

Los streams tienen la reputación de ser uno de los aspectos más complicados de la API de Node, y honestamente su interfaz bruta basada en eventos sí parece poco intuitiva, por lo que esa reputación no está del todo injustificada. Pero el concepto subyacente es simple: evitar cargar todo el conjunto de datos en memoria, trabajar con él por partes y asegurarse de que un productor rápido no inunde silenciosamente a un consumidor lento al hacerlo. Una vez que internalices este modelo como el enfoque adecuado, .pipe(), Transform y la contrapresión dejan de parecer tres APIs independientes que hay que memorizar por separado. En realidad se trata de una sola idea, presentada a través de tres componentes relacionados de la API, que resuelve un problema con el que algo como readFileSync nunca tuvo que lidiar, ya que desde un principio no estaba diseñado para manejar nada más allá de archivos pequeños.

Lecturas relacionadas