Node.js 流:FileStream、管道和事件

⚡ 智能摘要

Node.js 流以小块形式读取和写入数据,而不是将整个文件加载到内存中。本页解释了文件流的创建、管道链、背压处理以及保持大规模数据传输高效的事件模型。

  • ???? 流类型: Node.js 将每个流组织成可读、可写、双工或转换,每个 Node.js API 都返回这四种类型之一。
  • 📂 文件流: fs.createReadStream() 和 fs.createWriteStream() 可以分块读取和写入文件,而无需将整个文件加载到内存中。
  • 🔗 Piping: pipe() 方法将可读流直接连接到可写流,并自动在它们之间移动数据。
  • 背压: write() 在缓冲区满时返回 false,而 drain 事件则表示可以安全地恢复写入。
  • 📡 直播活动: data、end、error 和 finish 标志着流生命周期中的关键时刻,并驱动着大多数流处理代码。
  • 自定义事件: events 模块和 EventEmitter 允许任何 Node.js 应用程序定义、发出和监听自己的命名事件。
  • 🛠️ 监听器控制: once()、listenerCount() 和 newListener 事件管理处理程序运行的次数以及监听器如何工作。 trac凯德

Node.js 中的流类型

Node.js 公开的每个流都属于四种基本类型之一,了解给定 API 返回的类型,就能在编写任何代码之前理解该流的行为方式。Node.js 自带的 stream 模块定义了这些类别,本文中使用的文件、网络和压缩 API 都基于这些类别构建。

流类型 数据方向 示例 API 常见用例
Readable 仅源端(数据流出) fs.createReadStream() 读取文件,接收 HTTP 请求体
可写的 仅限目的地(数据流入) fs.createWriteStream() 写入文件,发送 HTTP 响应
Duplex 两个方向,独立 net.Socket(TCP 套接字) 双向网络通信
改造 无论哪个方向,数据在传输过程中都会被修改。 zlib.createGzip() 在管道中途压缩或编码数据

本文通篇使用的两个基本构建模块分别是可读流和可写流;两者都为……提供了支持。 HTTP Web 服务器 在 Node.js 中,双工流将两种角色结合在一个对象中,而转换流是一种双工流,其输出会修改其输入,例如在流中应用 gzip 压缩。

Node.js 中的文件流

定义了四种流类别之后,本文的其余部分将使用可读和可写流 API 构建文件传输。

Node 广泛使用流作为数据传输机制。

例如,当您使用 console.log 函数将任何内容输出到控制台时,您实际上是在使用流将数据发送到控制台。

Node.js 还能够从文件中流式传输数据,以便可以正确地读取和写入数据。现在我们将看一个示例,了解如何使用流从文件中读取和写入数据。对于此示例,我们需要遵循以下步骤

步骤1) 创建一个名为 data.txt 的文件,其中包含以下数据。假设此文件存储在本地计算机的 D 盘上。

Node.js 教程

介绍

活动

Generators

数据连接

使用 Jasmine

步骤2) 编写相关代码,利用流从文件读取数据。

Node.js 中的文件流

var fs = require("fs");
var stream;
stream = fs.createReadStream("D://data.txt");

stream.on("data", function(data) {
    var chunk = data.toString();
    console.log(chunk);
}); 

Code 解释:-

  1. 我们首先需要包含包含创建流所需的所有功能的“fs”模块。
  2. 接下来,我们使用 createReadStream 方法创建一个可读流。作为输入,我们给出了 data.txt 文件的位置。
  3. steam.on 函数是一个事件处理程序,在其中,我们将第一个参数指定为“数据”。这意味着每当数据从文件流中进入时,就会执行回调函数。在我们的例子中,我们定义了一个回调函数,它将执行 2 个基本步骤。第一个是将从文件读取的数据转换为字符串。第二个是将转换后的字符串作为输出发送到控制台。
  4. 我们正在从数据流中读取每个数据块并将其转换为字符串。
  5. 最后,我们将每个字符串转换块的输出发送到控制台。

输出:

Node.js 中的文件流

  • 如果代码正确执行,您将在控制台中看到上述输出。此输出与 data.txt 文件中的输出相同。

写入文件

和创建读取流一样,我们也可以创建写入流来将数据写入文件。我们先创建一个没有内容的空文件,名为 data.txt。假设这个文件放在我们电脑的 D 盘。

下面的代码显示了如何将数据写入文件。

Node.js 中的文件流

var fs = require("fs");
var stream;
stream = fs.createWriteStream("D://data.txt");

stream.write("Tutorial on Node.js")
stream.write("Introduction")
stream.write("Events")
stream.write("Generators")
stream.write("Data Connectivity")
stream.write("Using Jasmine") 

Code 解释:-

  1. 我们使用 createWriteStream 方法创建可写流。作为输入,我们给出了 data.txt 文件的位置。
  2. 接下来我们使用 stream.write 方法将不同行的文本写入文本文件。流将负责将这些数据写入 data.txt 文件。

如果你打开 data.txt 文件,你将在文件中看到以下数据

Node.js 教程

介绍

活动

Generators

数据连接

运用 Jasmine

Node.js 中的管道

独立读写流很有用,但 Node.js 还允许可读流直接连接到可写流,无需手动编写缓冲代码即可移动数据。

在 Node 应用程序中,可以使用 pipe() 方法将流连接在一起,该方法接受两个参数:

  • 必需的可写流,作为数据的目标,以及
  • 用于传递选项的可选对象。

使用管道的一个典型例子是,如果您想将数据从一个文件传输到另一个文件。

让我们看一个例子,了解如何使用管道将数据从一个文件传输到另一个文件。

步骤1) 创建一个名为 datainput.txt 的文件,其中包含以下数据。假设此文件存储在我们本地机器的 D 盘上。

Node.js 教程

介绍

活动

Generators

数据连接

使用 Jasmine

步骤2) 创建一个名为 dataOutput.txt 的空白文件并将其放在本地机器的 D 盘上。

步骤3) 编写以下代码来执行从datainput.txt文件到dataOutput.txt文件的数据传输。

Node.js 中的管道

var fs = require("fs");
var readStream = fs.createReadStream("D://datainput.txt");
var writeStream = fs.createWriteStream("D://dataOutput.txt");
readStream.pipe(writeStream);

Code 解释:-

  1. 我们首先为我们的 datainput.txt 文件创建一个“readstream”,其中包含需要传输到新文件的所有数据。
  2. 然后,我们需要为我们的 dataOutput.txt 文件创建一个“写入流”,这是我们的空文件,也是从 datainput.txt 文件传输数据的目的地。
  3. 然后我们使用管道命令将数据从读取流传输到写入流。管道命令将获取进入读取流的所有数据,并将其推送到写入流。

如果您现在打开 dataOutput.txt 文件,您将看到 datainput.txt 文件中存在的所有数据。

Node.js 中的流事件和背压

如上所示,调用 pipe() 隐藏了两个在规模化时很重要的机制:每个流发出的事件,以及防止快速源淹没慢速目标的背压系统。

  • 数据: 每次有新数据块可用时,都会通过可读流发出该数据。
  • 结束: 当可读流传输完最后一个数据块后发出。
  • 完: 当所有排队的数据都刷新完毕后,由可写流发出。
  • 错误: 无论哪种流类型出现故障,都必须处理该故障,否则进程会崩溃。

下面的示例监听所有四个事件,同时手动管理反压,这是 pipe() 自动提供的行为。

var fs = require("fs");
var readStream = fs.createReadStream("D://data.txt");
var writeStream = fs.createWriteStream("D://dataOutput.txt");

readStream.on("data", function(chunk) {
    var ok = writeStream.write(chunk);
    if (!ok) {
        readStream.pause();
    }
});

writeStream.on("drain", function() {
    readStream.resume();
});

readStream.on("end", function() {
    writeStream.end();
});

writeStream.on("finish", function() {
    console.log("Write completed.");
});

readStream.on("error", function(err) {
    console.log(err);
});

当可读流产生数据的速度超过可写流消耗数据的速度时,就会发生背压。`write()` 方法在缓冲区填满后会返回 false,因此上面的代码会暂停源操作,直到目标发出 `drain` 事件,表明可以安全地恢复操作。`pipe()` 函数会自动执行此循环。

⚠提示: 在生产环境中,建议使用 stream 模块中的 pipeline() 函数,而不是手动创建 pipe() 链,因为它会转发每个流中的错误并自动关闭文件句柄。

Node.js 中的事件

Stream 是 Node.js 更广泛的事件驱动设计中的一个例子。本文的其余部分将探讨事件的一般工作原理,并使用 Stream 内部依赖的同一个 EventEmitter 类。

事件是 Node.js 中的关键概念之一,有时 Node.js 被称为事件驱动框架。

基本上,事件就是发生的事情。例如,如果建立了与数据库的连接,则会触发数据库连接事件。事件驱动编程是创建在触发特定事件时触发的函数。

让我们看一个在 Node.js 中定义事件的基本示例。

我们将创建一个名为“data_received”的事件。当触发此事件时,文本“数据已接收”将发送到控制台。

Node.js 中的事件

var events = require('events');
var eventEmitter = new events.EventEmitter();
eventEmitter.on('data_received', function() {
    console.log('data received succesfully.');
});

eventEmitter.emit('data_received'); 

Code 解释:-

  1. 使用 require 函数包含“events”模块。使用此模块,您将能够在 Node.js 中创建事件。
  2. 创建一个新的事件发射器。这用于将事件(在我们的例子中为“data_received”)绑定到在步骤 3 中定义的回调函数。
  3. 我们定义一个事件驱动函数,该函数表示如果触发“data_received”事件,那么我们应该将文本“data_received”输出到控制台。
  4. 最后,我们使用 eventEmiter.emit 函数手动触发事件。这将触发 data_received 事件。

程序运行时,文本“数据已接收”将发送到控制台,如下所示。

Node.js 中的事件

发射事件

除了发出和监听事件之外,Node.js 还提供了额外的方法来控制监听器的注册和检查方式。

定义事件时,可以调用不同的事件方法。本主题重点介绍每种方法的详细信息。

  1. 一次性事件处理程序

有时您可能只想对第一次发生的事件做出反应。在这种情况下,您可以使用 once() 方法。

让我们看看如何利用 once 方法处理事件。

发射事件

Code 解释:-

  1. 这里我们使用‘once’方法来表示对于事件‘data_received’,回调函数应该只执行一次。
  2. 这里我们手动触发“data_received”事件。
  3. 当 'data_received' 事件再次被触发时,这次什么也不会发生。这是因为在第一步我们说过该事件只能被触发一次。

如果代码正确执行,日志中将输出“data_received successful”。此消息只会在控制台中出现一次。

  1. 检查事件监听器

在其生命周期的任何时刻,事件发射器都可以附加零个或多个侦听器。可以通过多种方式检查每种事件类型的侦听器。

如果您只对确定附加监听器的数量感兴趣,那么 EventEmitter.listenerCount() 方法就是最好的选择。

(注意: 监听器很重要,因为主程序应该知道是否正在将监听器添加到事件中,否则程序将因调用其他监听器而出现故障。)

发射事件

Code 解释:-

  1. 我们正在定义使用事件相关方法所需的 eventEmitter 类型。
  2. 然后我们定义一个名为 emitter 的对象,它将用于定义我们的事件处理程序。
  3. 我们创建了 2 个事件处理程序,它们基本上不执行任何操作。我们的示例保持简单,只是为了展示 listenerCount 方法的工作原理。
  4. 现在,当您在 data_received 事件上调用 listenerCount 方法时,它将在控制台日志中发送附加到此事件的事件监听器的数量。

如果代码正确执行,控制台日志中将显示值 2。

  1. newListener 事件

每次注册新的事件处理程序时,事件发射器都会发出 newListener 事件。此事件用于检测新的事件处理程序。当您需要为每个新的事件处理程序分配资源或执行某些操作时,通常会使用 newListener 事件。

发射事件

var events = require('events');
var eventEmitter = events.EventEmitter;
var emitter = new eventEmitter();
emitter.on("newListener", function(eventName, listener) {
    console.log("Added listener for " + eventName + " events");
});
emitter.on('data_received', function() {});
emitter.on('data_received', function() {}); 

Code 解释:-

  1. 我们正在为“newListener”事件创建一个新的事件处理程序。因此,每当注册新的事件处理程序时,控制台中都会显示文本“已添加侦听器”+ 事件名称。
  2. 这里我们向控制台写入文本“已添加监听器”+每个注册事件的事件名称。
  3. 我们为事件“data_received”定义了 2 个事件处理程序。

如果上述代码正确执行,控制台中将显示以下文本。它仅显示“newListener”事件处理程序被触发了两次。

添加了 data_received 事件的监听器

添加了 data_received 事件的监听器

常见问题

`fs.readFile()` 会将整个文件加载到内存中再返回,而 `fs.createReadStream()` 则会将文件分成小块,在数据到达时逐一返回。对于大型文件,流式传输占用的内存更少,并且允许在读取完成之前就开始处理。

是的。每个目标地址调用一次 `.pipe()`,例如先调用 `readStream.pipe(destinationA)`,再调用 `readStream.pipe(destinationB)`。Node.js 会将相同的数据块传递给每个目标地址,并且每个可写流仍然独立地管理自身的反压。

是的。将可读流通过诸如 zlib.createGzip() 之类的转换流进行管道传输,然后再传输到可写流:readStream.pipe(zlib.createGzip()).pipe(writeStream)。文件在传输过程中进行压缩,而无需将整个文件加载到内存中。

他们可以快速编写可用的 readStream 和管道样板代码,但经常会忽略错误事件处理程序和反压处理。在生产环境中使用 AI 生成的流代码之前,务必检查是否存在这两个缺陷。

他们可以提出可能的原因,例如缺少排水监听器或未处理的错误事件,但要确认真正的原因,仍然需要手动检查流状态和应用程序日志。

总结一下这篇文章: