Kodokon kodokon.com

Streams et backpressure : traiter de gros fichiers

Traitez des fichiers de plusieurs gigaoctets avec une mémoire constante grâce aux streams, à pipeline et à la backpressure.

9 min · 3 questions

Ouvrir cette leçon dans Kodokon

readFile charge l'intégralité d'un fichier en mémoire : sur un log de 4 Go, votre processus explose. Les streams traitent les données par morceaux (chunks) avec une empreinte mémoire constante, bornée par le highWaterMark (64 Kio par défaut pour un flux fichier). Quatre types existent : Readable, Writable, Duplex et Transform. Le problème central qu'ils résolvent s'appelle la backpressure : que faire quand le producteur est plus rapide que le consommateur ?

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();
Le contrat write/drain, appliqué à la main.

Quand write() renvoie false, le buffer interne a dépassé highWaterMark : les données ne sont pas perdues, mais continuer à écrire accumule tout en mémoire - exactement ce que les streams devaient éviter. Le contrat : cesser d'écrire et attendre l'événement drain. C'est ce que pipe() fait pour vous, et que pipeline() fait mieux encore : en plus de la backpressure, il propage les erreurs et détruit proprement tous les flux, dans les deux sens.

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");
Transformation en flux, mémoire constante.

Les flux Readable sont des itérables asynchrones : for await consomme les chunks en respectant la backpressure. Attention, un chunk est une tranche d'octets arbitraire : rien ne garantit qu'il se termine sur une fin de ligne. Pour un traitement ligne à ligne, node:readline gère le recollage des morceaux, y compris une fin \r\n coupée entre deux chunks grâce à 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");
Lecture ligne à ligne d'un fichier volumineux.

Quiz de validation

Vérifiez que vous avez bien retenu les points clés de cette leçon.

  1. Que signifie le fait que writable.write(chunk) renvoie false ?
    • Le chunk a été perdu et doit être réécrit
    • Le buffer interne dépasse highWaterMark : cessez d'écrire jusqu'à l'événement drain
    • Le flux a été détruit par une erreur
  2. Quel avantage décisif pipeline() a-t-il sur pipe() ?
    • pipeline est plus rapide car multithreadé
    • pipe gère la backpressure, pas pipeline
    • pipeline propage les erreurs et détruit tous les flux impliqués
    • pipeline convertit automatiquement les encodages
  3. En objectMode, que représente highWaterMark ?
    • Un nombre d'objets en file d'attente
    • Un nombre d'octets, comme en mode binaire
    • La taille maximale d'un objet sérialisé