Accueil / Articles / Comprendre les flux Node.js : le problème qu’ils résolvent réellement

Comprendre les flux Node.js : le problème qu’ils résolvent réellement

Découvrez pourquoi les flux Node.js existent, comment le transfert par canal fonctionne en interne, et ce que signifie réellement la contrainte de retour pour gérer efficacement de grandes quantités de données.

1261 mots

La plupart des tutoriels sur les flux commencent par présenter l’API, .pipe(), Readable, Writable, Transform, sans jamais aborder le problème réel que ces outils sont conçus pour résoudre. Une telle approche fait que les flux semblent être une série de procédures arbitraires à mémoriser. Inversez l’ordre, commencez par le problème, et l’API s’explique presque d’elle-même.

Le problème : certaines données sont trop volumineuses pour être stockées en une seule fois

Imaginez une tâche qui consiste à convertir un fichier CSV de 4 Go en JSON. L’approche naïve serait la suivante :

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 ne restitue le contrôle au programme qu’après que tout le fichier de 4 Go n’a pas été chargé en mémoire sous forme de chaîne JavaScript unique, et cette représentation en chaîne consomme généralement beaucoup plus de mémoire que ce à quoi la taille brute du fichier pourrait le faire penser. Sur une machine disposant seulement de 2 Go de RAM, ce code non seulement fonctionne de manière inefficace, mais il plante également, ou bien il prive tous les autres processus qui luttent pour accéder au même pool mémoire. Il n’y a rien de mal dans la logique elle-même ; toutes les étapes de transformation sont correctes. Le véritable défaut réside dans une hypothèse cachée au sein de readFileSync : à savoir qu’il est acceptable de conserver tout un fichier en mémoire simultanément, quel que soit son volume.

L’idée réelle derrière un flux

Au lieu de demander l’ensemble des données d’un coup, le streaming fonctionne selon un principe différent : on demande une partie, on la traite, puis on demande la suivante. À aucun moment le fichier complet n’est stocké en mémoire. Seule une petite partie existe à un instant donné, et cette partie est traitée puis libérée avant que la suivante n’arrive.

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");
});

Remarquez que la taille totale du fichier n’est jamais prise en compte dans ce code. Que le fichier source soit de 4 Go ou de 4 Ko, ce fragment de code s’exécute de la même manière, bien que l’empreinte mémoire maximale diffère énormément dans les deux cas. C’est justement l’objectif du streaming : on remplace la nécessité de disposer des données complètes avant de commencer par la possibilité de démarrer immédiatement et d’achever le travail sans jamais conserver en mémoire plus qu’une petite portion de données.

Tuyautage : relier directement une source à une destination

Tirer manuellement les données par morceaux un à un constitue une méthode utile, mais en pratique il est bien plus courant de relier directement un flux lisible à un flux écrivable, permettant ainsi aux données de circuler automatiquement de la source au destin plutôt que de transmettre chaque morceau manuellement :

const fs = require("fs");

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

Deux lignes suffisent pour dupliquer un fichier de n’importe quelle taille sans charger l’intégralité du contenu en mémoire. .pipe() ne réalise ici aucun miracle : il se contente de rediriger l’événement data émis par la source vers une appel write sur le destin, en gérant en plus un détail supplémentaire qui s’avère être plus important que le mécanisme de copie lui-même.

La partie que presque personne n’explique bien : la contrepression

C’est là le véritable problème que .pipe() vise à résoudre, au-delà d’une simple commodité. Imaginez un scénario où la lecture se fait rapidement, par exemple depuis un disque local, tandis que l’écriture est lente, peut-être via une connexion réseau à bande passante limitée.

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

Appeler .write() plus rapidement que le destinataire ne peut effectivement absorber les données n’amène pas le flux écrivable à générer d’erreur ni à bloquer l’exécution. Au contraire, il accumule silencieusement l’excédent dans un buffer mémoire interne, le conservant jusqu’à ce qu’il ait l’occasion de le vider. Si cette inadéquation persiste suffisamment longtemps, avec un écart important entre la vitesse de lecture et celle d’écriture, le problème exact que le streaming était censé éliminer se reconstitue discrètement : l’utilisation de la mémoire augmente sans limite, simplement différée plutôt que sur le champ, et il est généralement beaucoup plus difficile de la détecter avant qu’elle ne provoque une panne.

La contrepression constitue la protection contre ce mode de défaillance précis : elle permet à un flux écrivable d’indiquer qu’il a atteint sa capacité et a besoin que le producteur ralentisse, et un producteur correctement implémenté respecte ce signal au lieu de continuer à envoyer des données.

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
});

Dès que le buffer interne d’un flux écrivable dépasse sa limite configurée, la méthode .write() renvoie false, ce qui sert de signal pour arrêter la production jusqu’à ce que le flux émette un événement drain indiquant qu’il a éliminé les données en attente et est prêt à en accepter davantage. Ce schéma d’arrêt temporaire puis de reprise est précisément ce que la méthode .pipe() gère automatiquement en arrière-plan :

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

C’est cela, et non la brièveté, qui constitue la véritable justification pour privilégier .pipe() par rapport à l’association manuelle des écouteurs data et write. Envoyer manuellement des blocs de données en ignorant la valeur de retour de .write() réintroduit le même problème de mémoire illimitée que les flux étaient conçus pour éviter à l’origine ; il s’agit simplement d’une étape de moins par rapport à la erreur évidente inhérente à readFileSync.

Flux de transformation : traitement des données en temps réel

Il existe des cas où il ne suffit pas de déplacer les données telles quelles d’un endroit à un autre, il est également nécessaire de les remodeler en cours de route. Un flux Transform a été conçu précisément pour cela : il se situe au milieu d’une chaîne de pipelines, reçoit les blocs entrants, applique une opération à chacun d’eux, puis envoie le résultat à ce qui vient ensuite :

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"));

Chaque morceau est converti en majuscules au fur et à mesure qu’il traverse le pipeline, et à aucun moment le fichier complet, que ce soit sous sa forme originale ou convertie, n’a besoin d’être entièrement stocké en mémoire. C’est précisément le mécanisme à l’origine de la fonction intégrée de Node zlib.createGzip(). Il s’agit simplement d’un flux Transform qui compresse chaque morceau au fur et à mesure de son arrivée, et comme toute transformation, il peut être intégré directement dans une chaîne de pipelines de la même manière que dans l’exemple des majuscules :

const zlib = require("zlib");

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

La lecture, la compression et l’écriture ont lieu en même temps, en traitant de petits fragments de données au fur et à mesure de leur arrivée, sans que l’une des trois étapes n’ait besoin du fichier entier chargé simultanément.

Pourquoi il est utile de vraiment comprendre cela

Les flux ont la réputation d’être l’un des aspects les plus complexes de l’API Node, et honnêtement, leur interface brute basée sur les événements semble effectivement peu pratique, ce qui explique en partie cette réputation. Mais le concept sous-jacent est simple : éviter de charger l’ensemble des données en mémoire, les traiter par étapes plutôt que d’un coup, et s’assurer qu’un producteur rapide ne submerge pas silencieusement un consommateur lent. Une fois que vous avez internalisé ce modèle, .pipe(), Transform et le mécanisme de contrainte de débit ne semblent plus être trois API indépendantes à mémoriser séparément. Il s’agit en réalité d’une seule idée, présentée à travers trois éléments liés de l’API, permettant de résoudre un problème auquel des outils comme readFileSync n’ont jamais eu à faire, car ils n’étaient pas conçus pour gérer autre chose que de petits fichiers.

Lectures complémentaires