Transform-потоки в Node.js: обработка данных на лету

Содержание

О проверке кода. Примеры не запускались в этой сессии (по договорённости — фокус на контенте, а не прогонах). Говорим прямо, а не ставим бейдж «проверено».

Потоки — одна из сильнейших идей Node.js: обрабатывать данные по мере поступления, не загружая целиком. Transform-поток — тот, что данные ещё и преобразует на лету. Разберём его на практике.

Идея: обработка по кускам

Обычный подход «загрузить всё → обработать → сохранить» ломается на больших данных. Потоки работают иначе:

// файл проходит через конвейер по кускам, не загружаясь целиком
fs.createReadStream('input.txt')   // читает кусками
  .pipe(transform)                 // преобразует каждый кусок
  .pipe(fs.createWriteStream('output.txt'));  // пишет результат

Данные текут через конвейер фрагментами. Пришёл кусок — преобразовали — отдали дальше. В памяти в каждый момент лишь текущий фрагмент, а не весь объём. Это ключ к экономии памяти при работе с большими файлами.

Конвейер потоков: данные текут по кускам
  1. 1 Readable — источникФайл/сеть отдаёт данные кусками
  2. 2 Transform — преобразованиеКаждый кусок сжимается/парсится/меняется
  3. 3 Writable — приёмникРезультат пишется в файл/ответ
  4. 4 В памяти — один фрагментНе весь файл, а текущий кусок

Готовые Transform-потоки

Многое уже встроено — не надо писать руками:

const zlib = require('node:zlib');
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');

// сжать файл gzip — на лету, без загрузки в память
await pipeline(
  fs.createReadStream('big.log'),
  zlib.createGzip(),                 // готовый Transform: сжатие
  fs.createWriteStream('big.log.gz')
);

zlib.createGzip() — Transform-поток сжатия. Файл читается кусками, каждый сжимается и пишется — гигабайтный лог сжимается на скромной памяти. Готовые Transform есть для сжатия, шифрования, хеширования.

pipeline вместо pipe

// ❌ pipe — при ошибке потоки остаются висеть (утечка)
source.pipe(transform).pipe(dest);

// ✅ pipeline — ловит ошибки, закрывает всю цепочку
const { pipeline } = require('node:stream/promises');
try {
  await pipeline(source, transform, dest);
  console.log('готово');
} catch (err) {
  console.error('ошибка в конвейере:', err);   // одно место для всех ошибок
}

pipeline — правильный способ соединять потоки. Старый pipe при ошибке в одном из потоков оставляет остальные открытыми — утечка ресурсов и файловых дескрипторов. pipeline (особенно промис-версия из stream/promises) при сбое корректно закрывает всю цепочку и даёт ошибку в одном месте — обычный try/catch. В новом коде используют только его.

Свой Transform-поток

Когда готового нет — пишут свой:

const { Transform } = require('node:stream');

// перевести текст в верхний регистр на лету
const upperCase = new Transform({
  transform(chunk, encoding, callback) {
    const result = chunk.toString().toUpperCase();
    callback(null, result);   // (ошибка, результат) — результат уходит дальше
  },
});

await pipeline(
  fs.createReadStream('input.txt'),
  upperCase,
  fs.createWriteStream('output.txt')
);

Метод transform(chunk, encoding, callback) вызывается для каждого куска. Внутри преобразуете фрагмент и отдаёте результат через callback(null, result) (первый аргумент — ошибка, как в колбэках Node). Каждый пришедший фрагмент проходит через этот метод, результат уходит в следующий поток. Так пишут парсеры CSV, фильтры логов, конвертеры форматов — построчно и экономно.

Пример посерьёзнее — фильтр строк:

const filterErrors = new Transform({
  transform(chunk, enc, cb) {
    const lines = chunk.toString().split('\n')
      .filter(line => line.includes('ERROR'));
    cb(null, lines.join('\n') + '\n');
  },
});

Противодавление (backpressure)

Что если источник читает быстрее, чем приёмник успевает писать? Без контроля память переполнится накопленными данными.

Потоки решают это автоматически через противодавление: когда получатель не успевает, он сигналит источнику притормозить, и тот ждёт. pipeline и pipe управляют этим сами — вам не нужно думать о ручной синхронизации скоростей. Это одно из главных преимуществ потоков перед ручной обработкой кусками: защита от переполнения памяти встроена.

Async-итерация как альтернатива

Для чтения потока часто удобнее for await, чем события:

// перебрать поток построчно через async-итерацию
for await (const chunk of fs.createReadStream('file.txt')) {
  process(chunk);   // обрабатываем каждый кусок
}

Потоки — асинхронные итераторы, поэтому перебираются for await. Для построчного чтения текста поверх этого есть readline. Это современный, читаемый способ потреблять поток без ручной подписки на события data и end.

Когда это нужно

  • большие файлы — логи, дампы, CSV, которые не влезут в память;
  • преобразование на лету — сжатие, шифрование, конвертация форматов;
  • проксирование данных — читать из одного источника, писать в другой без буфера;
  • потоковые ответы сервера — отдавать клиенту данные по мере готовности.

Правило: как только объём данных непредсказуем или велик — думайте потоками, а не «загрузить всё». Это разница между приложением, которое падает на большом файле, и тем, что обрабатывает любой размер.


Смежные темы: управление памятью — зачем потоки; zlib и сжатие — готовый Transform; readline — построчное чтение; итераторы — for await. Полный список — в уроках Node.js.

Частые вопросы

Что такое Transform-поток?
Поток, который одновременно читает входные данные и пишет преобразованные, обрабатывая их по частям на лету. Пример — сжатие, шифрование, парсинг: кусок пришёл, преобразовали, отдали дальше, не держа весь объём в памяти.
Чем pipeline лучше pipe?
pipeline корректно обрабатывает ошибки и закрывает все потоки в цепочке при сбое, тогда как цепочка pipe оставляет висящие потоки и утечки при ошибке. В современном коде для соединения потоков берут pipeline из модуля stream/promises.
Зачем нужны потоки, если есть readFile?
readFile загружает файл в память целиком, что невозможно для файлов больше доступной памяти. Потоки обрабатывают данные по кускам, держа в памяти лишь текущий фрагмент, поэтому справляются с файлами и данными любого размера.
Как написать свой Transform-поток?
Создать Transform с методом transform(chunk, encoding, callback), в котором преобразуют кусок и передают результат через callback или push. Каждый пришедший фрагмент проходит через этот метод, а результат уходит в следующий поток цепочки.
Что такое противодавление в потоках?
Механизм, при котором быстрый источник притормаживает, если медленный получатель не успевает обрабатывать данные. Потоки и pipeline управляют этим автоматически, не давая переполнить память, если чтение быстрее записи.