导读:本期聚焦于台湾程序员创作的《在 Node.js 中如何创建可读写的双工流来处理大规模数据转换?》,敬请观看详情。处理大规模数据转换时,如果一次性把数据全部读入内存,很容易造成内存暴涨甚至进程崩溃。Node.js 的流(Stream)机制提供了一种分块处理的思路,而双工流 Duplex 则把可读端和可写端组合在一起,特别适合加密解密、编码转换、协议解析这类既要接收输入又要产出输出的场景。本文将围绕 stream.Duplex 和简化版的 stream.Transform 展开,详细讲解双工流的底层原理、手动实现步骤、背压机制的处理方式,以及如何用 pipeline 组合多个流来搭建高性能的数据转换管道,并附上完整的可运行代码示例和常见踩坑点,帮助你写出既省内存又高吞吐的数据处理程序。

当数据量达到几百 MB 甚至几个 GB 时,用 fs.readFile 一次性加载再处理的方式几乎必然导致内存溢出。正确做法是把数据切成小块,边读边处理边输出,这正是 Node.js 流的核心思想。在所有流类型中,双工流(Duplex Stream)是最容易让人困惑的一种:它同时具备可读和可写两个端,但两个端彼此独立,写入的数据并不会自动出现在读取端。理解这一点,是掌握双工流的第一步。

在 Node.js 中如何创建可读写的双工流来处理大规模数据转换?

双工流的底层原理:读写两端为何互不相通

Node.js 中的流有四种基本类型:可读流(Readable)、可写流(Writable)、双工流(Duplex)和转换流(Transform)。Duplex 本质上是同时继承了 Readable 和 Writable 的混合体,其内部维护着两个独立的缓冲区:一个用于暂存待写入下游的数据,另一个用于暂存上游写入但尚未被消费的数据。

这意味着当你调用 duplex.write(chunk) 时,数据进入的是可写端的处理逻辑;而当你调用 duplex.read() 或监听 data 事件时,读取的是可读端推送出来的数据。两者之间没有任何自动的通道。这个设计看似奇怪,实则是为了支持 TCP socket 这类场景:socket 收到的数据和发出去的数据本来就是两个方向独立的字节流。如果你的需求是「写入什么就转换出什么」,那更应该使用 stream.Transform,它是 Duplex 的一个特化子类,自动把可写端收到的数据经过 _transform 方法处理后推送到可读端。

从实现角度看,继承 stream.Duplex 时你通常需要同时实现 _write_read,或者实现 _write_final 来控制结束时机。而继承 stream.Transform 只需实现 _transform 和可选的 _flush,代码量少得多,也更不容易出错。

手动实现一个双工流:完整代码与逐行解析

下面用一个实际场景来演示:假设我们要处理一个超大文本文件,把其中所有的小写字母转换为大写,同时统计转换的字节数。这类任务用 Transform 最合适,但为了讲清楚 Duplex 的工作方式,我们先手写一个 Duplex 版本,再对比 Transform 版本。

const { Duplex } = require('stream');
const fs = require('fs');

// 手动实现双工流:把写入的数据转成大写后推送到可读端
class UpperDuplex extends Duplex {
  constructor(options) {
    super(options);
    this.processedBytes = 0;
  }

  // 可写端:收到数据块后被调用
  _write(chunk, encoding, callback) {
    const text = chunk.toString('utf8');
    const upper = text.toUpperCase();
    this.processedBytes += chunk.length;
    // 把处理结果推送到可读端的缓冲区
    // 第二个参数是编码,第三个是回调,push 返回 false 表示可读缓冲区已满
    this.push(upper);
    // 必须调用 callback 表示本次写入已完成,否则流会卡住
    callback();
  }

  // 可读端:消费者请求数据时被调用
  _read(size) {
    // 这里不需要主动做事,因为数据是 _write 里 push 进来的
    // 如果 _read 什么都不做,流会等待新数据被 push
  }

  // 可写端结束:所有数据写完后的收尾工作
  _final(callback) {
    // 通知可读端数据已全部输出完毕
    this.push(null);
    callback();
  }
}

// 使用方式:读取大文件 -> 双工流转换 -> 写入目标文件
const source = fs.createReadStream('input.txt', {
  highWaterMark: 64 * 1024 // 每次 64KB,控制内存占用
});
const dest = fs.createWriteStream('output.txt');

const upper = new UpperDuplex();
source.pipe(upper).pipe(dest);

dest.on('finish', () => {
  console.log('处理完成,共转换字节数:', upper.processedBytes);
});

upper.on('error', (err) => {
  console.error('转换出错:', err.message);
});

这段代码里有三个关键点。第一,_write 中的 callback() 必须被调用,且只能在数据处理完成后调用,否则 Node 会认为这次写入还没结束,后续数据会被阻塞。第二,push() 返回 false 时表示可读端缓冲区已经超过 highWaterMark,此时应该暂停处理,等 _read 被触发后再继续,这就是背压的手动管理。第三,_final 里调用 this.push(null) 表示可读端的数据全部输出完毕,缺少这一步下游流会永远等待。

可以看到,手动管理双工流的状态比较繁琐。如果你发现自己在 Duplex 里做的只是「写入、转换、push 出去」,那直接用 Transform 会简洁很多。

用 Transform 简化:更安全的数据转换实现

Transform 在 Duplex 的基础上帮你封装了数据流转的逻辑,你只需要专注写转换函数本身。下面是同样功能的 Transform 实现,代码量几乎减半。

const { Transform } = require('stream');
const fs = require('fs');
const { pipeline } = require('stream/promises');

class UpperTransform extends Transform {
  constructor(options) {
    super(options);
    this.processedBytes = 0;
  }

  _transform(chunk, encoding, callback) {
    this.processedBytes += chunk.length;
    // 第一参:转换后的数据;第二参:编码;第三参:错误对象(无错则为 null)
    callback(null, chunk.toString('utf8').toUpperCase());
  }

  // 可选:流结束前输出收尾数据,比如汇总信息
  _flush(callback) {
    this.push(`\n[统计] 共处理 ${this.processedBytes} 字节\n`);
    callback();
  }
}

// 用 pipeline 串联,自动处理错误传播和流销毁
async function main() {
  await pipeline(
    fs.createReadStream('input.txt', { highWaterMark: 64 * 1024 }),
    new UpperTransform(),
    fs.createWriteStream('output.txt')
  );
  console.log('全部处理完成');
}

main().catch((err) => {
  console.error('管道失败:', err);
  process.exitCode = 1;
});

_transformcallback 签名和 _write 不同:它多了一个参数用来直接传出转换结果,框架内部会自动帮你 push。需要注意的是,一个 chunk 进来不一定要对应一个 chunk 出去,你可以多次调用 this.push() 拆分输出,也可以暂时不输出(比如在做跨块的状态解析时),等凑够完整数据再输出,这正是流式解析大文件的关键技巧。

推荐使用 stream/promises 提供的 pipeline 而不是传统的 pipe 链式写法。原因在于 pipe 不会自动传播错误:一旦中间某个流出错,上游和下游可能不会被正确销毁,造成句柄泄漏。而 pipeline 会在任一流出错时销毁整条管道的所有流,并返回一个可捕获异常的 Promise,错误处理干净利落。

背压机制与常见踩坑点

背压(Backpressure)是大文件处理绕不开的话题。当数据生产速度快于消费速度时,如果不加控制,数据会在内存中不断堆积。Node.js 的处理方式是:write() 返回 false 表示写缓冲区已满,可读流的 pipe 机制检测到后自动暂停读取,等可写端触发 drain 事件再恢复。这个机制是全自动的,前提是你正确使用了 pipepipeline。如果你自己手写 write() 循环而忽略返回值,背压就完全失效了。

const fs = require('fs');

// 错误示范:忽略 write 返回值,背压失效
async function badCopy(src, dest) {
  const rs = fs.createReadStream(src);
  const ws = fs.createWriteStream(dest);
  for await (const chunk of rs) {
    ws.write(chunk); // 返回 false 时没有等待 drain,数据会在内存里堆积
  }
  ws.end();
}

// 正确示范:等待 drain 事件
function writeChunk(ws, chunk) {
  return new Promise((resolve) => {
    if (!ws.write(chunk)) {
      ws.once('drain', resolve);
    } else {
      resolve();
    }
  });
}

async function goodCopy(src, dest) {
  const rs = fs.createReadStream(src);
  const ws = fs.createWriteStream(dest);
  for await (const chunk of rs) {
    await writeChunk(ws, chunk);
  }
  ws.end();
}

除了背压之外,还有几个高频踩坑点值得注意。一是字符串和 Buffer 混用:如果在构造流时没有设置 objectMode,数据默认是 Buffer,调用 chunk.toString('utf8') 时如果一个多字节的中文字符恰好被切在两个 chunk 的边界上,会产生乱码,解决办法是保持 Buffer 一路传递,只在最终输出时再转字符串,或者使用 StringDecoder 处理边界。二是忘记调用 callback:这会让流静默挂起,没有任何报错,排查起来非常痛苦。三是在 objectMode: true 的流中,null 有特殊含义(表示流结束),不要把 null 当作普通数据 push 出去。

总结一下:处理大规模数据转换时,优先选择 stream.Transform 配合 pipeline 的组合,它封装了双工流的复杂性并内置了背压处理;只有当你确实需要读写两端独立工作时,才手动继承 stream.Duplex。合理设置 highWaterMark(一般 64KB 到 1MB 之间),全程使用 Buffer 传递二进制数据,并确保每个回调都被正确调用,就能写出既省内存又高吞吐的数据处理程序。

Node.js双工流Duplex Stream修改时间:2026-09-05 18:35:05

免责声明:已尽一切努力确保本网站所含信息的准确性。网站作品多为原创整理与精心创作,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们进行处理Email:chomcom@qq.com。
引用或转载本作品时,请注明当前出处:https://www.ipipp.com/html/20260905/51071.html,基于非商业用途的前提下,欢迎转载或二创本作品。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。