1. 流是 Node.js 的精髓,绕不开的那道坎
做 Node.js 开发的人,前期可以不懂流(Stream),但只要你碰过文件上传、日志写入、数据导出这类需要处理大量数据的场景,迟早会撞上内存暴涨、进程卡死这类问题。这时候回头翻源码,你大概率会在某个环节看到stream这个模块的身影。流类型说白了就是 Node.js 处理数据的一套抽象机制,它把数据当作水一样的东西,一段一段地流进流出,而不是一次性把整桶水倒进内存里。
举个例子,你写了一个把 5GB 日志文件读出来再写入另一个文件的功能。如果你用fs.readFileSync一次性读,进程内存立刻冲上去好几个 GB,小服务器直接 OOM 崩给你看。但如果你用流来读和写,内存占用始终徘徊在几十 MB 上下,因为数据是一块一块地搬运的,边读边写,不会在内存里囤积。这就是流类型存在的意义:让 Node.js 在面对大体积、高吞吐的数据时,依然能保持极低的内存占用量和极高的响应速度。
这篇文章我打算把 Node.js 里最容易让人犯迷糊的流类型拆开揉碎地讲一遍:四种流型(Readable、Writable、Duplex、Transform)各自是什么概念、适合承担什么任务;Node.js 内置类型里哪些模块本身就是流(fs 流、HTTP 请求响应流、zlib 压缩流、crypto 加密流);再写几个真实项目里一定会用到的实战场景和踩坑记录。如果你已经会用pipe()但说不清背压(backpressure)到底是怎么触发的,或者你刚学 Node 没几个月、正准备处理文件/网络场景时,这篇内容会非常对口。
2. 四种流类型,凭一张图搞定记忆
2.1 读与写的基础分类
Node.js 的 stream 模块把流分成四个大类,但底层最核心的就是 Readable(可读流)和 Writable(可写流)这两个基础型号。另外两个 Duplex(双工流)和 Transform(变换流)其实是前两者的组合和增强。
用生活中最直观的类比:Readable 就像一个水龙头,它不断往外出水(生产数据),另一端的人拿桶接(消费数据)。Writable 则像一个下水道口,你往里面倒水(写入数据),它负责把水送走。如果有一个设备既能出水、又能接水,那就是 Duplex 双工流,比如一个对讲机,既能说也能听。Transform 是更特殊的一种,它就像一个净水器:进水是原水,出水是净化后的水,中间经过了处理逻辑。
我刚入门那会儿总把 Duplex 和 Transform 搞混,一度以为它们是同一个东西。后来理解了核心差异:Duplex 的读写两头是独立的、互不干扰的,读入的数据不一定会经过写出的数据处理;Transform 则是读到的每一块数据必须经过 transform 处理逻辑之后,才能被输出,读和写是连通的、有依赖关系的。
2.2 流的三种状态:暂停态、流动态与结束
搞清楚流的“开关”状态比背一百个 API 都有用。Readable 流有两种内部状态:暂停态(paused)和流动态(flowing)。处于暂停态时,数据不会主动往外涌,它安静地待在缓冲区里等你来取,此时你可以主动调用read()方法去取数据,或者监听readable事件来感知"缓冲区里有货了"。而流动态就比较"豪放",数据一旦进来就立即往外发射data事件,你只要监听data事件就能不断收到数据块。
这两者之间怎么切换,是新人最容易懵掉的地方。简单记:如果你监听了data事件,或者调用了pipe()、resume(),流就会被切换到流动态;如果你调用了pause(),或者监听readable事件,流就会回到暂停态。实际开发里大部分时候我们通过pipe()或pipeline()来托管,不怎么手动切换状态,但排错时你必须能一眼看出流当前处于什么状态,否则你会被"为什么 data 事件不触发了?"这种问题折磨一下午。
const { Readable } = require('stream'); // 手动创建一个可读流,按需产出数据 const readable = Readable.from(['第一块', '第二块', '第三块']); // 流动态:监听 data 事件 readable.on('data', (chunk) => { console.log('收到来自可读流的数据:', chunk); }); readable.on('end', () => { console.log('可读流数据全部读取完毕'); });2.3 缓冲区的存在与 highWaterMark
流内部都有一个缓冲区,用来临时存放还没来得及被消费的数据。缓冲区的大小上限由highWaterMark参数控制,默认对于对象模式是 16 个对象,对于字节模式是 16384 字节(16KB)。这个参数没多少人会在意,但它直接决定了背压触发的频率以及内存占用的脉搏。
举个实际场景:你手头有一个读取速度远快于下游写入速度的生产者,这时候下游消费不过来,缓冲区就会涨满,一旦超过highWaterMark,write()就会返回false,这就是在提醒你"我的缓冲区已经快装不下了,你先别继续写"。如果你忽略这个返回值,数据照常往里面塞,缓冲区就会无限制地膨胀,内存随之飙升,最后整个进程轰然倒下。
把highWaterMark调大,意味着缓冲区能临时存更多数据,write 返回 false 的频率降低,连带调用链停顿变少;但代价是单个请求的内存峰值变高。反过来调小,内存蹭得稳,但 CPU 频繁停顿等待的次数会变多,整体吞吐量下降。这是个典型的取舍问题,需要根据项目的实际数据量和机器配置来权衡。
3. Node.js 内置类型中的流家族
3.1 fs 流:最常用也最容易误用
fs.createReadStream()和fs.createWriteStream()是绝大部分人接触的第一个"内置流类型"。你可能以为读文件就是fs.readFile()一把梭,实际上文件稍微大一点,readFile 的内存占用就会让人手心冒汗。用流的好处是显而易见的。下面这段代码是把一个大文件分块读完的写法:
const fs = require('fs'); // 以流的方式读取文件 const readStream = fs.createReadStream('./big-data.log', { encoding: 'utf8', // 默认 64KB,可按需调整 highWaterMark: 1024 * 1024 }); let lineCount = 0; readStream.on('data', (chunk) => { // 这里 chunk 是一块一段的字符串(因为指定了 encoding) lineCount += chunk.split('\n').length - 1; }); readStream.on('end', () => { console.log(`日志总行数(约): ${lineCount}`); });这里有个细节值得注意:如果你不传encoding,data 事件吐出来的是 Buffer 对象,而不是字符串。在需要拿流内容去做字符串匹配、正则提取时,建议在创建流时就把 encoding 定好,省得在回调里反复toString()。对于大文件的场景,创建流时传{ autoClose: true }是默认行为,文件读完会自动关闭 fd;如果你在多个地方共享同一个文件句柄,那就得手动管理close事件了。
fs.createWriteStream也差不多,只不过它面向的方向是"往外写"。写入时要注意write()的返回值和drain事件——后面背压那部分我会展开细讲。
3.2 HTTP 请求与响应:你天天在用却没意识到它是流
Node.js 的内置 http 模块中,req(请求对象)本身就是 Readable,res(响应对象)本身就是 Writable。甚至可以说,HTTP 这个传输层协议天然就适合基于流的方式处理,因为请求体和响应体都是边接收边解析的,一次性读入反而违背了 HTTP 的语义。
一个常见的需求是把客户端上传的文件保存到本地磁盘,用流来处理是如此丝滑:
const http = require('http'); const fs = require('fs'); const server = http.createServer((req, res) => { if (req.url === '/upload' && req.method === 'POST') { const writeStream = fs.createWriteStream('./uploads/upload.bin'); // req 是 Readable,直接管道接到文件写入流 req.pipe(writeStream); req.on('end', () => { res.end('上传完成'); }); // 中断处理:客户端断连时记得清理半截文件 req.on('aborted', () => { writeStream.destroy(); fs.unlink('./uploads/upload.bin', () => {}); }); } else { res.end('Not Found'); } }); server.listen(8080);用req.pipe(writeStream)处理上传时,能实现真正的"边收边写",磁盘上的文件逐渐从 0 字节长到完整大小,而内存始终静悄悄。如果你在这里用了req.on('data', ...)去拼接数据,再把整个数据 Buffer 一次性写入文件,那当并发上来时,单机内存消耗会变得极其吓人。
还有一个常被忽略的继承关系:http 模块中的IncomingMessage和ServerResponse分别继承自Readable和Writable,这意味着它们天然支持流式处理的所有方法。很多框架(比如 Express)底层其实就是调用了这些流方法,比如res.sendFile()内部就封装了fs.createReadStream()加pipe()的流程。
3.3 zlib 与 crypto:处理数据流的"处理器"
Node.js 内置的 zlib 模块提供了压缩流的实现,它继承自 Stream.Transform 类型。也就是说,你可以把一段普通文本的 Readable 流接入 gzip 压缩流,再从压缩流接出去,整个过程完全不落地到磁盘、不经过内存整体缓存在转换,数据在多个流之间一站一站地传送。
const zlib = require('zlib'); const fs = require('fs'); // 把一个大 JSON 文件边读边压缩成 .gz 文件 const inputStream = fs.createReadStream('./data.json'); const gzipStream = zlib.createGzip(); const outputStream = fs.createWriteStream('./data.json.gz'); inputStream.pipe(gzipStream).pipe(outputStream); outputStream.on('finish', () => { console.log('压缩完成'); });crypto 模块同样提供流式的哈希计算和加解密能力。如果你的项目需要生成文件的 MD5 或 SHA256 校验值,传统做法是读完整个文件再算,挺费内存。用流的方式,可以边读边累计哈希值:
const crypto = require('crypto'); const fs = require('fs'); const hash = crypto.createHash('sha256'); const fileStream = fs.createReadStream('./large-file.bin'); fileStream.on('data', (chunk) => hash.update(chunk)); fileStream.on('end', () => { console.log('文件 SHA256:', hash.digest('hex')); });这么写的好处和流的初衷完全一致:不在乎文件有多大,内存占用恒定,只取决于 highWaterMark 和当前处理块的大小。crypto 流式的哈希计算在下载大文件校验场景里非常适用,我在实际项目中生成更新包校验值的时候,就是用这一个方案解决了几 GB 级文件的哈希需求。
4. 实战中的流应用:从日志处理到接口代理
4.1 大规模日志文件的实时清理与过滤
单看流的概念会觉得抽象,落到具体需求才见真章。我有个需求:服务器每天产生几十 GB 的 nginx 访问日志,运维要求把含特定错误码的条目提取出来单独成文件,其余丢弃。如果每天手动用 grep 处理,显然不是工程技术该干的事。用 Node.js 的 Transform 流写一个过滤工具,一行行地处理,内存占用永远恒定。
const { Transform } = require('stream'); const fs = require('fs'); const readline = require('readline'); // Transform 流本质是"数据加工器" class ErrorFilter extends Transform { constructor(options = {}) { super(options); this.errorLines = []; } _transform(chunk, encoding, callback) { // 把块按行拆开处理,注意跨块边界问题 const text = chunk.toString(); const lines = text.split('\n'); for (const line of lines) { if (line.includes('"status": 500') || line.includes(' 500 ')) { this.push(line + '\n'); } } callback(); } } const filter = new ErrorFilter(); const source = fs.createReadStream('./access.log'); const dest = fs.createWriteStream('./error-500.log'); source.pipe(filter).pipe(dest); dest.on('finish', () => { console.log('筛选完成,500 错误已保存到 error-500.log'); });跨块边界问题值得单独提一嘴:如果原始文件行很长,一个 chunk(比如 64KB)可能刚好在一个完整行的中间截断,上面这个简单 split('\n') 的写法就会把一行的前半段和后半段拆成两条记录。严谨的做法是缓存上一个 chunk 的尾部,检测到末尾没有换行符,就等下一个 chunk 来再拼接。这也是为什么很多成熟工具直接使用readline模块——它已经处理好了这些边界细节。
4.2 流式接口代理:边接收边转发
有时候我们要做一个代理服务,A 客户端发大文件给 B 后端,如果中间不加处理,直接req.pipe(requestToBackend),整个链路里 Node.js 进程只是起了个"水管工"的角色,不额外缓冲数据。这样一个简单的流式代理,代码量极小:
const http = require('http'); const server = http.createServer((req, res) => { if (req.url.startsWith('/api/')) { const options = { hostname: 'backend.internal', port: 9000, path: req.url, method: req.method, headers: req.headers }; const proxyReq = http.request(options, (proxyRes) => { // 把后端的响应也通过流直接返给前端 res.writeHead(proxyRes.statusCode, proxyRes.headers); proxyRes.pipe(res); }); // 关键是这里:把前端的请求流转发到后端 req.pipe(proxyReq); req.on('error', (err) => { proxyReq.destroy(err); }); } }); server.listen(3000);这种模式在微服务网关里非常常见。它的优点有两个:第一,大体积请求体不会在网关进程里被完整缓存,内存占用极低;第二,响应也是边收边吐,客户端的首字节时间比"全部收完再发"快得多。不过注意,流式转发对后端服务的稳定性要求更高,如果后端在响应中途断掉,前端会收到一个截断的响应体,要知道自己加容错,比如把 socket 超时和错误处理都做完整。
4.3 使用 pipeline 化解内存危机
早年间用pipe()写流的时候,总被一个坑卡住:如果下游出错,上游流并不会自动停止。比如你从一个大文件读取,pipe 到写入流,写入的磁盘满了抛了错误,但读流还在继续产生数据,缓冲区就开始堆积,内存上涨到不可控。Node.js 后来提供了一个更安全的接口stream.pipeline(),它会在下游出错时自动销毁整个管道。
const { pipeline } = require('stream'); const fs = require('fs'); const zlib = require('zlib'); // 写法一:回调风格 pipeline( fs.createReadStream('./data.json'), zlib.createGzip(), fs.createWriteStream('./data.json.gz'), (err) => { if (err) { console.error('管道处理失败:', err); } else { console.log('管道压缩完成'); } } ); // 写法二:Promise 风格(Node 10+ 可用 stream/promises) const { pipeline: pipelineAsync } = require('stream/promises'); async function run() { try { await pipelineAsync( fs.createReadStream('./data.json'), zlib.createGzip(), fs.createWriteStream('./data.json.gz') ); console.log('一路顺利'); } catch (err) { console.error('出错了,但流已被清理干净:', err); } }我在自己的项目里基本已经不用裸pipe()了,除非是那种特别简单、且对错误处理要求不高的场景。pipeline能自动处理上游、变换、下游之间的联系,事件监听和资源销毁都替你做了,能让代码少掉几十行错误处理逻辑。强推。
5. 背压机制:流不稳的根源
5.1 背压从哪儿来,怎么观察
背压是流式架构里最核心的概念,简单说就是下游消费速度跟不上上游生产速度时,产生的一种"反向压力"。读流读得飞快,写流写得很慢,中间的缓冲区就会被塞满,然后读流开始变慢,最终可能停止读取,整个管道被迫降速。Node.js 靠的就是write()返回 false 以及drain事件这两个机制感知和响应背压。
手动用可写流写大量数据时,光注意write()是否返回 false 还不够。标准做法是:返回 false 就停止继续写入,并且监听一次drain事件,等它触发后继续写。你见过很多新人写死循环写入,跑一段时间内存爆炸,就是忽略了这层机制。
const fs = require('fs'); const writeStream = fs.createWriteStream('./large-output.log', { highWaterMark: 1024 * 1024 // 1MB 缓冲区 }); function writeData() { let canContinue = true; for (let i = 0; i < 1000000; i++) { const line = `第 ${i} 行日志内容\n`; // 如果 write 返回 false,表示缓冲区满了,应该暂停 canContinue = writeStream.write(line); if (!canContinue) { console.log('缓冲区已满,等待 drain 事件'); writeStream.once('drain', () => { console.log('缓冲区清空,继续写入'); writeData(); }); return; } } writeStream.end('全部写完'); } writeData();5.2 对象模式下的数据块语义
默认流是以 Buffer 为数据块,而如果你设置{ objectMode: true },流中的数据不会自动转成 Buffer,而是可以传递任意 JavaScript 对象。这在做一些数据管道任务时很爽,比如从数据库读出来的记录,一条一条地放入对象模式流,再经过变换后写入导出文件。
const { Readable, Transform } = require('stream'); // 转换 csv 行的对象模式流 class CsvTransform extends Transform { _transform(row, enc, cb) { // row 是对象,不是 Buffer const csvLine = Object.values(row).join(','); this.push(csvLine + '\n'); cb(); } } const input = Readable.from([ { name: '张三', age: 28 }, { name: '李四', age: 30 } ]); const transform = new CsvTransform({ objectMode: true }); input.pipe(transform).pipe(process.stdout); // 输出: // 张三,28 // 李四,30对象模式会让代码更直观,但要注意:对象模式下的 flow 不经过 Buffer 转换,数据在流转过程中被引用,如果下游处理慢,积压的对象仍然可能导致内存膨胀。所以对象模式也不能完全放开不管背压。
5.3 压测下的三种背压表现
以我在生产环境压测的经历看,背压问题主要有三种表现。第一种是最常见的:CPU 占用正常,但内存缓慢线性增长,最终 OOM。这种多是代码里忽略了write()返回值,写流无限接收数据。第二种:CPU 飙高,内存还算稳定,但处理吞吐量上不去。这往往是 highWaterMark 设得太小,缓冲区频繁满,引发大量上下文切换。第三种更隐蔽:所有流事件的回调都在正常触发,但总数据量就是不对,最后发现是某个中间 Transform 流的_flush方法没处理完尾部数据,导致丢掉了最后一点内容。
解决这些问题的通用思路,永远是先在代码里确认"谁在产生数据、谁在等待数据、谁的缓冲首当其冲"。你可以在关键流上挂readable,drain,close事件打点日志,观察它们的触发频率,再联系实际情况判断。
6. 常见问题与排查技巧实录
6.1 事件不触发:readable 还是 data?
很多人问我:"我用 fs.createReadStream 读文件,监听 data 事件为什么不触发?" 排查时第一件事是确认流当前是暂停态还是流动态,以及是否有人调用过pause()或者unpipe()。最常见的原因是你混用了两种读取方式:同时在监听readable事件,又监听了data事件。一旦你监听readable,流就会优先进入暂停态,此时data事件不会自行发射数据,两套机制互相打架,最终表现为数据没有如期输出。
解决建议:小文件想省事,用data事件没问题;大文件或者需要逐块精确控制的,用readable事件 + 手动read()。不要在同一个流上同时用两套读取模式,这属于自找麻烦。
6.2 流不结束:文件读完却没触发 end
另一种典型情况是 end 事件迟迟不来。检查方向有两个:一是你监听了data事件后又调用了pause(),之后没再resume(),这会让流停在暂停态,数据没读完,end 自然卡着不动;二是文件流创建后,有某个错误发生了,但你只监听了end,没有监听error事件,错误把流干掉了,end 也不会触发,你的回调就一直空等。所以流式代码的error监听是必须的,别图省事。
6.3 内存为何迟迟降不下来
还有一个高频问题:明明用了流,也设置了 highWaterMark,但内存就是降不下来。多数情况是流的缓冲里堆积了大的 Buffer 引用,或者 Transform 流的_transform函数里没有及时调用callback(),导致数据积压在流内部。排查时可以把中间流的数据块大小打出来,看看每次 push 的数据体量是否都超出了预期。如果某个变换流处理得非常慢,考虑给变换流一个自己的highWaterMark,避免上游塞爆它。
6.4 环境安装与脚本执行的那些坑
因为热词里大量出现 npm.ps1 禁止执行的报错,这里额外补一句安装环节的注意点。在 Windows 上安装 Node.js 后,如果 PowerShell 执行 npm 命令报"无法加载文件 npm.ps1,因为在此系统上禁止运行脚本",这不代表 Node.js 安装有问题,也不是电脑中病毒,而是 PowerShell 的执行策略默认约束了脚本运行。临时解决办法是以管理员身份打开 PowerShell,执行Set-ExecutionPolicy RemoteSigned改成当前用户允许本地脚本;或者你嫌麻烦,直接用 CMD 而不是 PowerShell 来跑 npm 命令,也能绕过去。环境配好之后,记得验证一下node -v和npm -v输出版本号,再进入正题。
这类环境问题看起来不算流主题的一部分,但很多刚入门的朋友卡在这里,迟迟进不到代码阶段,所以我顺手放在这里做个索引。
7. 关于流的最后一点心得体会
写完好几个月的流式处理代码之后,我最大的感受是:流真正难的地方不在 API 调用,而在于建模——你得想清楚你的数据链路里谁是生产者、谁是消费者、谁在中间做变换,以及当上下游速度不匹配时,你想让谁背这个压力。Node.js 把这一整套模型像积木一样搭好,开发者只需要正确地把流对象拼起来,就能用极小的内存处理一个看起来不可能完成的大文件任务。
再分享一个我个人的小习惯:代码里尽量用pipeline()而不是裸pipe();每条流都额外挂一个error监听;创建文件流时根据后续操作预先定好 encoding 或对象模式;打印日志时把 data 事件的块大小和触发时间记录一下。这些细节点对点的累积,能让你的流式代码在并发和压力下稳健很多。希望这篇梳理能帮你在 Node.js 流类型的路上少走一些弯路。