Transform-потоки в Node.js: обработка данных на лету
Содержание
О проверке кода. Примеры не запускались в этой сессии (по договорённости — фокус на контенте, а не прогонах). Говорим прямо, а не ставим бейдж «проверено».
Потоки — одна из сильнейших идей Node.js: обрабатывать данные по мере поступления, не загружая целиком. Transform-поток — тот, что данные ещё и преобразует на лету. Разберём его на практике.
Идея: обработка по кускам
Обычный подход «загрузить всё → обработать → сохранить» ломается на больших данных. Потоки работают иначе:
// файл проходит через конвейер по кускам, не загружаясь целиком
fs.createReadStream('input.txt') // читает кусками
.pipe(transform) // преобразует каждый кусок
.pipe(fs.createWriteStream('output.txt')); // пишет результат
Данные текут через конвейер фрагментами. Пришёл кусок — преобразовали — отдали дальше. В памяти в каждый момент лишь текущий фрагмент, а не весь объём. Это ключ к экономии памяти при работе с большими файлами.
- 1 Readable — источникФайл/сеть отдаёт данные кусками
- 2 Transform — преобразованиеКаждый кусок сжимается/парсится/меняется
- 3 Writable — приёмникРезультат пишется в файл/ответ
- 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.