Головна / Статті / Розуміння потоків Node.js: проблему, яку вони насправді вирішують

Розуміння потоків Node.js: проблему, яку вони насправді вирішують

Дізнайтеся, чому існують потоки Node.js, як внутрішньо працює передача даних через трубопровід, і що насправді означає зворотний тиск для ефективної обробки великих об’ємів даних.

1261 слів

Більшість посібників з потоків починаються з огляду API — функцій .pipe(), Readable, Writable, Transform — ще до того, як розглянути саму проблему, для вирішення якої створені ці інструменти. Такий підхід змушує сприймати потоки як якусь абстрактну процедуру, яку потрібно запам’ятати. Якщо змінити порядок та почати з проблеми, API переважно пояснюється саме собою.

Проблема: деякі дані занадто великі, щоб зберігати їх одночасно

Уявіть завдання, яке полягає у перетворенні експорту CSV розміром 4 ГБ у формат JSON. Наївний підхід виглядає так:

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 не поверне керування програмі, поки весь 4-гігабайтний файл не буде завантажений у пам’ять як єдиний рядок JavaScript, і така рядкова форма зазвичай використовує значно більше пам’яті, ніж можна було б припустити за розміром самого файлу. На пристрої з лише 2 ГБ вільної оперативної пам’яті цей код не лише працює неефективно, але й може повністю зупинитися, або ж призводить до завмерання інших процесів, які борються за ту саму пам’ять. У самій логіці немає жодних проблем; усі кроки перетворення є правильними. Справжньою проблемою є прихована передумова всередині readFileSync: припущення, що можна одночасно зберігати весь файл у пам’яті, незалежно від його розміру.

Справжня ідея потоку

Замість того, щоб вимагати всю набір даних заздалегідь, стрім працює за іншою схемою: запитується один фрагмент, він обробляється, після чого запитується наступний фрагмент. У жоден момент повний файл не знаходиться в пам’яті. У будь-який момент існує лише невелика його частина, яка обробляється та звільняється ще до того, як з’явиться наступна частина.

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

Зверніть увагу, що загальний розмір файлу зовсім не враховується в цьому коді. Незалежно від того, чи є джерельний файл 4 ГБ чи 4 КБ, цей фрагмент коду працює однаково, хоча пік використання пам’яті суттєво відрізняється у цих двох випадках. Саме в цьому суть стрімування: замість необхідності мати повний набір даних з самого початку ми отримуємо можливість негайно розпочати обробку та завершити її, не зберігаючи в пам’яті більше, ніж обмежений проміжок даних.

Пайплайнування: безпосереднє підключення джерела до пункту призначення

Ручне отримання частин по одній є корисним підходом, але на практиці набагато частіше використовується підключення потоку для читання безпосередньо до потоку для запису, що дозволяє даним автоматично текти від джерела до пункту призначення, замість того щоб вручну передавати кожну частину:

const fs = require("fs");

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

Для копіювання файлу будь-якого розміру достатньо двох рядків без необхідності завантажувати весь файл у пам’ять. Функція .pipe() тут не виконує жодних магічних дій — вона просто направляє подію data, яка генерується джерелом, у виклик write у пункті призначення, а також обробляє ще одну деталь, яка виявляється важливішою за сам механізм копіювання.

Частину, яку майже ніхто не пояснює належним чином: тиск назад

Це саме та справжня проблема, яку вирішує .pipe(), а не просто питання зручності. Уявіть собі сценарій, коли читання відбувається швидко, наприклад з локального диска, тоді як запис відбувається повільно, можливо, через мережеве з’єднання з обмеженою пропускною здатністю.

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

Виклик .write() швидше, ніж це може обробити пункт призначення, не спричиняє появи помилки чи зупинки виконання потоку для запису. Натомість надлишкові дані тихо накопичуються у внутрішньому буфері пам’яті, де залишаються до моменту можливості їх видалення. Якщо така невідповідність триватиме достатньо довго, з достатньо великою різницею між швидкістю читання та запису, саме проблема, яку мало усунути стрімінг, знову постає: використання пам’яті зростає без меж, просто відкладаючись замість того, щоб виникнути негайно, і зазвичай значно складніше її виявити до того, як це призведе до збою.

Бекпресур — це захист саме від такого способу збою: він дає потоку, який можна записувати, змогу сигналізувати про те, що він досяг своєї межі та потребує, щоб виробник уповільнив роботу, а належним чином реалізований виробник враховує цей сигнал замість того, щоб продовжувати надсилати дані.

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

Як тільки внутрішній буфер потоку, який можна записувати, перевищує встановлений ліміт, функція .write() повертає false, що слугує сигналом зупинити надсилання даних, поки потік не відправить подію drain, яка покаже, що він очистив чергу та готовий приймати нові дані. Саме цю схему переривання та відновлення роботи автоматично обробляє функція .pipe() у фоновому режимі:

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

Саме це, а не стислість, є справжньою причиною використання .pipe() замість ручного підключення слухачів data та write. Передавання блоків даних вручну з ігноруванням значення, яке повертає .write(), знову створює проблему необмеженої пам’яті, якої спочатку і мали уникнути стріми; це просто ще один крок від очевидної помилки, притаманної readFileSync.

Трансформація стрімів: обробка даних під час передачі

Існують випадки, коли недостатньо просто перемістити дані з одного місця в інше без змін — їх також потрібно перетворювати по дорозі. Саме для цього створений стрім Transform: він знаходиться посеред ланцюга стрімів, приймає надходячі блоки даних, застосовує до кожного певну операцію та передає результат наступному етапу обробки:

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

Кожен фрагмент перетворюється на верхні літери під час проходження через потік обробки, і жоден момент не вимагає, щоб весь файл у своїй початковій чи перетвореній формі знаходився повністю в пам’яті. Саме цей механізм лежить в основі вбудованої функції Node zlib.createGzip(). Це не що інше, як потік Transform, який стискає кожен фрагмент по мірі його надходження, і, як і будь-яка трансформація, його можна безпосередньо включити до ланцюга потоків так само, як це було зроблено у прикладі з верхніми літерами:

const zlib = require("zlib");

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

Читання, стиснення та запис відбуваються одночасно, причому обробка здійснюється над невеликими фрагментами даних по мірі їх надходження, і жоден з цих трьох етапів не вимагає одночасного завантаження всього файлу.

Чому це варто справді зрозуміти

Стріми мають репутацію однієї з найскладніших частин інтерфейсу API Node, і, чесно кажучи, сирий, заснований на подіях інтерфейс дійсно здається незручним, тож ця репутація не є абсолютно безпідставною. Проте концепція, яка стоїть за ними, проста: уникати завантаження всього набору даних у пам’ять, обробляти його по частинах та переконуватися, що швидкий виробник ніколи не зможе тихо перевантажити повільного споживача під час цього. Як тільки ви усвідомите це як основну модель, .pipe(), Transform та механізм зворотного тиску більше не будуть здаватися трьома незалежними API, які потрібно запам’ятовувати окремо. Насправді це єдиная ідея, яка проявляється через три взаємопов’язані частини API, вирішуючи проблему, з якою ніколи не доводилося стикатися таким інструментам, як readFileSync, адже він спочатку не був призначений для обробки чогось, крім невеликих файлів.

Пов’язана література

  • Node.js Streams Beyond the Basics: Memory, Backpressure, and Real Failures — Дізнайтеся, як стріми Node.js взаємодіють із Web Streams, яку економію пам’яті можна досягти за допомогою реальних тестів, а також про помилки в продакшені, які з’являються лише під навантаженням.
  • Node.js Streams Explained: Fixing Out-of-Memory File Crashes — Дізнайтеся, чому завантаження цілих файлів у пам’ять призводить до збоїв серверів Node.js, та як стріми типу readable, writable, duplex та transform допомагають вирішити цю проблему за допомогою механізму backpressure.