首页 / 文章 / 理解 Node.js 流:它们实际解决的问题

理解 Node.js 流:它们实际解决的问题

了解 Node.js 流存在的理由、管道机制的内部运作方式,以及背压在高效处理大量数据时的真正含义。

1261 词

大多数关于流处理的教程都是从 API 接口 .pipe()ReadableWritableTransform 开始讲解的,却从未涉及这些工具真正要解决的难题。这样的教学方式会让流处理显得像是一套必须记住的繁琐规则。如果反过来,从问题出发,API 的作用就会自然而然地显现出来。

问题所在:有些数据太大,无法一次性存储

想象一下,有一个任务需要将 4GB 大小的 CSV 文件转换为 JSON 格式。简单的处理方式如下:

const fs = require("fs");

const data = fs.readFileSync("export.csv", "utf8");
const rows = data.split("\n").map(parseRow);
fs.writeFileSync("export.json", JSON.stringify(rows));

readFileSync 会在将整个 4GB 的文件作为单个 JavaScript 字符串加载到内存中之前不会让程序继续执行,而这种字符串表示形式通常会占用比文件原始大小所暗示的更多内存。在只有 2GB 可用 RAM 的机器上,这段代码不仅运行效率低下,还会直接崩溃,或者导致其他争夺相同内存空间的进程无法正常工作。逻辑本身并没有问题,所有的转换步骤都是正确的。真正的缺陷在于 readFileSync 中的一个隐含假设:即无论文件大小如何,都可以同时将整个文件存储在内存中。

流式处理背后的真正理念

流式处理并非要求一次性获取整个数据集,而是基于不同的原理:先请求一部分数据进行处理,然后再请求下一部分。完整文件从未全部存在于内存中,任何时刻都只有一小部分数据存在,且这部分数据会在下一部分数据到来之前被处理并释放。

const fs = require("fs");

const readStream = fs.createReadStream("export.csv", { encoding: "utf8" });

readStream.on("data", (chunk) => {
  console.log(`Received ${chunk.length} characters`);
});

readStream.on("end", () => {
  console.log("Done reading the whole file, piece by piece");
});

请注意,此代码并未考虑文件的总大小。无论源文件是4GB还是4KB,这段代码的运行方式都完全相同,只不过两种情况下的内存占用峰值差异巨大。这正是流式处理的精髓:用无需预先获取完整数据集的优势,换取能够立即开始处理并在内存中仅保留少量数据即可完成工作的能力。

管道:将源直接连接到目标

手动逐个提取数据块确实是一种可行的方法,但实际上更常见的做法是将可读流直接连接到可写流,让数据自动从源头流向目的地,而无需手动传输每个数据块:

const fs = require("fs");

fs.createReadStream("export.csv")
  .pipe(fs.createWriteStream("export-copy.csv"));

只需两行代码就能复制任意大小的文件,而无需将整个文件加载到内存中。.pipe()并没有什么神奇之处,它只是将源端发出的data事件传递给目标端的write方法,同时还处理了一个比复制机制本身更为重要的细节。

几乎没人能讲清楚的部分:背压机制

这正是 .pipe() 所要解决的真正问题,而不仅仅是提供便利。想象这样一种情况:读取操作很快完成,比如从本地磁盘读取,而写入操作却很慢,可能是通过带宽受限的网络连接进行的。

readStream.on("data", (chunk) => {
  writeStream.write(chunk); // what happens if this can't keep up?
});

当调用 .write() 的速度超过目标端实际能够处理数据的速度时,可写入的流不会抛出错误或阻塞执行。相反,它会将多余的数据 silently 积累在内部内存缓冲区中,一直保存到有机会刷新为止。如果这种读取速度与写入速度之间的差距持续存在且足够大,流处理原本要解决的问题就会再次出现:内存使用量会无限制地增长,只是被推迟而非立即发生,而且通常在引发崩溃之前很难察觉。

背压机制正是为防止此类故障模式而设计的:它让可写入流能够发出信号,表明其已达到容量上限,需要生产者减少数据输出量;而实现良好的生产者会响应这一信号,而非强行继续写入。

readStream.on("data", (chunk) => {
  const canContinue = writeStream.write(chunk);
  if (!canContinue) {
    readStream.pause(); // stop reading until the writable side catches up
  }
});

writeStream.on("drain", () => {
  readStream.resume(); // writable side is ready for more
});

一旦可写入流的内部缓冲区超出预设限制,.write()方法就会返回false,以此作为暂停数据输出的信号,直到流触发drain事件,表明其已清空积压数据并准备好接收更多内容。这种暂停与恢复的机制正是.pipe()在后台自动处理的流程:

readStream.pipe(writeStream); // handles backpressure for you, silently, correctly

真正应该优先使用 .pipe() 而非手动连接 datawrite 监听器的理由在于此,而非简洁性。若手动转发数据块同时忽略 .write() 的返回值,就会重新出现流式处理原本旨在避免的无限制内存问题,这只不过是比 readFileSync 中固有的明显错误稍逊一筹罢了。

转换流:在传输过程中处理数据

有些情况下,仅仅将数据原封不动地从一处传输到另一处是不够的,还需要在传输过程中对数据进行重新格式化。Transform 流正是为这类需求而设计的:它位于管道链的中间,接收传入的数据块,对每个数据块执行特定操作,然后将处理结果转发给后续的处理环节:

const { Transform } = require("stream");

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

fs.createReadStream("input.txt")
  .pipe(upperCaseTransform)
  .pipe(fs.createWriteStream("output.txt"));

在数据处理流程中,每个数据块都会被转换为大写形式,而且无论是以原始格式还是转换后的格式,整个文件都无需完全加载到内存中。这正是 Node 内置的 zlib.createGzip() 所采用的机制。它本质上只是一个 Transform 流,会在数据块到达时立即对其进行压缩;和其它转换操作一样,它也可以像大写转换示例那样直接被接入处理链中:

const zlib = require("zlib");

fs.createReadStream("export.csv")
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream("export.csv.gz"));

读取、压缩和写入操作会同时进行,针对逐块到达的数据进行处理,这三个阶段都无需一次性加载整个文件。

为何值得真正理解这一点

Stream 被认为是 Node API 中较为棘手的领域之一,说实话,其原始的、基于事件的接口确实显得不够流畅,因此这种评价也并非毫无根据。但其背后的原理其实很简单:避免将整个数据集加载到内存中,而是逐部分处理,并确保快速的数据生成方不会在此过程中悄悄淹没缓慢的数据消费方。一旦你将这一理念内化为实际模型,.pipe()Transform 以及背压机制就不再像是需要单独记忆的三个互不相关的 API 了。它们实际上是一个整体概念,通过 API 中相互关联的三个部分体现出来,解决了像 readFileSync 这样的函数从未需要面对的问题——因为后者从一开始就只设计用于处理小文件而已。

相关阅读