Retour au cours

backend / nodejs

Streams avancés et backpressure

Leçon 161 exercice

Explication

Ce que vous allez apprendre

  • Comprendre le backpressure : ce qui se passe quand une source produit plus vite qu'une destination n'absorbe
  • Réagir correctement à la valeur de retour de write() et à l'événement drain
  • Préférer .pipe()/pipeline() à la gestion manuelle du backpressure
  • Ajuster highWaterMark pour équilibrer overhead et consommation mémoire

Dans quel contexte ?

Un service d'export de données copie des fichiers volumineux depuis un disque local rapide vers un stockage réseau plus lent. En développement, avec de petits fichiers de test, tout fonctionne parfaitement. En production, avec des fichiers de plusieurs gigaoctets, le processus Node.js crashe après quelques minutes avec une erreur de mémoire (heap out of memory). Cette leçon explique précisément pourquoi, et comment le backpressure évite ce genre d'incident.

Un problème invisible qui explose en production

La leçon sur les streams a montré comment traiter des fichiers volumineux sans tout charger en mémoire. Mais il reste un piège subtil à découvrir.

Que se passe-t-il si la SOURCE produit des données plus vite que la DESTINATION ne peut les absorber ? Par exemple, lire un fichier local très rapidement pendant qu'on l'envoie vers un stockage réseau lent.

Sans protection, les données s'accumulent dans un buffer en mémoire qui grossit indéfiniment. Jusqu'au crash du processus, un bug qui n'apparaît souvent qu'en production, sous charge réelle.

Piège fréquent

Ignorer la valeur de retour de writeStream.write(chunk) fonctionne très bien en test avec de petits fichiers, mais peut faire exploser la mémoire en production avec des fichiers volumineux ou un destinataire lent. Vérifiez toujours ce retour, ou utilisez .pipe()/pipeline() qui s'en charge pour vous.

Le backpressure est le mécanisme qui empêche exactement ce scénario. Concrètement, writeStream.write(chunk) renvoie false quand son buffer interne est plein.

C'est un signal explicite disant "ralentis, je ne peux plus rien absorber pour l'instant". Un code bien écrit répond à ce signal en mettant en pause la lecture, puis reprend uniquement quand l'événement drain indique que le buffer s'est vidé.

Gérer ce mécanisme de pause/reprise manuellement est possible, mais fastidieux et facile à mal implémenter. Heureusement, il existe une solution plus simple.

.pipe(), et surtout pipeline() déjà vu à la leçon 3, gèrent le backpressure AUTOMATIQUEMENT. C'est la raison principale pour laquelle ces méthodes sont recommandées en pratique plutôt que la gestion manuelle des événements.

Une fois ce mécanisme compris, un paramètre permet de l'ajuster finement : highWaterMark. Il détermine la taille du buffer interne avant que le stream ne signale un backpressure.

Un buffer trop petit multiplie les pauses et reprises, ce qui ajoute de l'overhead. Un buffer trop grand retarde la détection d'un déséquilibre et consomme plus de mémoire — un vrai compromis à ajuster selon la charge attendue.

Taille de highWaterMarkEffet
Trop petitepauses/reprises fréquentes, overhead CPU
Trop grandedétection tardive du déséquilibre, mémoire consommée en plus
Par défaut (16 Ko)bon compromis pour la plupart des cas

Le piège fréquent à connaître avant de pratiquer : ignorer la valeur de retour de write() fonctionne parfaitement, jusqu'au jour où le volume de données ou la lenteur du destinataire dépasse ce que la mémoire disponible peut absorber.

Commandes & code

Streams avancés et backpressure

js
// Le problème : un writable plus lent qu'un readable fait exploser la mémoire sans backpressure
import { createReadStream, createWriteStream } from "node:fs";

const readStream = createReadStream("./huge-file.bin");
const writeStream = createWriteStream("./destination.bin");

// Mauvais : ignore le retour de write(), peut charger tout en mémoire si writeStream est lent
readStream.on("data", (chunk) => {
  writeStream.write(chunk); // pas de gestion du backpressure
});
js
// Bon : respecter le retour de write() — false signifie "le buffer interne est plein"
readStream.on("data", (chunk) => {
  const canContinue = writeStream.write(chunk);

  if (!canContinue) {
    readStream.pause(); // stoppe la lecture tant que le writable n'a pas absorbé son buffer
  }
});

writeStream.on("drain", () => {
  readStream.resume(); // le writable a vidé son buffer, on peut reprendre
});

// En pratique, TOUJOURS préférer .pipe() ou pipeline() qui gèrent ça automatiquement :
readStream.pipe(writeStream);
js
// Transform stream avec contrôle du highWaterMark (taille du buffer interne)
import { Transform } from "node:stream";

class JsonLinesParser extends Transform {
  constructor(options) {
    super({ ...options, objectMode: true, highWaterMark: 16 }); // buffer 16 objets max
    this.buffer = "";
  }

  _transform(chunk, encoding, callback) {
    this.buffer += chunk.toString();
    const lines = this.buffer.split("\n");
    this.buffer = lines.pop();

    for (const line of lines) {
      if (line.trim()) {
        try {
          this.push(JSON.parse(line));
        } catch (err) {
          return callback(err); // erreur propagée dans le pipeline
        }
      }
    }
    callback();
  }

  _flush(callback) {
    if (this.buffer.trim()) {
      try {
        this.push(JSON.parse(this.buffer));
      } catch (err) {
        return callback(err);
      }
    }
    callback();
  }
}
js
// Traiter un flux JSONL volumineux avec contrôle de concurrence limité
import { pipeline } from "node:stream/promises";
import { Writable } from "node:stream";

async function processInBatches(records, batchSize, processFn) {
  for (let i = 0; i < records.length; i += batchSize) {
    await Promise.all(records.slice(i, i + batchSize).map(processFn));
  }
}

const batchWriter = new Writable({
  objectMode: true,
  highWaterMark: 100,
  async write(record, encoding, callback) {
    try {
      await db.records.insert(record); // insertion asynchrone par record
      callback();
    } catch (err) {
      callback(err);
    }
  },
});

await pipeline(
  createReadStream("./data.jsonl"),
  new JsonLinesParser(),
  batchWriter
);
js
// Stream Web API (fetch) vs stream Node — interopérabilité via Readable.fromWeb
import { Readable } from "node:stream";
import { createWriteStream } from "node:fs";

const response = await fetch("https://example.com/large-dataset.csv");
const nodeStream = Readable.fromWeb(response.body); // conversion ReadableStream Web -> Node

await pipeline(nodeStream, createWriteStream("./dataset.csv"));

Résumé

  • write() retourne false quand le buffer interne est plein : c'est le signal de backpressure.
  • .pipe() et pipeline() gèrent le backpressure automatiquement — les préférer à une gestion manuelle.
  • highWaterMark contrôle la taille du buffer interne d'un stream, ajustable selon la charge attendue.
  • Readable.fromWeb fait le pont entre les streams Web standards (fetch) et l'API stream de Node.

Exercices pratiques

1 disponible
1

Mission : un export qui plante uniquement en production, jamais en local

Objectif : Diagnostiquer un crash mémoire causé par l'absence de gestion du backpressure, puis corriger le code en respectant le signal de write() ou en migrant vers pipeline().

Contexte

Le service d'export copie des fichiers volumineux depuis un disque local rapide vers un stockage réseau (S3-compatible, plus lent). Le code actuel écoute readStream.on('data', chunk => writeStream.write(chunk)) sans jamais vérifier ce que write() retourne. En local avec des fichiers de test de quelques Ko, tout fonctionne. En production avec des fichiers de plusieurs gigaoctets vers un stockage réseau, le processus crashe après quelques minutes avec JavaScript heap out of memory.

Résoudre l’exercice →