Главная / Статьи / Понимание потоков 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, и как потоки читаемого, записываемого, двунаправленного типа и потоки преобразования решают эту проблему с помощью механизма обратного давления.