Kodokon kodokon.com

Потоки и backpressure: работа с большими файлами

Обрабатывай многогигабайтные файлы с постоянным расходом памяти благодаря потокам, pipeline и backpressure.

9 мин · 3 вопросов

Открыть этот урок в Kodokon

readFile загружает файл целиком в память: на логе в 4 ГБ твой процесс упадёт. Потоки обрабатывают данные кусками с постоянным расходом памяти, ограниченным highWaterMark (64 КиБ по умолчанию для файлового потока). Существует четыре типа: Readable, Writable, Duplex и Transform. Главная задача, которую они решают, называется backpressure: что делать, когда производитель быстрее потребителя?

JAVASCRIPT
import { createWriteStream } from "node:fs";

const out = createWriteStream("big.txt");
let i = 0;

function writeChunks() {
  let ok = true;
  while (i < 1e6 && ok) {
    ok = out.write(`line ${i}\n`);
    i += 1;
  }
  if (i < 1e6) out.once("drain", writeChunks);
  else out.end();
}

writeChunks();
Контракт write/drain, реализованный вручную.

Когда write() возвращает false, внутренний буфер превысил highWaterMark: данные не потеряны, но если продолжать писать, всё будет копиться в памяти - именно то, чего потоки и должны были избежать. Контракт таков: прекрати запись и дождись события drain. Это то, что делает за тебя pipe(), и то, что ещё лучше делает pipeline(): помимо backpressure он пробрасывает ошибки и корректно уничтожает все потоки, в обоих направлениях.

JAVASCRIPT
import { createReadStream, createWriteStream }
  from "node:fs";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";

const upper = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, chunk.toString().toUpperCase());
  },
});

await pipeline(
  createReadStream("big.txt"),
  upper,
  createWriteStream("big-upper.txt"),
);
console.log("done");
Потоковое преобразование, постоянный расход памяти.

Потоки Readable являются асинхронными итерируемыми объектами: for await потребляет куски, соблюдая backpressure. Осторожно: кусок - это произвольный срез байтов, ничто не гарантирует, что он заканчивается на границе строки. Для построчной обработки node:readline берёт на себя склейку кусков, включая случай, когда \r\n разорван между двумя кусками, благодаря crlfDelay: Infinity.

JAVASCRIPT
import { createReadStream } from "node:fs";
import { createInterface } from "node:readline";

const rl = createInterface({
  input: createReadStream("big.txt"),
  crlfDelay: Infinity,
});

let count = 0;
for await (const line of rl) {
  if (line.includes("42")) count += 1;
}
console.log(count, "matching lines");
Построчное чтение большого файла.

Проверка знаний

Убедись, что запомнил ключевые моменты этого урока.

  1. Что означает, если writable.write(chunk) возвращает false?
    • Кусок потерян и его нужно записать заново
    • Внутренний буфер превысил highWaterMark: прекрати запись до события drain
    • Поток был уничтожен ошибкой
  2. Какое решающее преимущество даёт pipeline() по сравнению с pipe()?
    • pipeline быстрее, потому что он многопоточный
    • pipe обрабатывает backpressure, а pipeline нет
    • pipeline пробрасывает ошибки и уничтожает все задействованные потоки
    • pipeline автоматически конвертирует кодировки
  3. Что представляет собой highWaterMark в режиме objectMode?
    • Количество объектов в очереди
    • Количество байтов, как и в бинарном режиме
    • Максимальный размер сериализованного объекта