Node系列 · Node基础:文件流
大文件(GB 级别)一次性
readFile会撑爆内存;网络请求的 body 也是按块到达的——这两种"数据分批到达"的场景都靠流处理。理解流的四种类型和背压机制,就理解了 Node I/O 的精髓。
一、什么是流
流(Stream) 是一种按时间分批处理数据的抽象接口:
- 数据不是一次性到达,而是分块(chunk)逐次可用
- 内存占用只与"单 chunk 大小"有关,与"总数据量"无关
- 生产者(写入端)和消费者(读取端)通过事件解耦
为什么需要流
| 场景 | 用 readFile | 用流 |
|---|---|---|
| 10KB 配置 | ✅ 简单 | ❌ 杀鸡用牛刀 |
| 500MB 日志 | ❌ 一次性占用 500MB 内存 | ✅ 始终 64KB 缓冲区 |
| 1GB 视频 | ❌ 直接 OOM | ✅ 边读边处理 |
| 网络大文件下载 | ❌ 全部缓冲在内存 | ✅ 流式写盘 |
二、四种流类型
Node 内置 4 种抽象流:
| 类型 | 抽象类 | 方向 | 示例 |
|---|---|---|---|
| Readable | stream.Readable | 只读 | fs.createReadStream、http.IncomingMessage |
| Writable | stream.Writable | 只写 | fs.createWriteStream、http.ServerResponse |
| Duplex | stream.Duplex | 读写 | net.Socket、zlib.createGzip |
| Transform | stream.Transform | 读写 + 转换 | zlib.createGzip、crypto.createCipheriv |
三、文件流
3.1 创建文件流
javascript
const fs = require('node:fs');
const reader = fs.createReadStream('./big.txt', {
encoding: 'utf-8', // 不指定则 chunk 是 Buffer
highWaterMark: 64 * 1024, // 默认 64KB;可调
});
const writer = fs.createWriteStream('./out.txt', {
encoding: 'utf-8',
flags: 'w', // 默认 'w',会覆盖
});highWaterMark 是"水位线"——内部缓冲超过这个值就触发背压。
3.2 消费流:事件 vs async iterator
事件方式(经典):
javascript
const reader = fs.createReadStream('./big.txt', { encoding: 'utf-8' });
reader.on('data', (chunk) => {
console.log(`收到 ${chunk.length} 字节`);
});
reader.on('end', () => {
console.log('读取完成');
});
reader.on('error', (err) => {
console.error('读取失败:', err);
});async iterator(推荐):
javascript
const fs = require('node:fs');
async function countLines(filePath) {
const reader = fs.createReadStream(filePath, { encoding: 'utf-8' });
let count = 0;
for await (const chunk of reader) {
count += chunk.split('\n').length - 1;
}
return count;
}
console.log(await countLines('./access.log'));for await...of 自动处理 data / end / error 事件,代码更直观。
3.3 写入流
javascript
const writer = fs.createWriteStream('./out.txt');
writer.write('第一行\n');
writer.write('第二行\n');
writer.end(); // 结束写入;end 之后再 write 会报错
writer.on('finish', () => {
console.log('全部写入完成');
});
writer.on('error', (err) => {
console.error('写入失败:', err);
});四、pipe:流的串联
pipe 把 readable 的输出直接灌进 writable,省去手动监听 data 事件:
javascript
const fs = require('node:fs');
const { pipeline } = require('node:stream/promises');
await pipeline(
fs.createReadStream('./source.txt'),
fs.createWriteStream('./dest.txt')
);
console.log('复制完成');pipeline(推荐用 stream/promises 版本)相比 reader.pipe(writer) 的优势:
| 方案 | 错误传播 | 资源清理 |
|---|---|---|
reader.pipe(writer) | 错误需各自监听 | 失败时不会自动关另一端 |
pipeline(reader, writer) | 任一环节错误都会 reject | 失败时自动关闭所有流 |
多级 pipe 组合
javascript
await pipeline(
fs.createReadStream('./access.log'),
new Transform({
transform(chunk, encoding, callback) {
// 只保留包含 ERROR 的行
const lines = chunk.toString().split('\n').filter((l) => l.includes('ERROR'));
callback(null, lines.join('\n'));
},
}),
fs.createWriteStream('./errors.log'),
);五、背压(Backpressure)
背压是流处理最关键的概念。当写入端消化速度跟不上读取端产生速度,会发生:
读取端不断产生 chunk
↓
写入端内部缓冲不断增长
↓
超过 highWaterMark 水位
↓
内存膨胀 → 可能 OOMpipe 会自动处理背压:write() 返回 false 时暂停 readable,drain 事件触发后再恢复。手动实现就要自己监听:
javascript
const writer = fs.createWriteStream('./out.txt');
function write(data) {
if (!writer.write(data)) {
// 缓冲已满,等 drain 再写
return new Promise((resolve) => writer.once('drain', resolve));
}
// 缓冲还有空间,立即返回
return Promise.resolve();
}
await write('chunk 1\n');
await write('chunk 2\n');
await new Promise((resolve) => writer.end(resolve));WARNING
不要忽略 write() 的返回值。直接 writer.write(chunk) 不管返回值,等同于背压机制不存在——最终会撑爆内存或触发 EPIPE 错误。
六、对象模式(Object Mode)
默认流传输 Buffer / string。可以开启 object mode 传输任意 JS 值:
javascript
const { Transform } = require('node:stream');
const csv = new Transform({
readableObjectMode: true,
writableObjectMode: true,
transform(row, encoding, callback) {
// 输入是 JS 对象,输出也是 JS 对象
callback(null, { ...row, processedAt: Date.now() });
},
});适用场景:
- 逐行解析 CSV,每行是对象
- 数据库查询流式消费,每条记录是对象
- JSON 数组按对象迭代
七、流的事件完整对照
| 事件 | 触发时机 | 流类型 |
|---|---|---|
data | 每个 chunk 可读时 | Readable |
end | 数据读完 | Readable |
error | 出错(必须监听,否则崩溃) | 全部 |
drain | 写入缓冲降到水位以下,可继续写 | Writable |
finish | 全部数据已写入底层 | Writable |
close | 流及其底层资源已关闭 | 全部 |
pipe / unpipe | 被 / 被取消 pipe | Readable |
pause / resume | 暂停 / 恢复推送 | Readable |
八、流的两种工作模式
| 模式 | 默认 | 行为 |
|---|---|---|
| 流动模式(flowing) | data 事件被监听时自动进入 | 数据尽快推给消费者 |
| 暂停模式(paused) | 默认 | 数据暂存内部缓冲,等显式 read() 或 resume() |
for await...of 内部会自动切换两种模式,开发者无感。
九、最佳实践
| 场景 | 推荐 |
|---|---|
| 大文件处理(GB 级) | 必用流 |
| HTTP body 接收 | 已经是流(req 是 Readable) |
| HTTP body 发送 | 已经是流(res 是 Writable) |
| 多级流(解压 → 解密 → 解析) | pipeline 串接 |
| 错误处理 | stream/promises 的 pipeline / finished Promise 版本 |
| 背压 | pipe / pipeline 自动处理;手动实现要监听 drain |
| 不要做的 | await reader.text() 把整个流读成字符串 / JSON.parse(await reader.text()) |
十、小结
- 流的核心价值:内存占用与数据总量无关
- 4 种类型:Readable(读)/ Writable(写)/ Duplex(读写)/ Transform(读写 + 转换)
- 消费方式:
on('data')事件 或for await...of(推荐) - 多级流用
pipeline串联,自动处理错误传播和资源清理 - 背压是流处理的灵魂——
pipe/pipeline自动处理,手写必须监听drain - 文件系统流:
fs.createReadStream/fs.createWriteStream,默认highWaterMark64KB
