Потоки, pipe и backpressure
Потоки гонят данные по chunk'ам, так что 4ГБ-файл никогда не оседает в RAM. Когда быстрый producer обгоняет медленный consumer, тормоз — write(), возвращающий false до drain. Подключай через pipeline(), а не ручной pipe.
Сервис год работал нормально. Потом клиент залил экспорт на 6ГБ, и под умер с OOMKilled, утянув ещё три тенанта на той же ноде. Хендлер выглядел невинно: прочитать загрузку, сжать gzip, записать в объектное хранилище — аккуратный for await-цикл, который читал chunk и звал upload.write(chunk). На медленном соединении с S3 каждый chunk, который не успевал сброситься, копился во внутреннем буфере writable. Цикл ни разу не паузился, потому что никто не проверял, что вернул write(). RSS вырос с 200МБ до 4ГБ за девяносто секунд. Код был не медленным — он был неограниченным.
Зачем нужны потоки: байты, которые ты никогда не держишь целиком
Что если твоя задача на 4ГБ падает в OOM не из-за бага в логике, а просто потому что ты загрузил файл до того, как прикоснулся к нему? Именно эту проблему решают потоки — и стоит один раз это увидеть, ты будешь тянуться за ними всякий раз, когда размер данных не ограничен.
Поток — это абстракция для данных, которые приходят или уходят во времени, обрабатываемых маленькими chunk’ами вместо материализации целиком в памяти. Выигрыш — ограниченная память: пропуская 4ГБ-файл через gzip на диск, ты держишь лишь несколько chunk’ов разом — килобайты, не гигабайты, — так что задача на 4ГБ и задача на 4ТБ занимают почти одинаковую память. Второй выигрыш — задержка до первого байта: потоковый HTTP-ответ может начать слать голову результата, пока его хвост ещё считается, вместо буферизации всего тела до того, как уйдёт первый байт.
Node моделирует четыре типа потоков. Readable — источник, из которого ты тянешь (чтение файла, тело HTTP-запроса, сокет). Writable — сток, в который ты пушишь (запись файла, HTTP-ответ, сокет). Duplex — оба сразу, с независимыми сторонами чтения и записи (TCP-сокет). Transform — это Duplex, чей выход есть функция от входа: gzip, шифр, парсер CSV, JSON.parse на строку. Transform — это как ты делаешь потоковый map: данные текут в одну сторону, изменённые данные текут из другой, и ты никогда не держишь весь датасет.
У Readable два режима чтения. В режиме paused ты явно зовёшь read(), чтобы вытянуть следующий chunk. В режиме flowing chunk’и пушатся в тебя через события 'data' так быстро, как приходят, — удобно, но если ты не успеваешь, встроенного тормоза нет, и именно так ты теряешь данные или взрываешь память. Навешивание обработчика 'data' или вызов pipe() переключает поток в flowing. Есть ещё object mode: вместо chunk’ов Buffer/string поток несёт произвольные JS-объекты (одна разобранная запись на chunk), что и делает Transform-конвейеры записей эргономичными.
Проблема: быстрый producer, медленный consumer
OOM из хука — каноничный провал потоков. Producer (чтение файла, быстрое — локальный диск) кормит consumer (загрузку в S3, медленную — сеть). Когда ты зовёшь writable.write(chunk), Node копирует chunk во внутренний буфер writable и пытается сбросить его в нижележащий ресурс. Если ресурс не успевает, буфер растёт. Ничто не мешает ему расти — write() всегда принимает chunk и возвращает управление. Если твой producer игнорирует возвращаемое значение и пишет дальше, буфер растёт без границы, пока процессу не кончится heap.
// БАГ: игнорирует backpressure. Буфер растёт без границы при медленном стоке.
import { createReadStream, createWriteStream } from "node:fs";
async function copyNaive(src, dst) {
const read = createReadStream(src);
const write = createWriteStream(dst);
for await (const chunk of read) {
write.write(chunk); // возвращаемое значение отброшено — нет тормоза
}
write.end();
}for await тянет chunk’и из быстрого источника так быстро, как может, и грузит их в writable, который может быть куда медленнее. Память теперь — это разница в скорости, умноженная на длительность, то есть неограниченна.
Backpressure: тормоз, встроенный в write()
Node даёт producer’у сигнал. writable.write(chunk) возвращает false, как только объём в буфере превышает highWaterMark потока (по умолчанию 16КБ для байтовых потоков, 16 объектов в object mode). false не значит, что запись провалилась — chunk всё равно принят, — это значит перестань слать; буфер достаточно полон. Корректный producer, увидев false, паузится и ждёт, пока writable не испустит событие 'drain', которое срабатывает, когда буфер сброшен ниже high-water mark. Потом он продолжает. Этот цикл — писать до false, ждать 'drain', повторять — и есть backpressure, и он ограничивает память примерно highWaterMark.
Ты почти никогда не пишешь этот цикл руками. readable.pipe(writable) подключает его за тебя: читает, пишет, а когда write() вернул false, зовёт readable.pause(), потом продолжает на 'drain'. Producer и consumer остаются связанными, а память — ограниченной, независимо от размера источника.
// КОРРЕКТНО: pipeline() подключает backpressure И проброс ошибок И очистку.
import { pipeline } from "node:stream/promises";
import { createReadStream, createWriteStream } from "node:fs";
import { createGzip } from "node:zlib";
async function gzipFile(src, dst) {
await pipeline(
createReadStream(src), // Readable (быстрый)
createGzip(), // Transform (потоковый map)
createWriteStream(dst), // Writable (медленный)
);
// резолвится только когда сброшен последний байт; реджектится, если упал ЛЮБОЙ этап
}▸Почему это работает
Почему pipeline(), а не .pipe().pipe()? Ручной pipe() пробрасывает данные и backpressure, но не ошибки, и он не чистит ресурсы. Если createGzip() бросит на полпути, readable никогда не уничтожается — его файловый дескриптор течёт, а ошибка всплывает как необработанное событие 'error', которое может уронить процесс. stream.pipeline() (и pipeline из node:stream/promises) добавили ровно для этого: он пробрасывает ошибки с любого этапа в callback/promise и уничтожает каждый поток в цепочке при провале или завершении, освобождая файловые дескрипторы и сокеты. Сеньорское правило прямое: используй pipeline(), никогда цепочку ручных .pipe(). Ручной pipe — это утечка ресурсов, ждущая своей первой ошибки.
Где это бьёт в проде
Везде, где быстрая сторона встречает медленную. Обратный прокси, гонящий ответ upstream’а клиенту на медленном мобильном канале: без backpressure прокси буферизует всё тело upstream’а в RAM на соединение, и пара тысяч медленных клиентов кладут машину в OOM. Большие загрузки файлов прямо в объектное хранилище, как в хуке. Конвейеры CSV/NDJSON, разбирающие миллионы строк: Transform в object mode держит память плоской независимо от числа строк, тогда как JSON.parse(fs.readFileSync(...)) загрузил бы весь файл. Server-sent events и хвост логов: клиент здесь — медленный consumer, и backpressure не даёт серверу копить мегабайты на подписчика.
Современный Node также говорит на Web Streams (ReadableStream/WritableStream, WHATWG-API, используемый fetch и браузером). У них та же семантика backpressure, выраженная иначе — ReadableStream не тянет из источника, пока consumer не попросит, — а stream.Readable.fromWeb() / .toWeb() мостят два мира, когда ты склеиваешь тела fetch() со стоками Node.
| Измерение | Загрузить всё в память (прочитать целиком, потом обработать) | Стримить по chunk’ам (с backpressure) |
|---|---|---|
| Пиковая память | Масштабируется с размером входа — 4ГБ-файлу нужно ~4ГБ+ RAM | Ограничена highWaterMark — несколько chunk’ов (~десятки КБ) независимо от размера |
| Задержка до первого байта | Высокая — ничего не выходит, пока не прочитан весь вход | Низкая — вывод стартует, как только преобразован первый chunk |
| Режим отказа | OOMKilled / смерть в GC на большом или враждебном входе | Плоский RSS; пропускная способность просто следует за самым медленным этапом |
writable.write(chunk) вернул false. Что это значит и что должен сделать корректный producer?
Надо скопировать Readable через gzip-Transform в Writable в долгоживущем сервере. Выбери способ подключения.
Когда ты сам пишешь Transform, переопредели _transform(chunk, encoding, callback), чтобы пушить преобразованный вывод и звать callback(), когда закончил с этим chunk’ом — фреймворк сам берёт на себя буферизацию и backpressure с обеих сторон, включая соблюдение highWaterMark твоего вывода. Поднятие highWaterMark меняет память на меньшее число drain-кругов (полезно для высокопропускных байтовых потоков); понижение — жёстче ограничивает память. Большинству кода стоит оставить значение по умолчанию.
- 01Пройди по механизму backpressure: что возвращает writable.write(), когда и что обязан сделать producer?
- 02Почему предпочесть pipeline() цепочке ручных .pipe()?
Поток обрабатывает данные маленькими chunk’ами во времени, а не материализует их целиком, что держит память ограниченной и снижает задержку до первого байта. В Node четыре вида — Readable (источник), Writable (сток), Duplex (оба) и Transform (потоковый map вроде gzip или парсера CSV), — а Readable читают в режиме flowing или paused, с object mode для конвейеров запись-за-записью. Опасность — быстрый producer, кормящий медленный consumer: write() всегда принимает chunk во внутренний буфер writable, поэтому producer, игнорирующий возвращаемое значение, растит этот буфер без границы и кладёт процесс в OOM — классический провал больших загрузок или прокси. Backpressure — встроенный тормоз: write() возвращает false, как только буфер проходит highWaterMark, и корректный producer паузится до события ‘drain’ перед отправкой следующего, ограничивая память около high-water mark. Этот цикл редко пишут руками — pipe() автоматизирует его, но сеньорский выбор — stream.pipeline(), потому что в отличие от ручного pipe он пробрасывает ошибки с любого этапа и уничтожает каждый поток при завершении или провале, так что ошибка на полпути не утечёт файловые дескрипторы и не уронит через необработанное ‘error’. Везде, где быстрая сторона встречает медленную — прокси, загрузки, CSV/NDJSON, SSE — стримь через pipeline() и дай backpressure держать оборону.
Практика
Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.