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。这个参数没多少人会在意但它直接决定了背压触发的频率以及内存占用的脉搏。举个实际场景你手头有一个读取速度远快于下游写入速度的生产者这时候下游消费不过来缓冲区就会涨满一旦超过highWaterMarkwrite()就会返回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}); });这里有个细节值得注意如果你不传encodingdata 事件吐出来的是 Buffer 对象而不是字符串。在需要拿流内容去做字符串匹配、正则提取时建议在创建流时就把 encoding 定好省得在回调里反复toString()。对于大文件的场景创建流时传{ autoClose: true }是默认行为文件读完会自动关闭 fd如果你在多个地方共享同一个文件句柄那就得手动管理close事件了。fs.createWriteStream也差不多只不过它面向的方向是往外写。写入时要注意write()的返回值和drain事件——后面背压那部分我会展开细讲。3.2 HTTP 请求与响应你天天在用却没意识到它是流Node.js 的内置 http 模块中req请求对象本身就是 Readableres响应对象本身就是 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 流类型的路上少走一些弯路。