聊聊Node.js stream 模块,看看如何构建高性能的应用

本篇文章带大家了解 node stream 模块,介绍一下如何使用 stream 构建高性能的 node.js 应用,希望对大家有所帮助!

聊聊Node.js stream 模块,看看如何构建高性能的应用

当你在键盘上输入字符,从磁盘读取文件或在网上下载文件时,一股信息流(bits)在流经不同的设备和应用。

如果你学会处理这些字节流,你将能构建高性能且有价值的应用。例如,试想一下当你在 YouTube 观看视频时,你不需要一直等待直到完整的视频下载完。一旦有一个小缓冲,视频就会开始播放,而剩下的会在你观看时继续下载。

Nodejs 包含一个内置模块 stream 可以让我们处理流数据。在这篇文章中,我们将通过几个简单的示例来讲解 stream 的用法,我们也会描述在面对复杂案例构建高性能应用时,应该如何构建管道去合并不同的流。

在我们深入理解应用构建前,理解 Node.js stream 模块提供的特性很重要。

让我们开始吧!

Node.js 流的类型

Node.js stream 提供了四种类型的流

可读流(Readable Streams)可写流(Writable Streams)双工流(Duplex Streams)转换流(Transform Streams)

更多详情请查看 Node.js 官方文档https://nodejs.org/api/stream.html#stream_types_of_streams

让我们在高层面来看看每一种流类型吧。

可读流

可读流可以从一个特定的数据源中读取数据,最常见的是从一个文件系统中读取。Node.js 应用中其他常见的可读流用法有:

process.stdin -通过 stdin  在终端应用中读取用户输入。http.IncomingMessage – 在 HTTP 服务中读取传入的请求内容或者在 HTTP 客户端中读取服务器的 HTTP 响应。

可写流

你可以使用可写流将来自应用的数据写入到特定的地方,比如一个文件。

process.stdout  可以用来将数据写成标准输出且被 console.log 内部使用。

接下来是双工流和转换流,可以被定义为基于可读流和可写流的混合流类型。

双工流

双工流是可读流和可写流的结合,它既可以将数据写入到特定的地方也可以从数据源读取数据。最常见的双工流案例是 net.Socket,它被用来从 socket 读写数据。

有一点很重要,双工流中的可读端和可写端的操作是相互独立的,数据不会从一端流向另一端。

转换流

转换流与双工流略有相似,但在转换流中,可读端和可写端是相关联的。

crypto.Cipher 类是一个很好的例子,它实现了加密流。通过 crypto.Cipher 流,应用可以往流的可写端写入纯文本数据并从流的可读端读取加密后的密文。之所以将这种类型的流称之为转换流就是因为其转换性质。

附注:另一个转换流是 stream.PassThroughstream.PassThrough 从可写端传递数据到可读端,没有任何转换。这听起来可能有点多余,但 Passthrough 流对构建自定义流以及流管道非常有帮助。(比如创建一个流的数据的多个副本)

从可读的 Node.js 流读取数据

一旦可读流连接到生产数据的源头,比如一个文件,就可以用几种方法通过该流读取数据。

首先,先创建一个名为 myfile  的简单的 text 文件,85 字节大小,包含以下字符串:

Lorem ipsum dolor sit amet, consectetur adipiscing elit. Curabitur nec mauris turpis.

现在,我们看下从可读流读取数据的两种不同方式。

1. 监听 data 事件

从可读流读取数据的最常见方式是监听流发出的 data 事件。以下代码演示了这种方式:

const fs = require('fs')const readable = fs.createReadStream('./myfile', { highWaterMark: 20 });readable.on('data', (chunk) => {    console.log(`Read ${chunk.length} bytes\n"${chunk.toString()}"\n`);})

highWaterMark 属性作为一个选项传递给 fs.createReadStream,用于决定该流中有多少数据缓冲。然后数据被冲到读取机制(在这个案例中,是我们的 data 处理程序)。默认情况下,可读 fs 流的 highWaterMark 值是 64kb。我们刻意重写该值为 20 字节用于触发多个 data 事件。

如果你运行上述程序,它会在五个迭代内从 myfile 中读取 85 个字节。你会在 console 看到以下输出:

Read 20 bytes"Lorem ipsum dolor si"Read 20 bytes"t amet, consectetur "Read 20 bytes"adipiscing elit. Cur"Read 20 bytes"abitur nec mauris tu"Read 5 bytes"rpis."

2. 使用异步迭代器

从可读流中读取数据的另一种方法是使用异步迭代器:

const fs = require('fs')const readable = fs.createReadStream('./myfile', { highWaterMark: 20 });(async () => {    for await (const chunk of readable) {        console.log(`Read ${chunk.length} bytes\n"${chunk.toString()}"\n`);    }})()

如果你运行这个程序,你会得到和前面例子一样的输出。

可读 Node.js 流的状态

当一个监听器监听到可读流的 data 事件时,流的状态会切换成”流动”状态(除非该流被显式的暂停了)。你可以通过流对象的 readableFlowing  属性检查流的”流动”状态

我们可以稍微修改下前面的例子,通过 data 处理器来示范:

AppMall应用商店 AppMall应用商店

AI应用商店,提供即时交付、按需付费的人工智能应用服务

AppMall应用商店 56 查看详情 AppMall应用商店

const fs = require('fs')const readable = fs.createReadStream('./myfile', { highWaterMark: 20 });let bytesRead = 0console.log(`before attaching 'data' handler. is flowing: ${readable.readableFlowing}`);readable.on('data', (chunk) => {    console.log(`Read ${chunk.length} bytes`);    bytesRead += chunk.length    // 在从可读流中读取 60 个字节后停止阅读    if (bytesRead === 60) {        readable.pause()        console.log(`after pause() call. is flowing: ${readable.readableFlowing}`);        // 在等待 1 秒后继续读取        setTimeout(() => {            readable.resume()            console.log(`after resume() call. is flowing: ${readable.readableFlowing}`);        }, 1000)    }})console.log(`after attaching 'data' handler. is flowing: ${readable.readableFlowing}`);

在这个例子中,我们从一个可读流中读取  myfile,但在读取 60 个字节后,我们临时暂停了数据流 1 秒。我们也在不同的时间打印了 readableFlowing 属性的值去理解他是如何变化的。

如果你运行上述程序,你会得到以下输出:

before attaching 'data' handler. is flowing: nullafter attaching 'data' handler. is flowing: trueRead 20 bytesRead 20 bytesRead 20 bytesafter pause() call. is flowing: falseafter resume() call. is flowing: trueRead 20 bytesRead 5 bytes

我们可以用以下来解释输出:

当我们的程序开始时,readableFlowing 的值是 null,因为我们没有提供任何消耗流的机制。

在连接到 data 处理器后,可读流变为“流动”模式,readableFlowing 变为 true

一旦读取 60 个字节,通过调用 pause()来暂停流,readableFlowing 也转变为 false

在等待 1 秒后,通过调用 resume(),流再次切换为“流动”模式,readableFlowing 改为 `true’。然后剩下的文件内容在流中流动。

通过 Node.js 流处理大量数据

因为有流,应用不需要在内存中保留大型的二进制对象:小型的数据块可以接收到就进行处理。

在这部分,让我们组合不同的流来构建一个可以处理大量数据的真实应用。我们会使用一个小型的工具程序来生成一个给定文件的 SHA-256。

但首先,我们需要创建一个大型的 4GB 的假文件来测试。你可以通过一个简单的 shell 命令来完成:

On macOS: mkfile -n 4g 4gb_fileOn Linux: xfs_mkfile 4096m 4gb_file

在我们创建了假文件 4gb_file 后,让我们在不使用 stream 模块的情况下来生成来文件的 SHA-256 hash。

const fs = require("fs");const crypto = require("crypto");fs.readFile("./4gb_file", (readErr, data) => {  if (readErr) return console.log(readErr)  const hash = crypto.createHash("sha256").update(data).digest("base64");  fs.writeFile("./checksum.txt", hash, (writeErr) => {    writeErr && console.error(err)  });});

如果你运行以上代码,你可能会得到以下错误:

RangeError [ERR_FS_FILE_TOO_LARGE]: File size (4294967296) is greater than 2 GB    at FSReqCallback.readFileAfterStat [as oncomplete] (fs.js:294:11) {  code: 'ERR_FS_FILE_TOO_LARGE'}

以上报错之所以发生是因为 JavaScript 运行时无法处理随机的大型缓冲。运行时可以处理的最大尺寸的缓冲取决于你的操作系统结构。你可以通过使用内建的 buffer 模块里的 buffer.constants.MAX_LENGTH 变量来查看你操作系统缓存的最大尺寸。

即使上述报错没有发生,在内存中保留大型文件也是有问题的。我们所拥有的可用的物理内存会限制我们应用能使用的内存量。高内存使用率也会造成应用在 CPU 使用方面性能低下,因为垃圾回收会变得昂贵。

使用  buffer.constants.MAX_LENGTH 减少 APP 的内存占用

现在,让我们看看如何修改应用去使用流且避免遇到这个报错:

const fs = require("fs");const crypto = require("crypto");const { pipeline } = require("stream");const hashStream = crypto.createHash("sha256");hashStream.setEncoding('base64')const inputStream = fs.createReadStream("./4gb_file");const outputStream = fs.createWriteStream("./checksum.txt");pipeline(    inputStream,    hashStream,    outputStream,    (err) => {        err && console.error(err)    })

在这个例子中,我们使用 pipeline() 函数提供的流式方法。它返回一个“转换”流对象 crypto.createHash,为随机的大型文件生成 hash。

为了将文件内容传输到这个转换流中,我们使用 hashStreamfs.createReadStream 创建了一个可读流 4gb_file。我们将 inputStream  转换流的输出传递到可写流 hashStream 中,而 outputStream 通过 checksum.txt 创建的。

如果你运行以上程序,你将看见在 fs.createWriteStream 文件中看见 4GB 文件的 SHA-256 hash。

对流使用 checksum.txtpipeline() 的对比

在前面的案例中,我们使用 pipe() 函数来连接多个流。另一种常见的方法是使用 pipeline 函数,如下所示:

inputStream  .pipe(hashStream)  .pipe(outputStream)

但这里有几个原因,所以并不推荐在生产应用中使用 .pipe()。如果其中一个流被关闭或者出现报错,.pipe() 不会自动销毁连接的流,这会导致应用内存泄露。同样的,pipe() 不会自动跨流转发错误到一个地方处理。

因为这些问题,所以就有了 pipe(),所以推荐你使用 pipeline() 而不是 pipeline() 来连接不同的流。 我们可以重写上述的 pipe() 例子来使用  pipe() 函数,如下:

pipeline(    inputStream,    hashStream,    outputStream,    (err) => {        err && console.error(err)    })

pipeline() 接受一个回调函数作为最后一个参数。任何来自被连接的流的报错都将触发该回调函数,所以可以很轻松的在一个地方处理报错。

总结:使用 Node.js 流降低内存并提高性能

在 Node.js 中使用流有助于我们构建可以处理大型数据的高性能应用。

在这篇文章中,我们覆盖了:

四种类型的 Node.js 流(可读流、可写流、双工流以及转换流)。如何通过监听 pipeline() 事件或者使用异步迭代器来从可读流中读取数据。通过使用 data 连接多个流来减少内存占用。

一个简短的警告:你很可能不会遇到太多必须使用流的场景,而基于流的方案会提高你的应用的复杂性。务必确保使用流的好处胜于它所带来的复杂性。

更多node相关知识,请访问:nodejs 教程!

以上就是聊聊Node.js stream 模块,看看如何构建高性能的应用的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/552159.html

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
如何根据我的技术栈选择最佳的 Java 框架?
上一篇 2025年11月9日 20:37:30
电信行业应用ChatGPT的四大示例
下一篇 2025年11月9日 20:37:36

相关推荐

  • VS Code微服务开发:Docker与Kubernetes集成

    VS Code通过Docker扩展实现本地容器化开发,支持自动生成Dockerfile、一键构建镜像及devcontainer环境一致性;2. Kubernetes扩展可连接集群并管理资源,结合Bridge to Kubernetes实现本地调试与集群网络集成;3. 使用Skaffold自动化构建部…

    2026年9月24日
    000
  • 使用正则表达式从JSON数组中提取JSON对象

    本文旨在提供一种使用Java正则表达式从包含多个JSON对象的JSON数组中提取单个JSON对象的方法。我们将详细介绍如何构建合适的正则表达式,并提供示例代码演示如何在Java中使用该表达式来实现JSON对象的提取,并对提取后的字符串进行优化处理,移除不必要的空白字符。 从JSON数组中提取JSON…

    2026年9月24日
    000
  • 数据实时迁移同步工具 CloudCanal v5.2.0.0 发布,支持 SaaS 全托管

    cloudcanal 免费社区版 是 clougence 公司推出的一款全自研、可视化、自动化数据迁移同步工具,具备 结构迁移、数据迁移、数据同步、数据校验、数据订正 等功能,支持 60+ 款流行关系型数据库、实时数仓、消息中间件、缓存数据库和搜索引擎之间数据互通,其中包含国产数据库 oceanba…

    2026年9月24日
    000
  • Java Stream API:从嵌套集合中提取唯一值的两种高效方法

    本文详细介绍了如何利用Java Stream API中的flatMap()和mapMulti()操作,高效地从包含嵌套列表的复杂数据结构(如List中包含List)中提取并收集唯一的元素(如城市名称),替代传统的嵌套循环,提升代码的简洁性和可读性。 在java编程中,我们经常会遇到处理复杂数据结构的…

    2026年9月24日
    100
  • 如何在PHP的require语句中传递参数并有效管理变量作用域

    本文探讨了在php中使用`require`或`include`语句时如何向被引入文件传递参数。文章详细阐述了通过直接变量作用域共享、利用`$_get`超全局变量(不推荐)以及将引入文件内容封装为函数或类(推荐最佳实践)这三种方法,并提供了相应的代码示例,旨在帮助开发者理解和选择最适合其场景的参数传递…

    2026年9月24日
    000
  • VSCode主题开发:创建动态色彩主题的进阶技术解析

    动态主题需通过外部插件监听系统事件实现,核心是利用vscode.themeColor API响应主题切换,结合语义化作用域与Semantic Highlighting精准控制配色逻辑,实现智能自适应视觉体验。 想让VSCode主题随环境自动切换色彩?动态主题不只是换个配色那么简单。核心在于理解VSC…

    2026年9月23日
    400
  • 使用 Mp4Parser Java API 创建可播放 MP4 文件的教程

    本文档旨在指导开发者使用 Mp4Parser Java API 创建可播放的 MP4 文件。通过一个简单的复制 MP4 文件结构的例子,深入理解 Mp4Parser 的核心概念和使用方法,帮助开发者避免常见错误,并为更复杂的 MP4 文件操作打下基础。本文将重点讲解如何正确复制 MP4 文件的关键 …

    2026年9月23日
    700
  • 如何在mysql中使用连接池提升并发

    连接池通过复用数据库连接减少开销,提升高并发下系统性能;需根据语言选择HikariCP、SQLAlchemy等组件,合理配置最大连接数、空闲连接等参数,并结合数据库优化与监控调优以充分发挥效果。 在高并发场景下,频繁创建和销毁数据库连接会带来显著的性能开销。MySQL本身不直接提供连接池功能,但可以…

    2026年9月23日
    300
  • Java中利用正则表达式高效提取JSON数组中的独立对象

    本文探讨了如何使用Java的Pattern和Matcher配合正则表达式,从格式化的JSON数组字符串中精确提取出每个独立的JSON对象字符串。文章详细解析了核心正则表达式的工作原理及其对格式的依赖性,并提供了完整的Java代码示例,同时强调了在实际应用中处理JSON的注意事项和更健壮的替代方案。 …

    2026年9月23日
    600
  • Java Stream API:高效处理嵌套列表并获取唯一元素

    本文详细介绍了如何利用Java Stream API高效地从嵌套列表中提取并收集唯一的元素。通过对比flatMap()和mapMulti()两种核心操作,文章演示了如何将多层数据结构扁平化,并最终将目标属性(如城市名称)收集到一个Set中,从而避免了传统嵌套循环的复杂性,提升代码的简洁性和可读性。 …

    2026年9月23日
    500
  • Java中简易聊天室项目实现

    先运行服务器再启动多个客户端实现群聊。服务器监听8888端口,为每个客户端创建线程,接收消息并广播给其他客户端;客户端输入昵称后发送消息,通过独立线程接收广播消息,输入exit退出。 实现一个简易的Java聊天室项目,主要涉及网络编程中的Socket通信、多线程处理多个客户端连接以及简单的I/O操作…

    2026年9月23日
    100
  • 使用Java Stream高效提取嵌套集合中的唯一元素

    本教程深入探讨如何利用Java Stream API高效处理嵌套集合,从包含多层列表的对象中提取并收集唯一的元素。我们将重点介绍flatMap()和mapMulti()两种强大的流操作,演示如何将List中每个Employee对象内部的List扁平化为单一的地址流,进而简洁且高可读性地获取所有员工的…

    2026年9月23日
    100
  • win8系统要求最低配置_Win8最低系统要求

    win8系统要求最低配置_Win8最低系统要求win8系统要求最低配置_Win8最低系统要求win8系统要求最低配置_Win8最低系统要求win8系统要求最低配置_Win8最低系统要求

    如果您准备在电脑上安装Windows 8操作系统,但不确定硬件是否满足基本运行条件,则需要确认设备是否达到微软官方设定的最低配置标准。以下是确保系统可正常安装和启动的具体要求。 本文运行环境:Dell XPS 13,Windows 11 一、处理器(CPU)要求 Windows 8系统要求处理器主频…

    2026年9月23日 用户投稿
    100
  • Flink 1.16 Job Manager 重启后消息丢失问题排查及解决

    Flink 作业在遇到异常时,会根据配置的重启策略进行自动重启。但如果整个 Job Manager 重启,可能会出现消息丢失的情况。本文旨在帮助你排查和解决 Flink 1.16 中 Job Manager 重启后消息丢失的问题,涵盖了可能的原因和相应的解决方案,确保数据处理的完整性。 问题分析 当…

    2026年9月23日
    000
  • 嵌入式Linux开发-根文件系统本地挂载

    嵌入式Linux开发-根文件系统本地挂载嵌入式Linux开发-根文件系统本地挂载嵌入式Linux开发-根文件系统本地挂载嵌入式Linux开发-根文件系统本地挂载

    引言 前一篇文章介绍了根文件系统的制作与nfs网络挂载,本文将探讨如何通过本地挂载根文件系统来完成系统启动。本地挂载通常用于产品发布阶段,并且分为两种操作方式。 第一种方式:在PC机上制作好文件映像rootfs.img,然后通过uboot加载并直接烧写到EMMC中。这种方法最便捷,适用于产品批量生产…

    2026年9月23日 用户投稿
    300
  • 使用 Mp4Parser API 重构 MP4 文件:理解原子结构与常见陷阱

    本文深入探讨了如何使用 Java 的 Mp4Parser API 进行 MP4 文件的低级操作,特别是在复制或重构文件时可能遇到的问题。通过一个实际案例,文章揭示了忽略关键 MP4 原子(如 uuid)可能导致文件无法播放的原因,并提供了修复后的代码示例,强调了理解 MP4 规范和原子完整性的重要性…

    2026年9月23日
    600
  • Java中利用正则表达式从JSON数组中提取独立JSON对象

    本文详细介绍了如何利用Java正则表达式从格式化的JSON数组中提取独立的JSON对象字符串。通过一个具体的代码示例,文章展示了如何构建一个精确的正则表达式模式来匹配并分离数组中的每个JSON实体,并提供了Java代码实现,包括去除多余空白字符的步骤,最终实现将JSON数组解析为可操作的独立对象字符…

    2026年9月23日
    200
  • windows8无法弹出usb设备怎么办_windows8安全移除U盘失败解决方法

    先重启Windows资源管理器,再依次排查占用进程、使用文件资源管理器弹出、确保Plug and Play服务运行、禁用USB选择性暂停、修复注册表通知项,可解决U盘无法安全移除问题。 如果您尝试从Windows 8电脑上安全移除U盘或其他USB设备,但系统提示设备正在使用中或没有任何反应,则可能是…

    2026年9月23日
    200
  • PHP高效读取大型GZ文件:揭示Gzip的顺序访问限制与实践方法

    本教程深入探讨了php中处理大型gz压缩文件的核心挑战:其固有的顺序访问特性。我们将解释为何无法对gz文件进行随机跳转读取,以及这意味着您必须从头开始按序解压数据。文章将提供一种实用的分块读取策略,并附带php示例代码,帮助开发者高效、安全地处理超大gz文件,同时讨论潜在的跨块数据处理问题及内存管理…

    2026年9月23日
    200
  • Java Optional与集合结合使用方法

    Optional与集合结合可避免空指针异常。1. 用Optional.ofNullable包装可能为null的集合元素;2. Stream中filter后接findFirst返回Optional,安全查找;3. 对象属性为Optional时,通过flatMap展开提取值;4. 方法返回Optiona…

    2026年9月23日
    300

发表回复

登录后才能评论
关注微信