Пошук уроків, статей та іншого контенту
Створіть потоковий Node.js-скрипт для читання, трансформації та запису великого файлу.
Потрібно створити Node.js-скрипт, який:
читає великий текстовий файл потоково;
обробляє його рядок за рядком;
видаляє порожні рядки;
маскує email-адреси;
прибирає пробіли в кінці рядків;
записує результат в інший файл.
Файл може бути настільки великим, що його не можна безпечно завантажити в пам’ять повністю.
readFileТакий підхід завантажує весь файл у пам’ять:
const fs = require('node:fs/promises');
const content = await fs.readFile('input.log', 'utf8');
const transformed = content.toUpperCase();
await fs.writeFile('output.log', transformed);Для невеликих файлів це нормально, але для файлу на кілька гігабайтів процес може:
використати всю доступну пам’ять;
почати працювати дуже повільно через часті операції зі збирання сміття;
завершитися з помилкою JavaScript heap out of memory.
Потокова обробка працює інакше: файл читається невеликими порціями, кожна порція одразу обробляється та передається на запис.
Для цього використаємо три потоки:
ReadStream → Transform → WriteStreamReadStream читає файл частинами;
Transform змінює дані;
WriteStream записує результат.
Важливо, що потоки Node.js працюють із частинами даних, які називаються chunks. Chunk не обов’язково відповідає одному рядку. Один рядок може:
повністю міститися в одному chunk;
бути розділеним між двома chunk;
містити кілька рядків в одному chunk.
Тому не можна просто викликати chunk.toString() і вважати отриманий текст набором повних рядків.
Для коректної роботи потрібно:
декодувати байти в текст;
додати до поточного chunk незавершений рядок із попереднього chunk;
знайти всі повні рядки;
зберегти останній незавершений фрагмент для наступного chunk.
Для декодування використаємо StringDecoder. Він правильно обробляє UTF-8-символи, які також можуть бути розділені між chunks.
Створіть файл process-log.js:
const fs = require('node:fs');
const { Transform, pipeline } = require('node:stream');
const { promisify } = require('node:util');
const { StringDecoder } = require('node:string_decoder');
const pipelineAsync = promisify(pipeline);
class LogTransform extends Transform {
constructor() {
super();
this.decoder = new StringDecoder('utf8');
this.remainder = '';
}
_transform(chunk, encoding, callback) {
try {
const decodedChunk = this.decoder.write(chunk);
const text = this.remainder + decodedChunk;
const lines = text.split(/\r?\n/);
// Останній елемент може бути неповним рядком.
this.remainder = lines.pop();
for (const line of lines) {
const transformedLine = this.transformLine(line);
if (transformedLine !== null) {
this.push(`${transformedLine}\n`);
}
}
callback();
} catch (error) {
callback(error);
}
}
_flush(callback) {
try {
// Декодер може зберігати частину багатобайтового символу.
this.remainder += this.decoder.end();
if (this.remainder.length > 0) {
const transformedLine = this.transformLine(this.remainder);
if (transformedLine !== null) {
this.push(`${transformedLine}\n`);
}
}
callback();
} catch (error) {
callback(error);
}
}
transformLine(line) {
const trimmedLine = line.replace(/\s+$/, '');
// Порожні та непотрібні рядки не передаємо далі.
if (trimmedLine.trim() === '') {
return null;
}
// Маскуємо email-адреси в логах.
return trimmedLine.replace(
/\b[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}\b/gi,
'[REDACTED_EMAIL]'
);
}
}
async function main() {
const [, , inputPath, outputPath] = process.argv;
if (!inputPath || !outputPath) {
console.error('Використання: node process-log.js <input> <output>');
process.exitCode = 1;
return;
}
if (inputPath === outputPath) {
console.error('Вхідний і вихідний файл мають бути різними.');
process.exitCode = 1;
return;
}
try {
await pipelineAsync(
fs.createReadStream(inputPath, {
encoding: null,
highWaterMark: 1024 * 1024
}),
new LogTransform(),
fs.createWriteStream(outputPath)
);
console.log(`Обробку завершено: ${outputPath}`);
} catch (error) {
console.error('Помилка потокової обробки:', error.message);
process.exitCode = 1;
}
}
main();Запуск:
node process-log.js application.log application.cleaned.logLogTransformКлас успадковується від Transform. Він реалізує два методи:
_transform() — обробляє кожен отриманий chunk;
_flush() — обробляє дані, які залишилися наприкінці потоку.
_transform_transform(chunk, encoding, callback) {
const decodedChunk = this.decoder.write(chunk);
const text = this.remainder + decodedChunk;
const lines = text.split(/\r?\n/);
this.remainder = lines.pop();
// Обробка повних рядків
}lines.pop() видаляє останній елемент масиву та повертає його. Цей елемент може бути неповним рядком, тому він зберігається в this.remainder.
Наприклад, якщо chunk містить:
first line
second liто після розділення:
[
'first line',
'second li'
]Останній елемент second li ще не можна обробляти, оскільки продовження може бути в наступному chunk.
Наступний chunk може містити:
ne
third lineПеред розділенням він об’єднується з залишком:
second line
third lineУ результаті рядки обробляються коректно незалежно від меж chunks.
_flushМетод _flush() викликається, коли вхідний потік завершився. Він потрібен для обробки останнього рядка, якщо файл не закінчується символом нового рядка.
Наприклад, файл може містити:
first line
last lineУ цьому випадку last line не матиме завершального \n, тому він залишиться в this.remainder до виклику _flush().
ReadStream за замовчуванням передає Buffer. Це корисно для потокової роботи, але межа chunk може припасти всередину UTF-8-символу.
Наприклад, кириличний символ займає кілька байтів. Якщо частина байтів опинилася в одному chunk, а решта — у наступному, прямий виклик:
chunk.toString('utf8')може некоректно декодувати символ.
StringDecoder зберігає неповну послідовність байтів і завершує її, коли отримує наступний chunk:
const decoder = new StringDecoder('utf8');
const text = decoder.write(chunk);Саме тому для потокової обробки UTF-8-тексту краще використовувати StringDecoder.
Потоки можуть працювати з різною швидкістю:
диск може швидко читати файл;
трансформація може обробляти дані повільніше;
диск призначення може повільно записувати результат.
Якщо читання не обмежувати, дані почнуть накопичуватися в пам’яті. Цю проблему вирішує backpressure — механізм, за допомогою якого повільний downstream-потік повідомляє upstream-потоку, що потрібно тимчасово сповільнитися.
Використання pipeline() забезпечує коректне з’єднання потоків і передавання backpressure:
await pipelineAsync(
readStream,
transformStream,
writeStream
);pipeline() також:
передає помилки між потоками;
закриває потоки після помилки;
завершується лише після повного запису результату;
не потребує ручного виклику pipe() для кожної частини ланцюжка.
highWaterMarkУ прикладі використано:
highWaterMark: 1024 * 1024Це приблизно 1 МБ даних на chunk для потоку читання.
Збільшення значення може зменшити кількість операцій читання, але збільшить кількість пам’яті, яку потік може використовувати. Зменшення значення зменшує розмір буферів, але може збільшити накладні витрати через частіші операції.
Не існує універсального оптимального значення. Воно залежить від:
швидкості диска;
розміру файлу;
складності трансформації;
доступної пам’яті;
середовища виконання.
highWaterMark не визначає розмір рядка. Рядок все одно може бути більшим за один chunk і потребувати накопичення в remainder.
Потоки можуть завершитися з помилкою через:
відсутній вхідний файл;
відсутність прав на читання;
відсутність прав на запис;
помилку файлової системи;
виняток під час трансформації.
У прикладі помилки передаються в callback(error):
try {
// Обробка chunk
callback();
} catch (error) {
callback(error);
}Після цього pipeline() відхиляє Promise, а помилка обробляється в catch:
try {
await pipelineAsync(readStream, transform, writeStream);
} catch (error) {
console.error(error.message);
}Не варто покладатися лише на подію finish у вихідного потоку. Вона означає, що запис завершився успішно, але сама по собі не забезпечує зручної обробки помилок усіх потоків у ланцюжку.
Створіть невеликий файл application.log:
2026-08-17 INFO User john@example.com logged in
2026-08-17 DEBUG Request completed
2026-08-17 ERROR Contact admin@example.org for detailsЗапустіть скрипт:
node process-log.js application.log application.cleaned.logВміст application.cleaned.log буде таким:
2026-08-17 INFO User [REDACTED_EMAIL] logged in
2026-08-17 DEBUG Request completed
2026-08-17 ERROR Contact [REDACTED_EMAIL] for detailsДля перевірки меж chunks можна тимчасово зменшити highWaterMark:
highWaterMark: 7Тоді рядки частіше розділятимуться між chunks, але результат має залишитися таким самим.
const content = await fs.promises.readFile(path, 'utf8');Для великих файлів потрібно використовувати createReadStream().
_transform(chunk, encoding, callback) {
const line = chunk.toString();
processLine(line);
callback();
}Chunk може містити частину рядка або кілька рядків. Потрібно зберігати незавершений фрагмент між викликами _transform().
Якщо обробляти тільки рядки, отримані після split('\n'), можна втратити останній рядок файлу без завершального переносу. Для цього потрібен _flush().
toString() для довільних UTF-8 chunksБагатобайтовий символ може бути розділений між chunks. Для тексту в UTF-8 використовуйте StringDecoder.
Виняток усередині _transform() потрібно передати через callback(error). Якщо викликати лише callback(), потік може завершитися некоректно або зависнути.
Не використовуйте один і той самий шлях для input і output. Відкриття WriteStream може очистити файл ще до того, як ReadStream прочитає його повністю.
Рядок, довший за доступну пам’ять, неможливо обробити звичайним рядковим алгоритмом без зміни формату обробки. Наведений підхід ефективний для великих файлів із рядками розумної довжини, але один окремий рядок усе одно тимчасово зберігається в пам’яті.
Великі файли потрібно обробляти потоково, а не через readFile().
ReadStream, Transform і WriteStream утворюють ефективний ланцюжок обробки.
Chunk не є гарантовано повним рядком.
StringDecoder коректно обробляє UTF-8 на межах chunks.
Незавершений рядок потрібно зберігати між викликами _transform().
_flush() потрібен для обробки даних наприкінці файлу.
pipeline() спрощує backpressure, завершення потоків і обробку помилок.
highWaterMark можна налаштовувати відповідно до характеристик файлу та середовища.