Розумеўце стрімамы Node.js: проблемы, якія яны насправдзе рашаюць
Дазнайце, чаму існуюць стрымы Node.js, як унутранаўна працюе перадача дадзенняя, і што на самай працэ павертнага тыску значыць для эфектыўнай обработкі вялікіх заўважэнняй дадзення.
Большасць нарадчыкаў па стрімама пачынаюцца з абгаворэння API, такіх як .pipe(), Readable, Writable, Transform, прычаму нават не трапляючыся з самай проблемай, якую павінны рашаць гэтыя інструменты. Такі падход заставляе стрімы здавацца простаю церемоніяй, якую трэба запам’ятаваць. Адваротнае падчыненне — пачаць з проблемы, і тады API большавяючай часткай сама сабе поясніцца.
Проблема: Дзеяныя даны занадта великія, ўсё яно не можна зберагчы за раз
Уявіце задачу, калі трэба перакануць экспорт у формате CSV размерам 4 ГБ у формат JSON. Наўмовіты падход выглядае так:
const fs = require("fs");
const data = fs.readFileSync("export.csv", "utf8");
const rows = data.split("\n").map(parseRow);
fs.writeFileSync("export.json", JSON.stringify(rows));
readFileSync не вярне кантроль програме пакуль весь 4 ГБ файл не будзе заванесены ў памяць як адна JavaScript-строка, і такая страковая форма зазвычай спрацоўвае значна больш памяці, чым можна было б спадзявацца на адважанні самага розмеру файла. На апаратзе з толькі 2 ГБ вялікай памяці такі код не толькі працуе неэфектыва, але і выклікае крах, або ж ускладнюе роботу іншых процэсаў, якія борацца за тую ж память. У самай логіцы няма ніч глыбокага; усі крокі трансфармацыі правильныя. Асэнсіяйная проблема — це захаваная падазрачнасць унутры readFileSync: ўяўленне, што можна адночасна трываць цэлы файл у памяці, незалежна ад таго, насколькі вялікі ён.
Асэнсіяйная ідея застосоўвання стрыму
У протыяжэнне даных работаець на іншай падставе, чым выкарыстоўванне всіх дадзеных з самага пачатку: запрашваецца адна частка, яя обрабоўвается, пасля чаго запрашваецца наступная. Ніколі весь файл не знаходзіцца ў памяці. У кожны момент існуе толькі невялікая частка, яя обрабоўвается і зваліваецца прыходзячы час наступнай.
const fs = require("fs");
const readStream = fs.createReadStream("export.csv", { encoding: "utf8" });
readStream.on("data", (chunk) => {
console.log(`Received ${chunk.length} characters`);
});
readStream.on("end", () => {
console.log("Done reading the whole file, piece by piece");
});
Заўважыце, што загальны размер файла ніколі не бераецца ў расчыт у гэтым кодзе. Незалежна ад таго, чы розмер выхіднага файла 4 ГБ чы 4 КБ, гэты фрагмент коду працюе аднакова, хоць пікавы расхлоп памяці сутэльна разніцяеся ў гэтых двух варыянтах. У гэму і заключаецца сенс пратокавання даных: замест трабоўкі да всіх дадзеных з самага пачатку вы отрымваеце можлівасць негараздо шыбка пачаць роботу і завершыць яе, не зберагаючы ў памяці больш за вузкі прымек дадзеных.
Пайпінг: паўязка выхаднага джерела безпосередна з пунктом прызначэння
Ручна перадача фрагментоў по аднаму ёсць корыстным спосабам, але на практыцы значна частае ўзьязджанне чытальнага стрэму безпосередна з чытальнага у запісны дозволяе даным тэчыць ад выхаднага месца да пункту призначэння автаматычна, замест таго каб кожны фрагмент перадаваць вручную:
const fs = require("fs");
fs.createReadStream("export.csv")
.pipe(fs.createWriteStream("export-copy.csv"));
Для копіювання файла будзь-якага розмеру патрабуюцца толькі две лініі, без неабяжнасці завантажваць увесь файл у память. Функцыя .pipe() тут не выканае нічога чырага — яна проста перадае змест data, выданы ў выхадным об’екте, у функцію write у пункті призначэння, а таксама керуе ўсьмі додатковыми аспектамі, якія виявляюцца важлівымі больш, чым сам механізм копіювання.
Частка, якую майже ніхто не пояснюе належна: тырчаўка
Гэта ўсё тая працоўная проблема, якую рашае .pipe(), і яка выходзіць за межы простай зручнасці. Уявіце сцэнарый, калі чытанне вядзецца быстра, напрыклад з локальнага дыска, а запіс — медленна, можа, через сецьярнавы з’ѐеднанне з обмежанай пропускной спроможнасцю.
readStream.on("data", (chunk) => {
writeStream.write(chunk); // what happens if this can't keep up?
});
Выклік .write() быстрэй, чым пункт прызначэння можа фактычна перабрацаваць даны, не вызывае таго, што стрім для запісу выдаў бы памылку чы прызначыў бы виконанне. У замест на гэта ён тыха скупляе надзейшыя даны ў внутранім буфере памяці, зберагаючы іх, пакуль не будзе можлівасці ўсунуць іх. Якщо такая неадэкватнасць трывае дастатньго часу, калі розклік між шыбкасцю чытання і запісу є вялікі, тады сама проблема, яку мала рашыць стрімаванне, знову тыха стварыцца: викорыстоўвання памяці расте без меж, проста з адкладаннем, а не негайна, і яе зазвычай значна важка выявіць прытаму, пакуль яна не спрычыніць збою.
Backpressure — гэта захад працы протыпа такога спосабу абярання: ён дае стрэму, які можа быць запісаны, магчымасць сигналізаваць пра тое, што ён досяг свайго ліміту і патрэбуе, каб продюсер зменшыў інтэнсівнасць запісу, а правільна реалізаваны продюсер уважае гэты сигнал і не працуе далей, незважаючы на яго.
readStream.on("data", (chunk) => {
const canContinue = writeStream.write(chunk);
if (!canContinue) {
readStream.pause(); // stop reading until the writable side catches up
}
});
writeStream.on("drain", () => {
readStream.resume(); // writable side is ready for more
});
Як толькі внутрашняй буфер стрэму, які можа быць запісаны, перакроюе свой настаўлены ліміт, метод .write() вяртае false, што служыць сігналам аб стаўканні процесу запісу, пакуль стрэм не выдае змяну drain, якая парадкавае пра тое, што ён расчыстіў застойны матэрыял і гатовы прымаць больш. Гэты патерн паўз і продажчыкавання якраз тое, што метод .pipe() автаматычна керуе на заднім плане:
readStream.pipe(writeStream); // handles backpressure for you, silently, correctly
Гэта, а не краткасць, ўсё-такі являе сабою практычную прычыну для выбору .pipe() замест адмашовага падключэння слухачаў data і write. Перадача фрагментаў вручную, ігнаруючы значэнне, якое вяртае .write(), зноў стварае тую ж проблему неактуальнага выкарыстоўвання памяці, якой спачатку і было задумана західвацца праз стрімы; гэта проста ўскладнення явныяй памылкі, прыхованай у readFileSync.
Трансфармацыя стрімоў: обробка дадзеных у руху
Є ситуацыі, калі простая перадача дадзеных з аднаго месца ў іншае недастатня — трэба таксама зменіць ўжо ў паходзе іх форму. Стрым Transform створаны самэ для гэтага: ён знаходзіцца сярод ланцуга стрімоў, прыймае надходзячыя фрагменты, застосоўвае да кожнага певную операцыю і перадае рэзультат таму, што знаходзіцца далей:
const { Transform } = require("stream");
const upperCaseTransform = new Transform({
transform(chunk, encoding, callback) {
callback(null, chunk.toString().toUpperCase());
},
});
fs.createReadStream("input.txt")
.pipe(upperCaseTransform)
.pipe(fs.createWriteStream("output.txt"));
Кожны частак ператвараецца у верхнія літеры пад час працы ў ланцюгу перадачы, і ніколі не патрабуецца, каб цэлы файл, у яго первісным чы рэкордаваным варыянце, застаўся ў всій меры ў памяці. Гэта самэ прынцып, які стоіць за вбудованым у Node функцыям zlib.createGzip(). Це проста поток Transform, які стискае кожны частак, як толькі той прыходзіць, і, як і будзь-яка іншая трансформацыя, яго можна без проблем дадаць у ланцюг пайпаў так сама, як і прыклад з верхнімі літерамі:
const zlib = require("zlib");
fs.createReadStream("export.csv")
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream("export.csv.gz"));
Чытанне, стысненне і запіс вадзейцца адразу, працуючы з малымі фрагментамі дадзейнаў, якія прыходзяць, пры чым ні адна з трох паэтапаў не патрабуецца, каб весь файл быў завантажаны адразу.
Чаму це варта розумець
Стрімы здобулі репутацыю аднаго з найболей складных аспектаў API Node, і, чыста правду кажучы, ўсунутая, абмовленая на змаганнях інтерфейс дзяйсна выглядае незручным, таму такая репутацыя не ўсё ж такі безпадзейна. Але ідея, якая стоіць за ўсім гэтым, простая: утримвацца ад завантажэння цэлага набора дадзеных у памяць, працаваць з ямі по частках і старацца, каб швыдкі виробнік ніколі не могаў таямна засыпаць повольнага спожывача пад час гэтага. Калі вы усвядоміце гэта як асалодны модель, .pipe(), Transform і прытиск зворачнай сілы больш не будуць здавацца трыма незв’язанымі API, якія трэба запамятаваць окрема. Насправды гэта адна ідея, якая выражаецца чераз тры часткі API, якія ўзаецься межы сабою і рашаюць проблему, з якою ніколі не давалося стаць такім функцыям, як readFileSync, таму што ён з самага пачатку не быў прызначаны для обробкі чаго-небудзь большага, чым маленькія файлы.
Спадневаная літэратура
- Node.js Streams Beyond the Basics: Memory, Backpressure, and Real Failures — Дазвольце дазнацца, як струмені Node.js взаімаюцься з Web Streams, які ўскладнення ў викорыстоўванні памяці можна зафіксаваць за дапамогою рэальных тестаў, а таксама пра проблэмы, якія выявляюцца толькі пад вялікай навантажэннем.
- Node.js Streams Explained: Fixing Out-of-Memory File Crashes — Дазвольце дазнацца, чаму завантажэнне цэлых файлоў у памяць спрычынае зупінку сервераў Node.js, і як струмені чытання, запісу, двунаправленага абмэну і трансфармацыі дапамагаюць рашыць гэту проблэму за дапамогою механізма backpressure.