На этой странице

Есть общая проблема, возникающая при работе с данными, называемая backpressure («обратное давление»); она описывает скопление данных за буфером во время их передачи. Когда принимающая сторона передачи выполняет сложные операции или по какой-то причине медленнее, у данных из входящего источника есть склонность накапливаться, как затор.

Чтобы решить эту проблему, должна существовать система делегирования, обеспечивающая плавный поток данных из одного источника в другой. Разные сообщества решали этот вопрос по-своему, применительно к своим программам; хорошие примеры — Unix-пайпы и TCP-сокеты, и это часто называют управлением потоком (flow control). В Node.js принятым решением стали потоки (streams).

Цель этого руководства — подробнее рассказать, что такое backpressure, и как именно потоки решают эту проблему в исходном коде Node.js. Во второй части руководства будут представлены рекомендуемые лучшие практики, чтобы код вашего приложения был безопасным и оптимизированным при реализации потоков.

Мы предполагаем небольшое знакомство с общим определением backpressure, Buffer и EventEmitter в Node.js, а также некоторый опыт с Stream. Если вы не читали эту документацию, неплохо сначала заглянуть в документацию API, поскольку это поможет расширить понимание при чтении этого руководства.

В компьютерной системе данные передаются от одного процесса к другому через пайпы, сокеты и сигналы. В Node.js есть похожий механизм под названием Stream. Потоки — это здорово! Они так много делают для Node.js, и почти каждая часть внутренней кодовой базы использует этот модуль. Вам как разработчику более чем рекомендуется использовать их тоже!

const readline = require('node:readline');

// process.stdin и process.stdout — оба экземпляры потоков (Streams).
const rl = readline.createInterface({
  input: process.stdin,
  output: process.stdout,
});

rl.question('Why should you use streams? ', answer => {
  console.log(`Maybe it's ${answer}, maybe it's because they are awesome! :)`);

  rl.close();
});

Хороший пример того, почему механизм backpressure, реализованный через потоки, — это отличная оптимизация, можно продемонстрировать, сравнив внутренние системные инструменты из реализации Stream в Node.js.

В одном сценарии мы возьмём большой файл (примерно ~9 ГБ) и сожмём его знакомым инструментом zip(1).

zip The.Matrix.1080p.mkv

Пока это будет выполняться несколько минут, в другом шелле мы можем запустить скрипт, который берёт модуль Node.js zlib, оборачивающий другой инструмент сжатия — gzip(1).

const fs = require('node:fs');
const gzip = require('node:zlib').createGzip();

const inp = fs.createReadStream('The.Matrix.1080p.mkv');
const out = fs.createWriteStream('The.Matrix.1080p.mkv.gz');

inp.pipe(gzip).pipe(out);

Чтобы проверить результаты, попробуйте открыть каждый сжатый файл. Файл, сжатый инструментом zip(1), уведомит вас, что файл повреждён, тогда как сжатие, выполненное через Stream, распакуется без ошибок.

В этом примере мы используем .pipe(), чтобы доставить данные из одного конца в другой. Однако обратите внимание, что здесь не прикреплены подходящие обработчики ошибок. Если какой-то чанк данных не будет корректно получен, источник Readable или поток gzip не будут уничтожены. pump — служебный инструмент, который корректно уничтожит все потоки в пайплайне, если один из них падает или закрывается, и в этом случае он обязателен!

pump нужен только для Node.js 8.x или ранее; для Node.js 10.x или новее представлен pipeline на замену pump. Это метод модуля для соединения (pipe) потоков, пробрасывающий ошибки, корректно выполняющий очистку и предоставляющий колбэк по завершении пайплайна.

Вот пример использования pipeline:

const fs = require('node:fs');
const { pipeline } = require('node:stream');
const zlib = require('node:zlib');

// Используем API pipeline, чтобы легко соединить серию потоков
// вместе и получить уведомление, когда пайплайн полностью завершён.
// Пайплайн для эффективного gzip-сжатия потенциально огромного видеофайла:

pipeline(
  fs.createReadStream('The.Matrix.1080p.mkv'),
  zlib.createGzip(),
  fs.createWriteStream('The.Matrix.1080p.mkv.gz'),
  err => {
    if (err) {
      console.error('Pipeline failed', err);
    } else {
      console.log('Pipeline succeeded');
    }
  }
);

Можно также использовать модуль stream/promises, чтобы применять pipeline с async / await:

const fs = require('node:fs');
const { pipeline } = require('node:stream/promises');
const zlib = require('node:zlib');

async function run() {
  try {
    await pipeline(
      fs.createReadStream('The.Matrix.1080p.mkv'),
      zlib.createGzip(),
      fs.createWriteStream('The.Matrix.1080p.mkv.gz')
    );
    console.log('Pipeline succeeded');
  } catch (err) {
    console.error('Pipeline failed', err);
  }
}

Бывают случаи, когда поток Readable может отдавать данные потоку Writable слишком быстро — гораздо быстрее, чем потребитель может обработать!

Когда это происходит, потребитель начнёт ставить все чанки данных в очередь для последующего потребления. Очередь записи будет становиться всё длиннее, и из-за этого больше данных придётся держать в памяти, пока весь процесс не завершится.

Запись на диск гораздо медленнее чтения с диска, поэтому, когда мы пытаемся сжать файл и записать его на жёсткий диск, возникнет backpressure, потому что диск записи не будет успевать за скоростью чтения.

// Втайне поток говорит: «эй, эй! погоди, это уже слишком много!»
// Данные начнут накапливаться на стороне чтения буфера данных, пока
// `write` пытается угнаться за входящим потоком данных.
inp.pipe(gzip).pipe(outputFile);

Вот почему механизм backpressure важен. Если бы системы backpressure не было, процесс израсходовал бы память вашей системы, фактически замедляя другие процессы и монополизируя большую часть вашей системы до завершения.

Это приводит к нескольким вещам:

  • Замедление всех остальных текущих процессов
  • Сильно перегруженный сборщик мусора
  • Исчерпание памяти

В следующих примерах мы возьмём возвращаемое значение функции .write() и поменяем его на true, что фактически отключает поддержку backpressure в ядре Node.js. Везде, где упоминается «модифицированный» бинарник, мы говорим о запуске бинарника node без строки return ret;, а вместо неё с заменённым return true;.

Давайте взглянем на быстрый бенчмарк. Используя тот же пример, что и выше, мы провели несколько замеров времени, чтобы получить медианное время для обоих бинарников.

   trial (#)  | `node` binary (ms) | modified `node` binary (ms)
=================================================================
      1       |      56924         |           55011
      2       |      52686         |           55869
      3       |      59479         |           54043
      4       |      54473         |           55229
      5       |      52933         |           59723
=================================================================
average time: |      55299         |           55975

Оба выполняются около минуты, так что разницы почти никакой, но давайте присмотримся, чтобы подтвердить, верны ли наши подозрения. Мы используем Linux-инструмент dtrace, чтобы оценить, что происходит со сборщиком мусора V8.

Измеренное время GC (garbage collector, сборщик мусора) указывает интервалы полного цикла одной «зачистки» (sweep), выполняемой сборщиком мусора:

approx. time (ms) | GC (ms) | modified GC (ms)
=================================================
          0       |    0    |      0
          1       |    0    |      0
         40       |    0    |      2
        170       |    3    |      1
        300       |    3    |      1

         *             *           *
         *             *           *
         *             *           *

      39000       |    6    |     26
      42000       |    6    |     21
      47000       |    5    |     32
      50000       |    8    |     28
      54000       |    6    |     35

Хотя оба процесса стартуют одинаково и, кажется, нагружают GC с одинаковой частотой, становится очевидно, что через несколько секунд при корректно работающей системе backpressure она распределяет нагрузку GC по стабильным интервалам в 4–8 миллисекунд до конца передачи данных.

Однако, когда системы backpressure нет, сборка мусора V8 начинает затягиваться. Обычный бинарник вызывает GC примерно 75 раз в минуту, тогда как модифицированный бинарник вызывает его лишь 36 раз.

Это медленный и постепенно накапливающийся «долг» от растущего использования памяти. По мере передачи данных, без системы backpressure, на каждую передачу чанка используется больше памяти.

Чем больше памяти выделяется, тем больше приходится обрабатывать GC за одну зачистку. Чем больше зачистка, тем больше GC нужно решать, что можно освободить, а сканирование отсоединённых указателей в большем пространстве памяти будет потреблять больше вычислительной мощности.

Чтобы определить потребление памяти каждым бинарником, мы засекли каждый процесс командой /usr/bin/time -lp sudo ./node ./backpressure-example/zlib.js по отдельности.

Вот вывод на обычном бинарнике:

Respecting the return value of .write()
=============================================
real        58.88
user        56.79
sys          8.79
  87810048  maximum resident set size
         0  average shared memory size
         0  average unshared data size
         0  average unshared stack size
     19427  page reclaims
      3134  page faults
         0  swaps
         5  block input operations
       194  block output operations
         0  messages sent
         0  messages received
         1  signals received
        12  voluntary context switches
    666037  involuntary context switches

Максимальный размер в байтах, занятый виртуальной памятью, оказывается примерно 87,81 МБ.

А теперь, изменив возвращаемое значение функции .write(), получаем:

Without respecting the return value of .write():
==================================================
real        54.48
user        53.15
sys          7.43
1524965376  maximum resident set size
         0  average shared memory size
         0  average unshared data size
         0  average unshared stack size
    373617  page reclaims
      3139  page faults
         0  swaps
        18  block input operations
       199  block output operations
         0  messages sent
         0  messages received
         1  signals received
        25  voluntary context switches
    629566  involuntary context switches

Максимальный размер в байтах, занятый виртуальной памятью, оказывается примерно 1,52 ГБ.

Без потоков, делегирующих backpressure, выделяется на порядок больше пространства памяти — огромная разница для одного и того же процесса!

Этот эксперимент показывает, насколько оптимизирован и экономичен механизм backpressure в Node.js для вашей вычислительной системы. Теперь давайте разберём, как он работает!

Есть разные функции для передачи данных от одного процесса к другому. В Node.js есть внутренняя встроенная функция .pipe(). Есть и другие пакеты, которые вы тоже можете использовать! Но в конечном счёте, на базовом уровне этого процесса, у нас есть два отдельных компонента: источник данных и потребитель.

Когда .pipe() вызывается от источника, он сигнализирует потребителю, что есть данные для передачи. Функция pipe помогает настроить подходящие замыкания backpressure для триггеров событий.

В Node.js источник — это поток Readable, а потребитель — поток Writable (оба можно заменить на Duplex или поток Transform, но это за рамками этого руководства).

Момент, когда срабатывает backpressure, можно сузить в точности до возвращаемого значения функции .write() потока Writable. Это возвращаемое значение определяется, конечно, несколькими условиями.

В любом сценарии, где буфер данных превысил highWaterMark или очередь записи в данный момент занята, .write() вернёт false.

Когда возвращается false, включается система backpressure. Она приостановит входящий поток Readable от отправки данных и подождёт, пока потребитель снова не будет готов. Как только буфер данных опустеет, будет эмитировано событие 'drain', и входящий поток данных возобновится.

Как только очередь закончена, backpressure снова позволит отправлять данные. Использовавшееся пространство памяти освободится и подготовится к следующей порции данных.

Это фактически позволяет использовать фиксированный объём памяти в любой момент времени для функции .pipe(). Не будет утечки памяти, не будет бесконечной буферизации, и сборщику мусора придётся иметь дело лишь с одной областью памяти!

Так что, если backpressure так важен, почему вы (вероятно) о нём не слышали? Ну, ответ прост: Node.js делает всё это за вас автоматически.

Это здорово! Но не так уж здорово, когда мы пытаемся понять, как реализовать свои собственные потоки.

На большинстве машин есть размер в байтах, определяющий, когда буфер заполнен (он варьируется от машины к машине). Node.js позволяет задать свой highWaterMark, но обычно значение по умолчанию — 16 КБ (16384, или 16 для потоков в objectMode). В случаях, когда вам может захотеться поднять это значение, дерзайте, но делайте это с осторожностью!

Чтобы лучше понять backpressure, вот блок-схема жизненного цикла потока Readable, который соединяется (pipe) с потоком Writable:

                                                     +===================+
                         x-->  Piping functions   +-->   src.pipe(dest)  |
                         x     are set up during     |===================|
                         x     the .pipe method.     |  Event callbacks  |
  +===============+      x                           |-------------------|
  |   Your Data   |      x     They exist outside    | .on('close', cb)  |
  +=======+=======+      x     the data flow, but    | .on('data', cb)   |
          |              x     importantly attach    | .on('drain', cb)  |
          |              x     events, and their     | .on('unpipe', cb) |
+---------v---------+    x     respective callbacks. | .on('error', cb)  |
|  Readable Stream  +----+                           | .on('finish', cb) |
+-^-------^-------^-+    |                           | .on('end', cb)    |
  ^       |       ^      |                           +-------------------+
  |       |       |      |
  |       ^       |      |
  ^       ^       ^      |    +-------------------+         +=================+
  ^       |       ^      +---->  Writable Stream  +--------->  .write(chunk)  |
  |       |       |           +-------------------+         +=======+=========+
  |       |       |                                                 |
  |       ^       |                              +------------------v---------+
  ^       |       +-> if (!chunk)                |    Is this chunk too big?  |
  ^       |       |     emit .end();             |    Is the queue busy?      |
  |       |       +-> else                       +-------+----------------+---+
  |       ^       |     emit .write();                   |                |
  |       ^       ^                                   +--v---+        +---v---+
  |       |       ^-----------------------------------<  No  |        |  Yes  |
  ^       |                                           +------+        +---v---+
  ^       |                                                               |
  |       ^               emit .pause();          +=================+     |
  |       ^---------------^-----------------------+  return false;  <-----+---+
  |                                               +=================+         |
  |                                                                           |
  ^            when queue is empty     +============+                         |
  ^------------^-----------------------<  Buffering |                         |
               |                       |============|                         |
               +> emit .drain();       |  ^Buffer^  |                         |
               +> emit .resume();      +------------+                         |
                                       |  ^Buffer^  |                         |
                                       +------------+   add chunk to queue    |
                                       |            <---^---------------------<
                                       +============+

Если вы настраиваете пайплайн, чтобы сцепить несколько потоков для манипуляции данными, вы, скорее всего, будете реализовывать поток Transform.

В этом случае вывод из вашего потока Readable попадёт в Transform и будет направлен (pipe) в Writable.

Backpressure будет применён автоматически, но учтите, что и входящий, и исходящий highWaterMark потока Transform можно менять, и это повлияет на систему backpressure.

Начиная с Node.js v0.10, класс Stream предлагает возможность менять поведение .read() или .write(), используя версии этих функций с подчёркиванием (._read() и ._write()).

Есть документированные рекомендации по реализации читаемых потоков и реализации записываемых потоков. Будем считать, что вы их прочитали, и следующий раздел зайдёт немного глубже.

Золотое правило потоков — всегда уважать backpressure. Что считается лучшей практикой — это непротиворечивая практика. Пока вы внимательно избегаете поведения, конфликтующего с внутренней поддержкой backpressure, вы можете быть уверены, что следуете хорошей практике.

В общем,

  1. Никогда не делайте .push(), если вас об этом не просят.
  2. Никогда не вызывайте .write() после того, как он вернул false; вместо этого ждите события 'drain'.
  3. Потоки меняются между разными версиями Node.js и в зависимости от используемой библиотеки. Будьте внимательны и тестируйте.

Что касается пункта 3, невероятно полезный пакет для построения браузерных потоков — readable-stream. Родд Вэгг (Rodd Vagg) написал отличный пост в блоге, описывающий полезность этой библиотеки. Вкратце, она обеспечивает своего рода автоматическую плавную деградацию (graceful degradation) для потоков Readable и поддерживает старые версии браузеров и Node.js.

Пока что мы посмотрели, как .write() влияет на backpressure, и много сосредоточились на потоке Writable. В силу устройства Node.js, данные технически текут вниз по течению — от Readable к Writable. Однако, как мы можем наблюдать в любой передаче данных, вещества или энергии, источник так же важен, как и приёмник, и поток Readable жизненно важен для того, как обрабатывается backpressure.

Оба этих процесса полагаются друг на друга для эффективной коммуникации: если Readable игнорирует, когда поток Writable просит его прекратить присылать данные, это может быть так же проблематично, как и когда возвращаемое значение .write() некорректно.

Так что, помимо уважения к возвращаемому значению .write(), мы должны также уважать возвращаемое значение .push(), используемого в методе ._read(). Если .push() возвращает false, поток прекратит читать из источника. Иначе он продолжит без паузы.

Вот пример плохой практики использования .push():

// Это проблематично, поскольку полностью игнорирует возвращаемое значение push,
// которое может быть сигналом backpressure от потока-приёмника!
class MyReadable extends Readable {
  _read(size) {
    let chunk;
    while (null !== (chunk = getNextChunk())) {
      this.push(chunk);
    }
  }
}

Вот пример хорошей практики, где поток Readable уважает backpressure, проверяя возвращаемое значение this.push():

class MyReadable extends Readable {
  _read(size) {
    let chunk;
    let canPushMore = true;
    while (canPushMore && null !== (chunk = getNextChunk())) {
      canPushMore = this.push(chunk);
    }
  }
}

Кроме того, извне пользовательского потока есть подводные камни при игнорировании backpressure. В этом контрпримере хорошей практики код приложения проталкивает данные всякий раз, когда они доступны (о чём сигнализирует событие 'data'):

// Это игнорирует механизмы backpressure, которые Node.js настроил,
// и безусловно проталкивает данные, независимо от того, готов ли
// поток-приёмник к ним или нет.
readable.on('data', data => writable.write(data));

Вот пример использования .push() с читаемым потоком.

const { Readable } = require('node:stream');

// Создаём пользовательский читаемый поток
const myReadableStream = new Readable({
  objectMode: true,
  read(size) {
    // Проталкиваем некоторые данные в поток
    this.push({ message: 'Hello, world!' });
    this.push(null); // Отмечаем конец потока
  },
});

// Потребляем поток
myReadableStream.on('data', chunk => {
  console.log(chunk);
});

// Вывод:
// { message: 'Hello, world!' }

В этом примере мы создаём пользовательский читаемый поток, который проталкивает один объект в поток с помощью .push(). Метод ._read() вызывается, когда поток готов потреблять данные, и в этом случае мы сразу проталкиваем некоторые данные в поток и отмечаем конец потока, протолкнув null.

Затем мы потребляем поток, слушая событие 'data' и логируя каждый чанк данных, протолкнутый в поток. В этом случае мы проталкиваем лишь один чанк данных, поэтому видим только одно сообщение в логе.

Напомним, что .write() может вернуть true или false в зависимости от некоторых условий. К счастью для нас, при построении своего потока Writable конечный автомат потока обработает наши колбэки и определит, когда обрабатывать backpressure и оптимизировать поток данных за нас.

Однако, когда мы хотим использовать Writable напрямую, мы должны уважать возвращаемое значение .write() и внимательно относиться к этим условиям:

  • Если очередь записи занята, .write() вернёт false.
  • Если чанк данных слишком велик, .write() вернёт false (предел задаётся переменной highWaterMark).
// Этот writable некорректен из-за асинхронной природы колбэков JavaScript.
// Без инструкции return для каждого колбэка, кроме последнего,
// велика вероятность, что будет вызвано несколько колбэков.
class MyWritable extends Writable {
  _write(chunk, encoding, callback) {
    if (chunk.toString().indexOf('a') >= 0) {
      callback();
    } else if (chunk.toString().indexOf('b') >= 0) {
      callback();
    }
    callback();
  }
}

// Правильно это следует написать так:
if (chunk.contains('a')) {
  return callback();
}

if (chunk.contains('b')) {
  return callback();
}
callback();

Есть также вещи, на которые стоит обратить внимание при реализации ._writev(). Эта функция связана с .cork(), но есть распространённая ошибка при написании:

// Использование .uncork() дважды здесь делает два вызова на уровне C++, что делает
// технику cork/uncork бесполезной.
ws.cork();
ws.write('hello ');
ws.write('world ');
ws.uncork();

ws.cork();
ws.write('from ');
ws.write('Matteo');
ws.uncork();

// Правильно это написать, используя process.nextTick(), который срабатывает
// на следующей итерации event loop.
ws.cork();
ws.write('hello ');
ws.write('world ');
process.nextTick(doUncork, ws);

ws.cork();
ws.write('from ');
ws.write('Matteo');
process.nextTick(doUncork, ws);

// Как глобальная функция.
function doUncork(stream) {
  stream.uncork();
}

.cork() можно вызывать сколько угодно раз, нужно лишь быть внимательным и вызвать .uncork() столько же раз, чтобы поток снова потёк.

Потоки — часто используемый модуль в Node.js. Они важны для внутренней структуры, а для разработчиков — для расширения и связывания в рамках экосистемы модулей Node.js.

Надеемся, теперь вы сможете диагностировать проблемы и безопасно писать свои потоки Writable и Readable с оглядкой на backpressure, а также делиться знаниями с коллегами и друзьями.

Обязательно почитайте подробнее про Stream, чтобы узнать о других функциях API, которые помогут улучшить и раскрыть ваши возможности работы с потоками при построении приложения на Node.js.