Node.js 流中的背压 Backpressuring in Streams

数据处理中的一个常见问题叫作背压:数据传输时,缓冲区后方的数据不断积压。当接收端需要执行复杂操作,或因其他原因处理速度较慢时,来自数据源的数据就容易像堵塞一样堆积。 backpressure

要解决这个问题,必须有一种协调机制,保证数据从一个来源平稳流向另一个目标。不同社区根据各自程序的特点采用了不同方案,Unix 管道和 TCP 套接字就是典型例子,这类机制通常称为流量控制。Node.js 采用流来解决这个问题。

本指南详细说明背压是什么,以及 Node.js 源码中的流究竟如何处理背压。后半部分介绍实现流时应遵循的最佳实践,使应用代码更加安全、高效。

这里假定你已经了解 Node.js 中背压、Buffer 和 EventEmitter 的基本概念,并有一定的 Stream 使用经验。如果尚未阅读相关文档,建议先查看 API 文档,这会帮助你理解本指南。 backpressure Buffer EventEmitters Stream

数据处理的问题 The Problem with Data Handling

计算机系统通过管道、套接字和信号在进程之间传递数据。Node.js 中有一种类似的机制,称为 Stream。流为 Node.js 提供了许多能力,内部代码几乎各处都会使用这个模块;开发者也非常值得使用它。 Stream

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

// process.stdin and process.stdout are both instances of 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();
});
JavaScript

import readline from 'node:readline';

// process.stdin and process.stdout are both instances of 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();
});
JavaScript

要说明基于流的背压机制为何是一种有效优化,可以将系统工具与 Node.js 的 Stream 实现进行比较。 Stream

第一种情况是用熟悉的 zip(1) 工具压缩一个大文件,大小约为 9 GB。 zip(1)

zip The.Matrix.1080p.mkv

这可能需要几分钟。在另一个 shell 中,可以运行使用 Node.js zlib 模块的脚本;该模块封装了另一种压缩工具 gzip(1)。 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);
JavaScript

import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';

const gzip = createGzip();

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

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

为了检查结果,可以尝试打开两个压缩文件。原文报告:zip(1) 生成的文件提示损坏,而通过 Stream 完成的压缩文件能够正常解压。 zip(1) Stream

这个例子通过 .pipe() 将数据从一端送到另一端,但没有配置适当的错误处理器。如果某个数据块接收失败,Readable 数据源或 gzip 流不会被销毁。pump 工具会在流水线中的某个流失败或关闭时,正确销毁所有流,因此这种情况下需要它。 pump

pump 仅在 Node.js 8.x 或更早版本中有必要。Node.js 10.x 及以后引入了 pipeline,可以替代 pump。这个模块方法能够连接多个流、转发错误、正确清理资源,并在流水线结束时调用回调。 pump pipeline pump

下面是使用 pipeline 的例子:

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

// Use the pipeline API to easily pipe a series of streams
// together and get notified when the pipeline is fully done.
// A pipeline to gzip a potentially huge video file efficiently:

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');
    }
  }
);
JavaScript

import fs from 'node:fs';
import { pipeline } from 'node:stream';
import zlib from 'node:zlib';

// Use the pipeline API to easily pipe a series of streams
// together and get notified when the pipeline is fully done.
// A pipeline to gzip a potentially huge video file efficiently:

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');
    }
  }
);
JavaScript

也可以通过 stream/promises 模块,配合 async/await 使用 pipeline: stream/promises

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);
  }
}
JavaScript

import fs from 'node:fs';
import { pipeline } from 'node:stream/promises';
import zlib from '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);
  }
}
JavaScript

数据过多、来得过快 Too Much Data, Too Quickly

有时 Readable 流向 Writable 提供数据的速度过快,远远超过消费者能够处理的速度。 Readable Writable

这时,消费者会将所有数据块排入队列,等待之后处理。写入队列越来越长,整个过程结束前必须在内存中保存越来越多的数据。

写磁盘通常比读磁盘慢。因此,压缩文件并写入硬盘时,如果写入速度跟不上读取速度,就会产生背压。

// Secretly the stream is saying: "whoa, whoa! hang on, this is way too much!"
// Data will begin to build up on the read side of the data buffer as
// `write` tries to keep up with the incoming data flow.
inp.pipe(gzip).pipe(outputFile);
JavaScript

这就是背压机制的重要性。如果缺少背压系统,进程会耗尽系统内存,拖慢其他进程,并在结束之前独占大量系统资源。

结果包括:

  • 拖慢其他正在运行的进程。
  • 让垃圾收集器承受过大的负荷。
  • 耗尽内存。

下面的原文实验移除了 .write() 函数原本的返回值,改为返回 true,从而禁用 Node.js 核心中的背压支持。文中所说的“修改版”二进制程序,是将 return ret; 替换为 return true; 后的 node 程序。 return value

垃圾收集的额外负担 Excess Drag on Garbage Collection

先看一个简单基准测试。使用前面的同一个例子,原作者多次计时,比较两个二进制程序的运行时间。

   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

二者都大约需要一分钟,看起来差异不大。进一步检查 V8 垃圾收集器的行为,可以验证这一判断;原文使用 Linux 工具 dtrace 进行观察。 dtrace

下面的 GC(垃圾收集)计时表示垃圾收集器完成一次完整扫描周期的时间:

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 的工作节奏也似乎一致。但运行几秒后,正常背压系统会将 GC 负荷分散到稳定的 4–8 毫秒间隔中,直到数据传输结束。

没有背压系统时,V8 垃圾收集耗时开始变长。正常二进制程序每分钟约触发 75 次 GC,而修改版仅触发 36 次。

这是内存使用不断增长所积累的缓慢债务。缺少背压时,每一次数据块传输都会使用更多内存。

分配的内存越多,GC 每次扫描需要处理的工作就越多。扫描范围越大,需要判断哪些对象可以释放的成本就越高;在更大的内存空间中扫描失去引用的指针,会消耗更多计算资源。

内存耗尽 Memory Exhaustion

为了确定两个二进制程序的内存消耗,原作者分别使用 /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 MB。

将 .write() 的返回值改掉后,得到: return value .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 GB。

缺少流的背压协调后,分配的内存空间增加了一个数量级。同样的处理过程,差距非常明显。

这个实验展示了 Node.js 背压机制对计算资源的优化效果。接下来分析其工作原理。

背压如何解决这些问题 How Does Backpressure Resolve These Issues?

在不同处理环节之间传递数据可以使用各种函数。Node.js 内置了 .pipe(),也有其他可用的软件包。但从最基本的层面看,这个过程只有两个不同组成部分:数据源和消费者。 .pipe() other packages

在数据源上调用 .pipe() 时,它会通知消费者有数据需要传输。pipe 函数为相应事件的触发建立背压相关闭包。 .pipe()

在 Node.js 中,数据源是 Readable 流,消费者是 Writable 流。二者也可以是 Duplex 或 Transform 流,但这种情况超出了本指南的讨论范围。 Readable Writable Duplex Transform

背压何时触发,可以精确定位到 Writable 的 .write() 函数返回值。当然,这个返回值由若干条件决定。 Writable .write()

数据缓冲区超过 highWaterMark,或写入队列正忙时,.write() 会返回 false。 highWaterMark .write()

返回 false 后,背压系统开始工作:暂停传入的 Readable 流,停止发送数据,直到消费者重新就绪。数据缓冲区清空后,会发出 ‘drain’ 事件,恢复传入的数据流。 Readable ‘drain’

队列处理完成后,背压机制再次允许发送数据。之前占用的内存空间会释放,为下一批数据做好准备。

这样,.pipe() 在任一时刻便能够使用一定范围内的内存,而不是无限缓冲;垃圾收集器也只需处理相应的内存区域。原文将这种效果表述为没有内存泄漏、没有无限缓冲。 .pipe()

背压如此重要,你却可能从未听说过它,原因很简单:Node.js 自动替你完成了这些工作。

这非常便利,但也使我们在实现自定义流时更难理解内部机制。

缓冲区达到某个字节阈值后就被视为已满,具体值可能因机器或环境而异。Node.js 允许自定义 highWaterMark。原文给出的常见默认值是 16 KB(16384 字节),objectMode 流则是 16 个对象。可以提高阈值,但要谨慎。 highWaterMark

.pipe() 的生命周期 Lifecycle of .pipe()

为了更清楚地理解背压,下面用流程图展示 Readable 流通过管道连接 Writable 流的生命周期: Readable piped 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 流。 Transform

这种情况下,Readable 的输出进入 Transform,再通过管道进入 Writable。 Readable Transform Writable

Readable.pipe(Transformable).pipe(Writable);

背压会自动应用。不过,Transform 流的输入和输出两侧都可以调整 highWaterMark,这些设置都会影响背压系统。 Transform

背压使用准则 Backpressure Guidelines

从 Node.js v0.10 开始,Stream 类允许通过带下划线的函数 ._read() 和 ._write(),修改 .read() 和 .write() 的行为。 Node.js v0.10 Stream .read() .write() ._read() ._write()

文档分别提供了实现 Readable 流和 Writable 流的指南。这里假定你已经读过;下一节会进一步深入。 implementing Readable streams implementing Writable streams

实现自定义流应遵循的规则 Rules to Abide By When Implementing Custom Streams

流的黄金法则是:始终尊重背压。好的实践应当与内部机制一致,只要避免与内部背压支持冲突的行为,就能遵循正确的做法。

一般来说:

  1. 没有读取请求时,不要调用 .push()。
  2. .write() 返回 false 后,不要继续调用它,应等待 ‘drain’。
  3. 流的行为会随 Node.js 版本和所用库而变化,要谨慎并进行测试。

对于第三点,构建浏览器流时,readable-stream 是非常有用的软件包。Rodd Vagg 的文章介绍了它的用途。它为 Readable 流提供自动平稳降级能力,并支持旧版浏览器和 Node.js。 readable-stream great blog post Readable

Readable 流的具体规则 Rules specific to Readable Streams

前面主要讨论了 .write() 对背压的影响,以及 Writable 流。Node.js 中的数据实际沿着 Readable 到 Writable 的方向流动。但无论传输的是数据、物质还是能量,来源与目标都同样重要;Readable 对背压处理至关重要。 .write() Writable Readable Writable Readable

这两个过程依靠相互通信才能有效工作。如果 Readable 无视 Writable 要求停止发送数据的信号,问题会与 .write() 返回值错误一样严重。 Readable Writable .write()

因此,除了尊重 .write() 的返回值,也必须尊重 ._read() 中 .push() 的返回值。.push() 返回 false 时,应停止从数据源读取;否则可以继续。 .write() .push() ._read() .push()

下面是错误使用 .push() 的例子: .push()

// This is problematic as it completely ignores the return value from the push
// which may be a signal for backpressure from the destination stream!
class MyReadable extends Readable {
  _read(size) {
    let chunk;
    while (null !== (chunk = getNextChunk())) {
      this.push(chunk);
    }
  }
}
JavaScript

下面是正确的做法:Readable 检查 this.push() 的返回值,尊重背压。

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

自定义流之外的代码也可能无视背压。下面的反例在数据可用时(触发 ‘data’ 事件),就强行将数据向下游发送。 ‘data’ event

// This ignores the backpressure mechanisms Node.js has set in place,
// and unconditionally pushes through data, regardless if the
// destination stream is ready for it or not.
readable.on('data', data => writable.write(data));
JavaScript

下面是 Readable 流使用 .push() 的例子。 .push()

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

// Create a custom Readable stream
const myReadableStream = new Readable({
  objectMode: true,
  read(size) {
    // Push some data onto the stream
    this.push({ message: 'Hello, world!' });
    this.push(null); // Mark the end of the stream
  },
});

// Consume the stream
myReadableStream.on('data', chunk => {
  console.log(chunk);
});

// Output:
// { message: 'Hello, world!' }
JavaScript

import { Readable } from 'node:stream';

// Create a custom Readable stream
const myReadableStream = new Readable({
  objectMode: true,
  read(size) {
    // Push some data onto the stream
    this.push({ message: 'Hello, world!' });
    this.push(null); // Mark the end of the stream
  },
});

// Consume the stream
myReadableStream.on('data', chunk => {
  console.log(chunk);
});

// Output:
// { message: 'Hello, world!' }
JavaScript

这个例子创建自定义 Readable 流,通过 .push() 推入一个对象。流准备好消费数据时会调用 ._read();这里立即推入数据,再推入 null,标记流已经结束。 .push() ._read()

接着通过监听 ‘data’ 事件消费该流,并记录每个被推入的数据块。因为只推入一个数据块,所以只看到一条日志。

Writable 流的具体规则 Rules specific to Writable Streams

前面提到,.write() 会根据条件返回 true 或 false。实现自定义 Writable 流时,流的状态机会处理回调,决定何时处理背压,并优化数据流动。 .write() Writable stream state machine

直接使用 Writable 时,必须尊重 .write() 的返回值,并注意以下条件: Writable .write()

  • 写入队列忙时,.write() 返回 false。 .write()
  • 数据块过大时,.write() 返回 false;阈值由 highWaterMark 表示。 .write() highWaterMark

// This writable is invalid because of the async nature of JavaScript callbacks.
// Without a return statement for each callback prior to the last,
// there is a great chance multiple callbacks will be called.
class MyWritable extends Writable {
  _write(chunk, encoding, callback) {
    if (chunk.toString().indexOf('a') >= 0) {
      callback();
    } else if (chunk.toString().indexOf('b') >= 0) {
      callback();
    }
    callback();
  }
}

// The proper way to write this would be:
if (chunk.contains('a')) {
  return callback();
}

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

实现 ._writev() 时也有需要注意的地方。这个函数与 .cork() 配合使用,但常见错误如下: ._writev() .cork()

// Using .uncork() twice here makes two calls on the C++ layer, rendering the
// cork/uncork technique useless.
ws.cork();
ws.write('hello ');
ws.write('world ');
ws.uncork();

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

// The correct way to write this is to utilize process.nextTick(), which fires
// on the next 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);

// As a global function.
function doUncork(stream) {
  stream.uncork();
}
JavaScript

可以多次调用 .cork(),只需确保调用相同次数的 .uncork(),让数据重新流动。 .cork() .uncork()

结语 Conclusion

流是 Node.js 中经常使用的模块。它对内部结构非常重要,也帮助开发者扩展功能、连接 Node.js 模块生态中的各个部分。

希望现在你能够在考虑背压的前提下,排查问题并安全地编写自己的 Writable 和 Readable 流,也能与同事朋友分享这些知识。 Writable Readable

继续阅读 Stream 文档,了解其他 API 函数,在使用 Node.js 构建应用时进一步发挥流的能力。 Stream

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容