数据处理中的一个常见问题叫作背压:数据传输时,缓冲区后方的数据不断积压。当接收端需要执行复杂操作,或因其他原因处理速度较慢时,来自数据源的数据就容易像堵塞一样堆积。 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();
});
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();
});
要说明基于流的背压机制为何是一种有效优化,可以将系统工具与 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);
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);
为了检查结果,可以尝试打开两个压缩文件。原文报告: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');
}
}
);
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');
}
}
);
也可以通过 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);
}
}
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);
}
}
数据过多、来得过快 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);
这就是背压机制的重要性。如果缺少背压系统,进程会耗尽系统内存,拖慢其他进程,并在结束之前独占大量系统资源。
结果包括:
- 拖慢其他正在运行的进程。
- 让垃圾收集器承受过大的负荷。
- 耗尽内存。
下面的原文实验移除了 .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
流的黄金法则是:始终尊重背压。好的实践应当与内部机制一致,只要避免与内部背压支持冲突的行为,就能遵循正确的做法。
一般来说:
- 没有读取请求时,不要调用 .push()。
- .write() 返回 false 后,不要继续调用它,应等待 ‘drain’。
- 流的行为会随 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);
}
}
}
下面是正确的做法:Readable 检查 this.push() 的返回值,尊重背压。
class MyReadable extends Readable {
_read(size) {
let chunk;
let canPushMore = true;
while (canPushMore && null !== (chunk = getNextChunk())) {
canPushMore = this.push(chunk);
}
}
}
自定义流之外的代码也可能无视背压。下面的反例在数据可用时(触发 ‘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));
下面是 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!' }
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!' }
这个例子创建自定义 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();
实现 ._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();
}
可以多次调用 .cork(),只需确保调用相同次数的 .uncork(),让数据重新流动。 .cork() .uncork()
结语 Conclusion
流是 Node.js 中经常使用的模块。它对内部结构非常重要,也帮助开发者扩展功能、连接 Node.js 模块生态中的各个部分。
希望现在你能够在考虑背压的前提下,排查问题并安全地编写自己的 Writable 和 Readable 流,也能与同事朋友分享这些知识。 Writable Readable
继续阅读 Stream 文档,了解其他 API 函数,在使用 Node.js 构建应用时进一步发挥流的能力。 Stream











暂无评论内容