Node.js Stream 流处理与 Buffer 操作详解
掌握流式数据处理核心,构建高效的大文件与数据管道
在处理大文件或高频数据流时,一次性将所有数据加载到内存会导致严重的性能问题。Node.js 提供了 Stream(流)和 Buffer(缓冲区)两大核心模块,让我们可以分块处理数据,实现高效的内存管理和数据管道。本文将深入讲解两者的用法,并通过实战案例掌握流式编程。
一、Buffer:二进制数据容器
Buffer 是 Node.js 中处理二进制数据的基础,它在 V8 堆内存之外分配,专门用于存储原始字节序列。
1.1 创建 Buffer
// 从字符串创建(默认 UTF-8 编码)
const buf1 = Buffer.from('Hello 世界', 'utf8');
console.log(buf1); // <Buffer 48 65 6c 6c 6f 20 e4 b8 96 e7 95 8c>
console.log(buf1.length); // 12(字节数,不是字符数)
// 从数组创建
const buf2 = Buffer.from([0x48, 0x65, 0x6c, 0x6c, 0x6f]);
console.log(buf2.toString()); // Hello
// 分配指定大小的空 Buffer
const buf3 = Buffer.alloc(10); // 填充 0
const buf4 = Buffer.allocUnsafe(10); // 不初始化,更快但有残留数据
注意:
Buffer.allocUnsafe不会清零内存,速度更快但可能包含敏感数据,仅在确定会完全覆写时使用。
1.2 读写操作
const buf = Buffer.alloc(16);
// 写入数据
buf.write('Hello', 0, 'utf8'); // 从偏移 0 写入字符串
buf.writeUInt8(255, 5); // 在偏移 5 写入单字节
buf.writeUInt16BE(1000, 6); // 大端序 16 位整数
buf.writeDoubleBE(Math.PI, 8); // 双精度浮点数
// 读取数据
console.log(buf.toString('utf8', 0, 5)); // Hello
console.log(buf.readUInt8(5)); // 255
console.log(buf.readUInt16BE(6)); // 1000
console.log(buf.readDoubleBE(8)); // 3.141592653589793
1.3 常用操作
const a = Buffer.from('Hello ');
const b = Buffer.from('World');
// 拼接
const combined = Buffer.concat([a, b]);
console.log(combined.toString()); // Hello World
// 比较
console.log(Buffer.compare(a, b)); // -1(a < b)
// 切片(共享内存)
const sliced = combined.subarray(0, 5);
console.log(sliced.toString()); // Hello
// 复制
const copied = Buffer.alloc(5);
combined.copy(copied, 0, 0, 5);
Buffer 的常见场景:文件 I/O、网络传输、加密计算。以下是一个 Base64 编解码的示例:
const text = 'Hello 世界';
// 编码
const encoded = Buffer.from(text, 'utf8').toString('base64');
console.log(encoded); // SGVsbG8g5LiW55+l
// 解码
const decoded = Buffer.from(encoded, 'base64').toString('utf8');
console.log(decoded); // Hello 世界
二、四种流类型
Node.js 的 Stream 模块定义了四种流类型:
| 类型 | 描述 | 典型场景 |
|---|---|---|
| Readable | 可读流 | 文件读取、HTTP 请求体 |
| Writable | 可写流 | 文件写入、HTTP 响应体 |
| Duplex | 双工流 | TCP Socket |
| Transform | 转换流 | Gzip 压缩、加密 |
const { Readable, Writable, Duplex, Transform } = require('stream');
// Readable 示例
const readable = Readable.from(['Hello', ' ', 'World']);
// Writable 示例
const writable = new Writable({
write(chunk, encoding, callback) {
console.log('写入:', chunk.toString());
callback();
}
});
readable.pipe(writable);
// 输出: 写入: Hello 写入: World
三、pipe 管道链
pipe 方法将可读流连接到可写流,自动处理数据流动和背压。多个流可以通过管道串联,形成处理链。
const fs = require('fs');
const zlib = require('zlib');
const crypto = require('crypto');
// 读取文件 → Gzip 压缩 → 加密 → 写入文件
fs.createReadStream('input.txt')
.pipe(zlib.createGzip()) // 压缩
.pipe(crypto.createCipheriv('aes-256-cbc', key, iv)) // 加密
.pipe(fs.createWriteStream('output.enc')); // 写入
// 更推荐使用 pipeline(更好的错误处理)
const { pipeline } = require('stream/promises');
async function compressAndEncrypt() {
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('output.gz')
);
console.log('处理完成');
}
pipeline 相比 pipe 的优势:
| 特性 | pipe | pipeline |
|---|---|---|
| 错误处理 | 手动监听 error | 自动传递错误 |
| 资源清理 | 需手动销毁 | 自动销毁所有流 |
| 代码风格 | 链式 | 函数式 |
| Promise 支持 | 否 | 是 |
四、自定义流实现
4.1 自定义 Readable
const { Readable } = require('stream');
class CounterStream extends Readable {
constructor(max, options) {
super(options);
this.max = max;
this.current = 0;
}
_read() {
if (this.current < this.max) {
const chunk = Buffer.from(`${this.current}\n`, 'utf8');
this.push(chunk);
this.current++;
} else {
this.push(null); // 结束流
}
}
}
const counter = new CounterStream(10);
counter.on('data', data => process.stdout.write(data));
counter.on('end', () => console.log('计数完成'));
4.2 自定义 Transform
Transform 流是 Duplex 的特化,它接收输入、变换后输出。
const { Transform } = require('stream');
// 行转大写
class UpperCaseStream extends Transform {
_transform(chunk, encoding, callback) {
this.push(chunk.toString().toUpperCase());
callback();
}
}
// JSON 行解析器
class JsonLineParser extends Transform {
constructor() {
super({ objectMode: true });
this.buffer = '';
}
_transform(chunk, encoding, callback) {
this.buffer += chunk.toString();
const lines = this.buffer.split('\n');
this.buffer = lines.pop(); // 保留不完整的最后一行
for (const line of lines) {
if (line.trim()) {
try {
this.push(JSON.parse(line));
} catch {
callback(new Error(`无效的 JSON: ${line}`));
return;
}
}
}
callback();
}
_flush(callback) {
if (this.buffer.trim()) {
this.push(JSON.parse(this.buffer));
}
callback();
}
}
// 使用:日志文件解析
const logStream = fs.createReadStream('app.log');
logStream
.pipe(new JsonLineParser())
.on('data', obj => console.log('日志条目:', obj.level, obj.message));
五、背压处理
当数据生产速度超过消费速度时,会产生背压。如果忽视背压,数据会在内存中堆积,最终导致内存溢出。
// 错误示范:忽略背压
readable.on('data', chunk => {
writable.write(chunk); // 如果写入慢,数据会堆积
});
// 正确示范:监听 drain 事件
readable.on('data', chunk => {
const ok = writable.write(chunk);
if (!ok) {
readable.pause();
writable.once('drain', () => readable.resume());
}
});
// 推荐方式:使用 pipe 或 pipeline 自动处理背压
readable.pipe(writable);
write() 方法返回 false 表示内部缓冲已满,此时应暂停读取,等待 drain 事件后再继续。pipe 内部已实现这一机制。
六、大文件处理实战
6.1 逐行读取
const { createReadStream } = require('fs');
const { createInterface } = require('readline');
async function processLargeFile(filePath) {
const fileStream = createReadStream(filePath);
const rl = createInterface({
input: fileStream,
crlfDelay: Infinity
});
let lineCount = 0;
let totalBytes = 0;
for await (const line of rl) {
lineCount++;
totalBytes += Buffer.byteLength(line);
if (lineCount % 100000 === 0) {
console.log(`已处理 ${lineCount} 行`);
}
}
console.log(`完成: 共 ${lineCount} 行, ${totalBytes} 字节`);
}
processLargeFile('huge-data.csv');
6.2 流式 CSV 解析
const { pipeline, Transform } = require('stream');
const { parse } = require('csv-parse');
async function analyzeCSV(inputPath, outputPath) {
let rowCount = 0;
let totalAmount = 0;
const filter = new Transform({
objectMode: true,
transform(row, encoding, callback) {
const amount = Number(row[2]);
if (amount > 1000) {
rowCount++;
totalAmount += amount;
this.push(row);
}
callback();
}
});
await pipeline(
fs.createReadStream(inputPath),
parse({ columns: true, delimiter: ',' }),
filter,
fs.createWriteStream(outputPath)
);
console.log(`筛选 ${rowCount} 条记录,总金额 ${totalAmount}`);
}
6.3 HTTP 流式上传
const http = require('http');
const fs = require('fs');
const server = http.createServer((req, res) => {
if (req.method === 'POST' && req.url === '/upload') {
// req 本身就是 Readable 流
const writeStream = fs.createWriteStream('./uploads/uploaded.dat');
req.pipe(writeStream);
writeStream.on('finish', () => {
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ message: '上传成功' }));
});
writeStream.on('error', err => {
res.writeHead(500);
res.end('上传失败');
});
}
});
七、HTTP 流式响应
流式响应可以实现服务端推送(Server-Sent Events)或分块传输,避免长时间等待。
const http = require('http');
http.createServer((req, res) => {
if (req.url === '/stream') {
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});
let id = 0;
const interval = setInterval(() => {
const data = { time: new Date(), id: id++ };
res.write(`data: ${JSON.stringify(data)}\n\n`);
if (id > 10) {
res.write('event: close\n');
res.write('data: 结束\n\n');
res.end();
clearInterval(interval);
}
}, 1000);
req.on('close', () => clearInterval(interval));
}
}).listen(3000);
大文件下载示例:
const fileStream = fs.createReadStream('large-video.mp4');
res.writeHead(200, {
'Content-Type': 'video/mp4',
'Content-Disposition': 'attachment; filename="video.mp4"'
});
fileStream.pipe(res);
// 支持断点续传
const range = req.headers.range;
if (range) {
const [start, end] = range.replace(/bytes=/, '').split('-').map(Number);
fileStream.destroy();
const partial = fs.createReadStream('large-video.mp4', { start, end });
res.writeHead(206, {
'Content-Range': `bytes ${start}-${end}/*`,
'Accept-Ranges': 'bytes'
});
partial.pipe(res);
}
八、objectMode 对象流
默认流处理的是 Buffer 或字符串。设置 objectMode: true 后,流可以处理任意 JavaScript 对象。
const { Readable } = require('stream');
const users = [
{ id: 1, name: 'Alice' },
{ id: 2, name: 'Bob' },
{ id: 3, name: 'Charlie' }
];
const userStream = Readable.from(users); // 默认 objectMode
userStream.on('data', user => {
console.log(`用户: ${user.name}`); // 直接是对象,无需 JSON.parse
});
总结
Stream 和 Buffer 是 Node.js 处理 I/O 密集型任务的两大基石。Buffer 解决了二进制数据的存储与转换问题,Stream 则提供了分块处理数据的流式模型。掌握管道链、自定义流、背压处理和 objectMode,就能构建出高效的大文件处理管道和高性能的 HTTP 流式服务。在实际开发中,优先使用 pipeline 替代 pipe,利用 stream/promises 的 Promise 接口简化异步流程,能让代码更健壮、更易维护。