Конвейеры потоков, трансформы и ловушка .pipe()
.pipe() не пробрасывает ошибки и не разрушает цепочку, так что один сбой в середине течёт fd-шками до EMFILE. stream.pipeline() проводит ошибки и очистку; пиши Transform через _transform и мости Node- и Web-потоки на стыках.
Сервис ingest падал каждый вторник примерно в одно и то же время, всегда с EMFILE: too many open files, и никогда не воспроизводился в стейджинге. Открытые файловые дескрипторы, нарисованные за неделю, поднимались медленной лесенкой — по несколько сотен в день — пока процесс не упирался в ulimit 4096, и accept() не начинал отказывать. Горячий путь был в три строки: req.pipe(gunzip).pipe(parser). Большинство запросов были в порядке. Но часть присылала битое gzip-тело, gunzip бросал ошибку, и поскольку .pipe() пробрасывает данные, но не ошибки, входящий сокет так и не разрушался. Каждый плохой запрос утекал одним fd. Код был верен на счастливом пути — несчастливого в нём просто не было.
Ловушка .pipe(): утечка одного fd на каждый сбой
Урок 02 оставил тебе правило «используй pipeline(), а не ручной pipe()». Вот почему — по механизму. a.pipe(b).pipe(c) верно проводит поток данных и backpressure — эта часть работает, — но по умолчанию делает две опасные вещи. Первое: он не пробрасывает ошибки. Если b генерит 'error', эта ошибка остаётся на b; она не распространяется ни на a, ни на c. Второе: он не разрушает остаток цепочки при сбое. Когда b падает, a остаётся открытым и текущим, а c никто не велит закрыться. Все ресурсы, что держат a и c, — файловый дескриптор, сокет, буфер в полёте — утекают.
// THE TRAP: a malformed gzip body throws in gunzip; the socket leaks.
import zlib from "node:zlib";
function handler(req, res) {
const gunzip = zlib.createGunzip();
const parser = createNdjsonParser();
req.pipe(gunzip).pipe(parser); // ❌ no error wiring, no cleanup
parser.on("data", (record) => save(record));
parser.on("end", () => res.end("ok"));
// gunzip throws on bad input → 'error' on gunzip only.
// req (the socket) is never destroyed. One leaked fd per bad request.
}В одиночном запуске CLI это невидимо: процесс завершается и ОС забирает fd обратно. На долгоживущем сервере это смертельно. Каждый утёкший дескриптор постоянен на всё время жизни процесса, и они накапливаются по одному на сбой. Стандартный ulimit -n часто 1024, поднятый до 4096 на нагруженных машинах; как только ты утёк столько, каждый новый сокет, открытие файла и accept() отказывают с EMFILE, и сервис валится — ровно та вторничная лесенка. Хуже того, 'error' на потоке без слушателя сам по себе становится uncaughtException, так что второй режим сбоя голого pipe — жёсткое падение. У ловушки два выхода, и оба плохие.
Фикс: stream.pipeline()
stream.pipeline(...streams, callback) добавили, чтобы закрыть именно эту дыру. Он проводит тот же поток данных и backpressure, что и pipe(), и дополнительно: пробрасывает ошибку из любого этапа в один callback и разрушает каждый поток в цепочке — при сбое и при нормальном завершении. Одно место узнаёт исход; ничто не утекает. Промис-форма из node:stream/promises даёт те же гарантии за await — это та форма, что нужна в современном async-коде.
// THE FIX: pipeline wires errors AND destroys the whole chain on any outcome.
import { pipeline } from "node:stream/promises";
import zlib from "node:zlib";
async function handler(req, res) {
try {
await pipeline(
req, // Readable: the inbound socket
zlib.createGunzip(), // Transform: throws on a malformed body
createNdjsonParser(), // Transform: objectMode records out
saveSink(), // Writable: persists each record
);
res.end("ok");
} catch (err) {
// gunzip's error lands HERE. req, gunzip, parser, sink are ALL destroyed.
res.statusCode = 400;
res.end("bad body");
}
}Когда createGunzip() падает на битом теле, pipeline реджектится, каждый поток — включая входящий сокет req — разрушается, и fd освобождается. Никакой лесенки, никакого uncaughtException. Callback-форма, pipeline(req, gunzip, parser, sink, (err) => …), идентична по гарантиям; используй её, когда ты не в async-функции. Для одного потока, а не цепочки, stream.finished(stream, cb) (или его промис-форма) сообщит, когда поток завершён или упал, без переписывания учёта слушателей вручную. Правило заостряется: всякий раз, когда задействован больше чем один поток, первым делом тянись за pipeline() и никогда не печатай .pipe( в серверном коде.
Авторинг Transform: _transform, _flush, object mode
Когда смотришь на дата-конвейер и ни один встроенный не подходит — нужно парсить строки, срезать поля или на лету перекодировать записи — писать свой Transform правильное решение. Правил немного, но нарушить любое из них значит получить сбои, которые сложно воспроизвести.
Transform — это Duplex, у которого читаемая сторона есть функция от записываемой: gzip, шифр, построчный парсер. Ты пишешь его, реализуя _transform(chunk, encoding, callback): сделай работу, отправь ноль или больше выходов через this.push(...), затем вызови callback(), сигнализируя, что с этим чанком ты закончил и готов к следующему. Опциональный _flush(callback) бежит один раз, когда записываемая сторона заканчивается, — это твой шанс отдать трейлер или то, что осталось во внутреннем накопителе. Ставь objectMode: true, когда чанки — это записи (распарсенные строки), а не байты; тогда this.push(obj) отдаёт JS-объект вместо Buffer.
import { Transform } from "node:stream";
// Поток байт на входе → один объект на строку NDJSON на выходе. objectMode на читаемой стороне.
class NdjsonParser extends Transform {
constructor() {
super({ readableObjectMode: true }); // bytes in, objects out
this.tail = "";
}
_transform(chunk, _enc, cb) {
const lines = (this.tail + chunk).split("\n");
this.tail = lines.pop(); // last fragment may be a partial line
try {
for (const line of lines) {
if (line) this.push(JSON.parse(line)); // emit one record per line
}
cb(); // success: ready for next chunk
} catch (err) {
cb(err); // failure: propagate, do NOT throw + cb
}
}
_flush(cb) {
if (this.tail) this.push(JSON.parse(this.tail)); // last line, no newline
cb();
}
}У Transform два highWaterMark — по одному на сторону — каждый со значением по умолчанию 16384 байта (16КБ) в байтовом режиме или 16 объектов в object mode, и он участвует в backpressure на обеих: он перестаёт тянуть сверху, когда его читаемая сторона полна, и сигналит верху притормозить. Тут кусаются три режима сбоя. Сигналь об ошибках через cb(err), а не throw синхронно внутри async-продолжения — throw там ускользает из машинерии потока и становится uncaughtException, ровно как ловушка EventEmitter. Никогда не вызывай cb дважды: двойной callback — это жёсткий ERR_MULTIPLE_CALLBACK, портящий внутреннее состояние потока. И кардинальный грех — Transform, который буферизует всё в _flush (накопить все чанки, отдать в конце), молча убивает стриминг: он держит весь датасет в памяти, так что вход 4ГБ, который должен был стоить ~3×16КБ на трёхэтапном конвейере, вместо этого стоит 4ГБ и роняет процесс по OOM. Весь смысл — пушить по ходу.
▸Почему это работает
Современный Node несёт два мира потоков. Рядом с собственными потоками Node есть Web Streams — WHATWG ReadableStream / WritableStream / TransformStream — тот же API, что у браузера и edge-рантаймов. Они важны, потому что fetch() отдаёт тебе Response.body как Web ReadableStream, а не Node Readable, так что везде, где ты делаешь fetch() и хочешь спайпить тело в Node-приёмник, нужно мостить. stream.Readable.fromWeb(webReadable) и nodeReadable.toWeb() адаптируют между ними: Readable.fromWeb(res.body) превращает тело fetch в Node Readable, который можно прямо засунуть в pipeline(). Компромисс реален. Web Streams переносимы — идентичный код бежит в браузере, Deno, Cloudflare Workers и Node, — но они медленнее и тоньше: меньше встроенного, нет эргономики objectMode, и адаптер добавляет копию/обёртку на чанк порядка нескольких микросекунд каждая, что складывается на горячем пути с мелкими чанками. Node-потоки быстрее и проверены боем, но только под Node. Эвристика: используй Web Streams на edge и на границах fetch ради переносимости; конвертируй в Node-потоки для тяжёлой трансформации внутри процесса и держи конвертацию на стыке, а не во внутреннем цикле.
Долгоживущий HTTP-сервер пайпит входящий запрос через gunzip-Transform и NDJSON-парсер в приёмник. Часть запросов несёт битое gzip-тело. Как ты соединяешь цепочку?
В req.pipe(gunzip).pipe(parser) на сервере gunzip падает на битом теле. Что станет с файловым дескриптором req и каков итоговый системный сбой?
- 01Почему req.pipe(gunzip).pipe(parser) течёт файловым дескриптором, когда gunzip падает, и как именно stream.pipeline() это предотвращает?
- 02Каковы правила авторинга _transform/_flush у Transform и какой режим сбоя по памяти надо избегать?
Урок 02 сказал «используй pipeline(), а не pipe()»; вот механизм и цена. a.pipe(b).pipe(c) пробрасывает данные и backpressure, но по умолчанию не пробрасывает ошибки и не разрушает цепочку при сбое — так что когда средний этап вроде gunzip падает на битом теле, верхний сокет остаётся открытым и его файловый дескриптор течёт, по одному на плохой запрос, лезя лесенкой, пока процесс не упрётся в ulimit (~1024–4096) и не умрёт с EMFILE; та же ошибка голого pipe без слушателя — ещё и uncaughtException. stream.pipeline(...streams, cb) — и pipeline из node:stream/promises ради await — закрывает обе дыры: он пробрасывает ошибку из любого этапа в один callback/await и разрушает каждый поток при сбое и при завершении, освобождая все ресурсы; stream.finished() делает то же для одного потока. Ты пишешь Transform, реализуя _transform(chunk, enc, cb) (пушь выходы, затем cb() или cb(err)) и опциональный _flush(cb), используя objectMode для потоков записей; у Transform два highWaterMark (16КБ / 16 объектов по умолчанию) и backpressure на обеих сторонах, так что никогда не делай throw вместо cb(err), никогда не вызывай cb дважды и никогда не буферизуй весь поток в _flush, иначе ты меняешь ~3×16КБ на весь файл и роняешь по OOM. Наконец, Node теперь несёт WHATWG Web Streams (ReadableStream/WritableStream, WHATWG API) рядом со своими; Readable.fromWeb() / .toWeb() мостят их для тел fetch() и edge-рантаймов — Web Streams переносимы, но медленнее и тоньше, так что конвертируй на стыке и держи Node-потоки в горячем внутреннем цикле. Теперь, когда встречаешь .pipe().pipe() в серверном коде, риск ясен сразу: ищи слушатель 'error' на каждом этапе — если хоть один пропущен, один битый запрос запустит медленную лесенку fd к EMFILE.
Практика
Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.