Streams и backpressure
Любой сервис рано или поздно встречается с данными, которые не влезают в память целиком: выгрузка из базы на миллион строк, импорт CSV, проксирование файла, gzip-ответ на скачивание. Копировать такое через readFile + writeFile — значит держать весь буфер в памяти и вылететь по OOM на первом же файле больше свободного RAM. Стримы решают это: данные идут кусками (chunks), обрабатываются по мере поступления, а потребитель регулирует скорость производителя. Этот механизм регулировки называется backpressure — и он главный предмет этой главы.
Краткая версия показала pipe и упомянула highWaterMark. Здесь — полная картина: четыре класса стримов, два режима чтения, разница между pipe и pipeline, ручное управление backpressure и три рабочих примера, которые можно сразу взять в проект.
Четыре класса
Заголовок раздела «Четыре класса»Всё построено на базовых классах модуля node:stream:
| Класс | Роль | Пример из ядра |
|---|---|---|
Readable |
Источник данных | fs.createReadStream, HTTP-запрос (входящий) |
Writable |
Приёмник данных | fs.createWriteStream, HTTP-ответ |
Transform |
Читает → преобразует → отдаёт | zlib.createGzip, crypto.createCipheriv |
Duplex |
Независимые каналы чтения и записи | TCP-сокет |
Transform — наследник и Readable, и Writable одновременно: пишешь в него чанками, читаешь трансформированные чанки. Именно поэтому цепочка «файл → gzip → файл» собирается в одну строку.
Режимы чтения: flowing и paused
Заголовок раздела «Режимы чтения: flowing и paused»Readable имеет два состояния:
- Paused (по умолчанию). Данные не текут, пока ты их не запросишь: через
read(),pipe()илиfor await. - Flowing. Данные текут сами, собираются через события
data. Включается подпиской наdataили вызовомresume().
import { createReadStream } from 'node:fs';
// Flowing: получаешь каждый чанк событием, сколько бы их ни былоcreateReadStream('big.iso') .on('data', (chunk) => console.log('получил', chunk.length, 'байт')) .on('end', () => console.log('конец'));
// Paused: ручное управление, удобно для протоколов с заголовкамиconst stream = createReadStream('file.bin');stream.on('readable', () => { let chunk; while ((chunk = stream.read()) !== null) { // обработать chunk }});Главное отличие — контроль скорости. В flowing-режиме данные льются как придут, и если обработка внутри data медленная, чанки копятся в памяти (Node буферизует их внутри стрима). В paused-режиме ты сам решаешь, когда забрать следующий кусок. Про for await — ниже, он делает paused-режим удобным.
Backpressure: механика
Заголовок раздела «Backpressure: механика»У каждого Writable есть внутренний буфер и порог — highWaterMark. Правила простые:
writable.write(chunk)возвращаетboolean:true— «могу принять ещё»,false— «буфер заполнен до highWaterMark, замедляйся».false— это не ошибка и не отказ: чанк принят в буфер. Но продолжать писать, игнорируяfalse, — значит неограниченно раздувать память.- Когда буфер дренируется (реальный I/O завершён), стрим испускает событие
drain. Продолжать писать послеfalseнужно только послеdrain.
import { once } from 'node:events';
async function writeLots(writable, chunks) { for (const chunk of chunks) { const ok = writable.write(chunk); if (!ok) { // буфер полон — ждём, пока не уйдёт хотя бы highWaterMark данных await once(writable, 'drain'); } } writable.end(); await once(writable, 'finish');}highWaterMark задаётся при создании стрима. Для файлов умолчание 16 КБ — мелко для быстрых дисков; для прокси через сеть 16–64 КБ нормально. Не ставь «на всякий случай» 10 МБ: backpressure перестанет работать раньше, чем тебе нужно. Подробная механика с диаграммами — в официальном гайде Node.js «Backpressuring in Streams».
pipe против pipeline
Заголовок раздела «pipe против pipeline»// Вариант 1: ручная склейкаreadStream.pipe(gzip).pipe(writeStream);pipe возвращает последний стрим в цепочке, поэтому a.pipe(b).pipe(c) выглядит красиво. Проблема в обработке ошибок: событие error всплывает только на том стриме, где оно случилось. Если readStream упал, writeStream останется висеть открытым, сокеты не закроются — утечка дескрипторов.
// Вариант 2: pipeline — ошибки и очистка из коробкиimport { pipeline } from 'node:stream/promises';
await pipeline(readStream, gzip, writeStream);// Ошибка в любом звене → все стримы уничтожены, промис отклонён.// Успех → все корректно закрыты.pipeline (в промисной версии из node:stream/promises) — единственный правильный способ собирать цепочки в новом коде. Он уничтожает стримы (destroy()), корректно обрабатывает abort через AbortSignal, и ошибка в одном звене приводит к чистой остановке всех — это зафиксировано и в официальной документации stream.pipeline.
Единственный кейс, где pipe ещё жив: мгновенная склейка двух стримов с последующей ручной обработкой ошибок на каждом. Но честно — таких кейсов почти не осталось.
Пример 1: файловый копировщик с замером памяти
Заголовок раздела «Пример 1: файловый копировщик с замером памяти»import { createReadStream, createWriteStream } from 'node:fs';import { pipeline } from 'node:stream/promises';
async function copyFile(src, dest) { await pipeline( createReadStream(src, { highWaterMark: 64 * 1024 }), createWriteStream(dest), ); console.log('скопировано:', dest);}
const src = process.argv[2];const dest = process.argv[3];
// смотрим пик памяти процессаcopyFile(src, dest).then(() => { const used = process.memoryUsage().heapUsed / 1024 / 1024; console.log(`heapUsed после копирования: ${used.toFixed(1)} МБ`);});Сравни с await writeFile(dest, await readFile(src)) на файле в 2 ГБ: второй вариант выделит 2 ГБ под входной буфер и ещё около 2 ГБ под выходной, а стримовый — держит в памяти два куска по 64 КБ. Это разница между «работает на VPS с 1 ГБ RAM» и OOM-киллером.
Пример 2: gzip-трансформ на лету
Заголовок раздела «Пример 2: gzip-трансформ на лету»import { createReadStream, createWriteStream } from 'node:fs';import { createGzip } from 'node:zlib';import { pipeline } from 'node:stream/promises';
// сжимаем файл, не создавая промежуточных близнецов в памятиawait pipeline( createReadStream('access.log'), createGzip({ level: 6 }), createWriteStream('access.log.gz'),);createGzip — это Transform: читает несжатые чанки, отдаёт сжатые. Степень сжатия (level: 1..9) — классический размен скорость/размер. Для HTTP-ответов gzip-стрим ставится прямо в цепочку ответа — и заголовок Content-Encoding: gzip ставится до того, как известен итоговый размер, поэтому ответ идёт chunked.
Пример 3: CSV-парсер в объектном режиме
Заголовок раздела «Пример 3: CSV-парсер в объектном режиме»Вот где стримы раскрываются полностью: читаем CSV любого размера, отдаём наружу объекты по одному.
import { createReadStream } from 'node:fs';import { Transform } from 'node:stream';import { pipeline } from 'node:stream/promises';
// мини-парсер CSV: предполагаем простые строки без кавычек-переносовfunction parseCsv() { let buffer = ''; let headers = null;
return new Transform({ objectMode: true, // на выходе объекты, а не байты
transform(chunk, _enc, callback) { buffer += chunk.toString('utf8'); const lines = buffer.split('\n'); buffer = lines.pop(); // последняя строка может быть неполной — в буфер
for (const line of lines) { if (!line.trim()) continue; const cells = line.split(','); if (!headers) { headers = cells.map((h) => h.trim()); continue; } const row = Object.fromEntries(headers.map((h, i) => [h, cells[i]?.trim()])); this.push(row); // отдаём объект дальше по цепочке } callback(); // этот чанк обработан, можно слать следующий },
flush(callback) { if (buffer.trim()) { const cells = buffer.split(','); this.push(Object.fromEntries(headers.map((h, i) => [h, cells[i]?.trim()]))); } callback(); // конец стрима }, });}
// использование: читаем построчно, обрабатываем с backpressure из коробкиawait pipeline( createReadStream('users.csv', { encoding: 'utf8' }), parseCsv(), async function* (source) { for await (const row of source) { // тут может быть запись в БД — await реально тормозит чтение файла console.log('обрабатываю:', row.email); yield row; } },);Два ключевых момента. Во-первых, objectMode: true — стрим носит объекты, и this.push(row) не блокируется на заполнении байтового буфера. Во-вторых, финальное звено — асинхронный генератор: for await + yield превращает любую асинхронную обработку в Transform, и backpressure работает автоматически — файл читается ровно с той скоростью, с какой обработка успевает.
Типичные ошибки и грабли
Заголовок раздела «Типичные ошибки и грабли»- Цепочка из
pipeбез обработки ошибок. Ошибка в середине оставляет хвостовые стримы открытыми: файловые дескрипторы, сокеты, таймеры — утечка до перезапуска процесса. Лечение:pipelineвсегда. - Игнорирование возврата
write().writable.write(chunk)в цикле без проверкиfalseиdrain— классическая утечка памяти: буфер растёт быстрее, чем успевает I/O. Симптом — растущийheapUsedи RSS без видимых причин. data-обработчик, который думает, что успевает. В flowing-режиме событияdataидут так быстро, что тяжёлая обработка внутри обработчика копит чанки. Лечение:for awaitили paused-режим с явнымread().- Событие
endвместоfinish.end— у Readable («мне больше нечего читать»),finish— у Writable («я всё записал»). В коде копированияwriteStream.on('end')никогда не сработает, иawaitзавершения теряется. - Смешение объектного и байтового режимов. Подключить
objectMode: trueTransform после файлового Readable без перекодировки — и получишьERR_INVALID_ARG_TYPEили молчаливую порчу данных. Вpipelineмежду байтовым и объектным звеном всегда должно быть явное звено-конвертер. - Потеря хвоста при ручном парсинге. Строки, разорванные границей чанков, — баг №1 самописных парсеров. Правило: всегда оставляй «хвост» в буфере до следующего чанка и добирай его во
flush.
Вопросы на собеседовании
Заголовок раздела «Вопросы на собеседовании»- Что такое backpressure и зачем он нужен? Механизм, при котором медленный потребитель тормозит быстрого производителя. У Writable есть буфер и highWaterMark;
write()возвращаетfalse— производитель обязан дождатьсяdrain. Без этого память растёт безгранично. - Чем
pipelineлучше цепочкиpipe?pipelineпробрасывает ошибку из любого звена, уничтожает все стримы (destroy), корректно закрывает цепочку и возвращает промис.pipeне передаёт ошибки по цепочке — хвостовые стримы висят. - В чём разница flowing и paused режимов? Flowing: данные текут событиями
data, скорость не контролируется без дополнительного кода. Paused: данные забираются явно (read(),for await), потребитель сам регулирует темп. По умолчанию Readable — paused. - Что такое
highWaterMark? Порог внутреннего буфера стрима. Для байтовых стримов — в байтах (16 КБ по умолчанию), для объектных — в объектах (16). Достижение порога переводитwrite()в возвратfalseи событиеreadableу Readable. - Как правильно дождаться окончания записи в Writable?
writable.end()затемawait once(writable, 'finish')— илиawait pipeline(...), который делает это сам. Событиеfinishозначает: все буферизованные данные реально ушли на диск/в сеть. - Что такое объектный режим и где он нужен?
objectMode: trueпозволяет стриму носить произвольные JS-объекты вместо Buffer/строк. Нужен в ETL, парсерах CSV/JSON Lines, очередях обработки — везде, где единица данных — запись, а не байт. for awaitнад Readable — flowing или paused? Paused: следующий чанк запрашивается только после завершения тела итерации. Асинхронная работа внутри цикла автоматически создаёт backpressure на источник.
Практика
Заголовок раздела «Практика»- Прокси с замером. HTTP-сервер, который по
GET /file/:nameстримит файл из директории в ответ, с HTTP-таймаутом иpipeline. Нагрузи егоautocannonи смотриprocess.memoryUsage(): память должна оставаться плоской независимо от размера файла. Сравни с вариантомreadFile. - Генератор + gzip + счётчик. Напиши
Readable, который генерирует N случайных строк, склей черезpipelineсcreateGzipи Writable, считающим байты на выходе. Убедись, что приN = 10_000_000память не растёт. - JSON Lines парсер. Расширь CSV-парсер до формата JSONL (каждая строка — JSON-объект). Обработай случай строки, разорванной между чанками, и невалидного JSON: невалидные строки должны уходить в отдельный error-стрим или счётчик, не роняя пайплайн.
- Backpressure вручную. Напиши Writable, который «записывает» с задержкой 100 мс на чанк, и наполняй его из цикла с проверкой возврата
write()и ожиданиемdrain. Замерь размер внутреннего буфера (writable.writableLength) и убедись, что он не превышаетhighWaterMark + размер одного чанка. Потом убери проверку и смотри, как буфер раздувается. - Рефакторинг на
pipeline. Найди в своём pet-проекте место, где читается/пишется файл или проксируется ответ, и переведи его наpipelineс корректной обработкой ошибок. Добавь тест с обрывом соединения на середине передачи.
Что почитать
Заголовок раздела «Что почитать»- Документация Node.js: Stream — API всех четырёх классов и событий.
- Backpressuring in Streams — официальный гайд по backpressure.
- Stream Handbook (substack) — классика, чуть устарела по API, но не по идеям.
- Node.js Streams: everything you need to know — подробный обзор с диаграммами.
- API Compatibility / DON’T USE pipe — официальная рекомендация по pipeline.