Node.js Stream 流式编程与 Buffer 二进制处理

深入 Node.js Stream 流式编程模式:Readable/Writable/Transform/Duplex 流、管道链、背压控制,以及 Buffer 二进制数据处理与编码转换。

1. 为什么需要 Stream

Stream 是 Node.js 的核心抽象,处理大量数据时无需一次性加载到内存,边读边处理,内存占用恒定。

无 Stream(全部加载):
  读取 1GB 文件 → 加载到内存 → 处理 → 写入新文件
  内存占用: 1GB+

有 Stream(边读边写):
  读取 chunk ──→ 处理 chunk ──→ 写入 chunk
  内存占用: 16KB (chunk size)

2. 四种流类型

类型说明示例
Readable可读流(数据源)fs.createReadStream、HTTP req
Writable可写流(数据目标)fs.createWriteStream、HTTP res
Duplex可读可写(双向)net.Socket、TCP 连接
Transform转换流(读入→转换→写出)zlib.createGzip、加密流

2.1 Readable 流

const fs = require('fs');

// 创建可读流
const readable = fs.createReadStream('large-file.txt', {
    encoding: 'utf8',
    highWaterMark: 16 * 1024  // 每次读取 16KB
});

// 消费数据
readable.on('data', (chunk) => {
    console.log(`Received ${chunk.length} bytes`);
});
readable.on('end', () => console.log('Done'));
readable.on('error', (err) => console.error(err));

// 或使用异步迭代器(更现代)
(async () => {
    for await (const chunk of readable) {
        console.log(chunk);
    }
})();

2.2 Writable 流

const writable = fs.createWriteStream('output.txt');

writable.write('Hello ');
writable.write('World\n');
writable.end('Done!');

writable.on('finish', () => console.log('写入完成'));

2.3 Transform 转换流

const { Transform } = require('stream');

// 自定义转换流:将所有输入转为大写
const upperCaseTransform = new Transform({
    transform(chunk, encoding, callback) {
        this.push(chunk.toString().toUpperCase());
        callback();
    }
});

// 使用管道链
fs.createReadStream('input.txt')
    .pipe(upperCaseTransform)
    .pipe(fs.createWriteStream('output.txt'));

3. 管道与背压

3.1 Pipe 管道链

const fs = require('fs');
const zlib = require('zlib');

// 压缩文件流式处理
fs.createReadStream('input.txt')
    .pipe(zlib.createGzip())      // 压缩
    .pipe(fs.createWriteStream('input.txt.gz'));

// 解压
fs.createReadStream('input.txt.gz')
    .pipe(zlib.createGunzip())
    .pipe(fs.createWriteStream('output.txt'));

3.2 背压(Backpressure)

当 Writable 消费速度慢于 Readable 生产速度时,需要背压机制防止内存无限增长。

const readable = fs.createReadStream('big-file.txt');
const writable = fs.createWriteStream('output.txt');

// pipe 自动处理背压:
// readable 检测到 writable.write() 返回 false → 暂停读取
// writable 的 drain 事件触发 → 恢复读取
readable.pipe(writable);

// 手动处理背压
readable.on('data', (chunk) => {
    const canContinue = writable.write(chunk);
    if (!canContinue) {
        readable.pause();  // 暂停读取
        writable.once('drain', () => readable.resume());  // 等 writable 就绪
    }
});

4. Buffer 二进制处理

Buffer 是 Node.js 处理二进制数据的类,固定大小、直接操作内存。

// 创建 Buffer
const buf1 = Buffer.alloc(10);           // 10 字节,初始化为 0
const buf2 = Buffer.from('Hello');       // 从字符串创建
const buf3 = Buffer.from([0x48, 0x65]);  // 从数组创建

// 读写
buf1.write('Node.js');
console.log(buf1.toString());  // 'Node.js'
console.log(buf1.toString('hex'));  // 十六进制
console.log(buf1.toString('base64'));  // Base64

// 切片(共享内存,不复制)
const slice = buf2.slice(0, 3);
slice[0] = 0x41;  // 修改 slice 会影响原 buf2!

// 拼接
const buf4 = Buffer.concat([buf2, buf3]);

// 编码转换
const utf8 = Buffer.from('你好', 'utf8');
console.log(utf8.toString('utf8'));  // '你好'

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「nodejs」更多文章

  1. Node.js ORM 深度对比:Prisma、TypeORM、Sequelize 与 Drizzle
  2. Node.js 设计模式与最佳实践:从 SOLID 到六边形架构
  3. Node.js 高级测试策略:从单元测试到混沌工程的完整实践