引言
把 10GB 的日志读进内存再处理,程序会直接 OOM;把 500MB 的 CSV 一次性 JSON.parse 成数组,进程会卡死几秒。流(Stream) 的意义就在于:用固定大小的缓冲区处理任意大小的数据,让内存占用与数据总量解耦。
但流也是 Node 里最容易出错的 API 之一:忘记监听 error 导致进程崩溃、把 pipe 当万能导致错误不传播、混用 Node Streams 与 Web Streams 导致类型不匹配、缓冲区设置不当导致吞吐骤降。TypeScript 能在其中帮上忙——前提是你会用它的类型参数。
本文聚焦 TypeScript 流式 I/O 的工程落地:从 Streams 的类型体系讲起,覆盖背压与 pipeline、大文件读写、CSV/JSON 流式解析、类型化 Transform、Web Streams 互操作、错误与资源释放,最后给出内存与吞吐调优。
目录
- 1. Node Streams 的类型体系
- 2. 背压与 pipeline
- 3. 大文件读写实践
- 4. CSV 与 JSON 流式解析
- 5. Transform 流与类型化
- 6. Web Streams 与 Node 互操作
- 7. 错误处理与资源释放
- 8. 内存与吞吐调优
- 9. 流式上传与下载
- 10. 生产实践与踩坑清单
1. Node Streams 的类型体系
1.1 四种基本流
Readable 是数据源(文件、HTTP 请求、标准输入);Writable 是数据汇(文件、HTTP 响应、标准输出);Duplex 既可读又可写(TCP socket);Transform 是 Duplex 的特化,写入什么就转换出什么。
1.2 类型参数与 objectMode
import { Readable, Writable } from "node:stream"
const src: Readable = Readable.from(["a", "b", "c"]) // Readable<string>
const sink: Writable = new Writable({ write(chunk, _enc, cb) { cb() } })
默认 objectMode: false 时元素类型是 Buffer;开启 objectMode 后可承载任意对象,此时泛型参数才真正生效。多数代码把流交给库内部处理,类型被 any 掉——要获得类型安全,需显式声明 Transform<I, O> 的输入输出类型。
1.3 Buffer 与 string 的取舍
const strStream = Readable.from(["hello"], { encoding: "utf8" }) // chunk 为 string
不设置 encoding 时 chunk 是 Buffer;设置后自动转成 string。文本处理设 encoding 更省事,但要注意多字节字符被缓冲区切分时的边界问题。
一句话总结:Readable 产、Writable 收、Duplex 双向、Transform 转换——泛型参数表达元素类型,
objectMode决定承载Buffer还是任意对象。
2. 背压与 pipeline
2.1 什么是背压
当写入端处理不过来时,读端仍源源不断产数据,缓冲区就会膨胀直至耗尽内存。背压(backpressure)是下游向上游传递「慢一点」信号的机制:write() 返回 false 表示缓冲区已满。
2.2 手动处理背压
async function pump(src: Readable, dst: Writable) {
for await (const chunk of src) {
if (!dst.write(chunk)) await once(dst, "drain") // 等缓冲区排空
}
dst.end()
}
2.3 用 pipeline 自动处理
import { pipeline } from "node:stream/promises"
await pipeline(
fs.createReadStream("big.csv"),
new Transform({ /* ... */ }),
fs.createWriteStream("out.csv"),
)
pipeline 自动串联背压、自动在出错时销毁所有流、返回 Promise 便于 try/catch。新代码应一律使用 pipeline 而非 .pipe()。
2.4 pipe 与 pipeline 的差别
| 维度 | pipe | pipeline |
|---|---|---|
| 错误传播 | 不传播 | 传播并销毁 |
| 资源清理 | 需手动 | 自动 |
| 返回值 | 目标流 | Promise |
.pipe() 的经典坑是:源流出错时目标流不会关闭,导致文件句柄泄漏。
一句话总结:背压是下游对上游的减速信号,
write()返回false时等drain——优先用pipeline,它自动处理背压、错误传播与资源销毁。
3. 大文件读写实践
3.1 不要用 readFile
// 反例:整个文件进内存
const text = await fs.readFile("10gb.log", "utf8")
// 正例:流式逐行
const rl = readline.createInterface({ input: fs.createReadStream("10gb.log") })
for await (const line of rl) if (line.includes("ERROR")) handle(line)
readFile 的内存占用与文件大小成正比,createReadStream 则恒定在一个缓冲区大小。
3.2 逐行处理与追加写
const rl = createInterface({ input: createReadStream("access.log"), crlfDelay: Infinity })
for await (const line of rl) { /* 处理每一行 */ }
const out = createWriteStream("out.log", { flags: "a" }) // 追加而非截断
readline 自动处理 \n 与 \r\n,crlfDelay: Infinity 保证 CRLF 被当作单个换行。日志类场景用 a 追加;多个进程同时以 a 写同一文件在多数文件系统上是安全的(O_APPEND 原子),用 w 则会互相覆盖。
3.3 大文件处理的常见模式
| 模式 | 适用 | 要点 |
|---|---|---|
| 逐行过滤 | 日志检索 | readline |
| 分块转换 | 编解码 | Transform |
| 分片并行 | 可独立处理 | 按偏移切分 |
| 流式聚合 | 统计 | 常数空间累加 |
一句话总结:大文件必须流式处理,
readFile的内存占用与文件大小成正比——逐行用readline,转换用Transform,需要并行时按偏移分片。
4. CSV 与 JSON 流式解析
4.1 CSV 流式解析
await pipeline(
createReadStream("users.csv"),
parse({ columns: true, skip_empty_lines: true }),
new Writable({
objectMode: true,
write(row, _enc, cb) { /* row: Record<string, string> */ cb() },
}),
)
csv-parse 逐行产出对象而非一次性构建整个数组。手写 split(",") 会在引号内包含逗号时出错,务必用成熟解析器。
4.2 JSON 流式解析
await pipeline(
createReadStream("events.json"),
parser(),
streamArray(), // 逐个产出数组元素
new Writable({ objectMode: true, write({ value }, _e, cb) { handle(value); cb() } }),
)
JSON.parse 要求完整字符串;stream-json 则能增量解析超大 JSON 数组,内存占用与数组长度无关。每行一个对象的 NDJSON 则是最省事的流式格式。
4.3 解析的坑
编码不是 UTF-8 时要显式指定;BOM 头会让第一列字段名带上不可见字符;超大单元格仍会占用内存;解析器报错后要终止整条 pipeline 而非跳过,否则数据静默丢失。
一句话总结:CSV 与 JSON 都要用流式解析器而非手写切分——
csv-parse与stream-json让内存占用与数据量解耦,NDJSON 是最省事的流式格式。
5. Transform 流与类型化
5.1 自定义 Transform
import { Transform, TransformCallback } from "node:stream"
class UpperCase extends Transform {
_transform(chunk: Buffer, _enc: BufferEncoding, cb: TransformCallback) {
cb(null, chunk.toString().toUpperCase())
}
}
5.2 类型化输入输出
class ParseLine extends Transform {
constructor() { super({ objectMode: true }) } // 承载对象
_transform(line: string, _enc: BufferEncoding, cb: TransformCallback) {
cb(null, JSON.parse(line)) // 运行时仍需校验
}
}
Transform 继承自 Duplex,可同时声明读端与写端的类型;在 objectMode 下 _transform 的 chunk 才是业务对象类型。
5.3 处理 flush
class Batch extends Transform {
private buf: Item[] = []
constructor(private size = 100) { super({ objectMode: true }) }
_transform(item: Item, _e: BufferEncoding, cb: TransformCallback) {
this.buf.push(item)
if (this.buf.length >= this.size) { const b = this.buf; this.buf = []; cb(null, b) }
else cb()
}
_flush(cb: TransformCallback) { // 流结束前冲刷残余
if (this.buf.length) this.push(this.buf)
cb()
}
}
忘记实现 _flush 会丢失最后不足一批的数据,这是批处理场景最常见的 bug。此外 _transform 里可以 await,但要确保任何分支都调用 cb,否则流会永久挂起。
一句话总结:自定义 Transform 要显式声明
objectMode与元素类型——批处理必须实现_flush冲刷残余,异步分支务必在所有路径上调用cb。
6. Web Streams 与 Node 互操作
6.1 两套 Streams 的差异
Node Streams 基于 EventEmitter,是 Node 的历史 API;Web Streams 是 WHATWG 标准,基于 ReadableStream/WritableStream/TransformStream,在浏览器、Deno、Bun 与 Node 中都可用。fetch 的响应体就是 Web Streams。
6.2 双向转换
import { Readable } from "node:stream"
const webStream: ReadableStream = Readable.toWeb(Readable.from(["a", "b"]))
const res = await fetch("https://example.com/big.json")
await pipeline(Readable.fromWeb(res.body as ReadableStream), createWriteStream("big.json"))
6.3 互操作的坑
Readable.fromWeb 需要 Web 的 ReadableStream 而非 NodeJS.ReadableStream 类型,混用时需断言;两套流的背压机制不同,转换点可能成为瓶颈;取消语义不同,一方 cancel 未必传递到另一方。转换应尽量少,最好只在边界做一次。
一句话总结:Node Streams 与 Web Streams 用
toWeb/fromWeb互转——fetch响应体是 Web Stream,转换点越少越好,取消与背压语义在两套体系间并不完全等价。
7. 错误处理与资源释放
7.1 error 事件必须监听
const stream = createReadStream("maybe-missing.txt")
stream.on("error", (err) => logger.error({ err }, "read failed"))
未监听的 error 事件会直接让进程崩溃。这是流最容易踩的坑之一,尤其在回调风格代码里。
7.2 pipeline 的自动清理
try {
await pipeline(src, transform, dst)
} catch (err) {
logger.error({ err }, "pipeline failed")
// pipeline 已自动销毁所有流,无需手动 close
}
pipeline 出错时会销毁沿途所有流并释放句柄,这是它相对 .pipe() 的最大优势。
7.3 finally 与取消
const handle = await fs.open("data.bin", "r")
try { await pipeline(handle.createReadStream(), dst) }
finally { await handle.close() }
即使 pipeline 抛错,finally 也保证文件句柄被关闭——忘记关闭句柄是长跑服务泄漏的主要来源。用户取消上传或进程收到 SIGTERM 时,要用 AbortController 主动销毁流。
一句话总结:流必须监听
error,否则进程崩溃——pipeline负责出错时的级联销毁,但外部资源(文件句柄、连接)仍要在finally中显式关闭。
8. 内存与吞吐调优
8.1 highWaterMark
const src = createReadStream("big.bin", { highWaterMark: 1024 * 1024 }) // 1MB
highWaterMark 是单个流的内部缓冲上限。默认 64KB 适合小数据;大文件顺序读写适当调大(256KB~1MB)能显著提升吞吐,但会提高内存占用。
8.2 避免字符串拼接
// 反例:反复拼接产生大量中间字符串
let all = ""
for await (const chunk of src) all += chunk
// 正例:累积 Buffer 数组,最后一次性合并
const parts: Buffer[] = []
for await (const chunk of src) parts.push(chunk)
const all = Buffer.concat(parts)
8.3 压缩加密与并行度
await pipeline(createReadStream("data.json"), zlib.createGzip(), createWriteStream("data.json.gz"))
压缩在写出前、加密在最外层;顺序错误会导致压缩率下降或加密无效。压缩本身是 CPU 密集操作,会降低吞吐,需要权衡。并行策略上:顺序读写用单流加大 highWaterMark,独立分片用多流并行加限并发,CPU 密集转换应移到 worker_threads。
一句话总结:吞吐调优的关键是
highWaterMark、避免字符串拼接、以及合理的并行度——压缩加密注意顺序,CPU 密集转换应移到 Worker 线程。
9. 流式上传与下载
9.1 HTTP 请求体是流
app.post("/upload", (req, res) => {
const dst = createWriteStream(`/tmp/${randomName()}`)
req.pipe(dst) // 或 await pipeline(req, dst)
dst.on("finish", () => res.json({ ok: true }))
})
不把请求体读进内存,即可支持超大文件上传;配合 content-length 限制可防滥用。
9.2 下载与流式转发
app.get("/download/:id", (req, res) => {
res.setHeader("Content-Disposition", `attachment; filename="${req.params.id}"`)
pipeline(createReadStream(path), res)
})
const upstream = await fetch(sourceUrl)
await pipeline(Readable.fromWeb(upstream.body as ReadableStream), res)
流式转发不需要落盘,适合代理与转码网关;但要处理上游断开与下游取消,避免连接泄漏。
9.3 断点续传
通过 Range 头与 Accept-Ranges 实现断点续传:客户端请求 Range: bytes=1000-,服务端用 createReadStream(path, { start: 1000 }) 返回 206 与 Content-Range。
一句话总结:HTTP 请求与响应都是流,直接
pipeline即可零内存转发——用Range头实现断点续传,注意处理上游断开与下游取消。
10. 生产实践与踩坑清单
10.1 决策表
| 需求 | 做法 | 避免 |
|---|---|---|
| 读大文件 | createReadStream | readFile |
| 逐行处理 | readline | 手动 split |
| 多步转换 | pipeline | 链式 .pipe() |
| 超大 JSON | stream-json | JSON.parse |
| 超大 CSV | csv-parse | 手写切分 |
10.2 踩坑清单
忘记监听 error 导致进程崩溃;.pipe() 出错时下游不关闭导致句柄泄漏;_flush 未实现导致批处理丢尾部数据;_transform 某分支忘记调用 cb 导致流挂起;highWaterMark 过小拖慢吞吐;混用两套 Streams 时类型断言掩盖真实不匹配;对象模式下把巨大数组塞进单个 chunk。
10.3 调试手段与纪律
用 --trace-gc 观察 GC 是否因流缓冲频繁触发;用 process.memoryUsage() 打点确认内存是否随处理量增长;用火焰图定位转换函数是否是瓶颈。纪律是:所有流操作走 pipeline;所有外部句柄在 finally 关闭;转换函数保持纯函数、无副作用;内存占用应与数据量无关——若有关,说明某处偷偷把流读成了数组。
一句话总结:流式 I/O 的生产化 =
pipeline+error监听 +finally释放 + 常数内存——只要内存占用随数据量增长,就一定有一处把流退化成了全量读取。
延伸阅读
- 异步控制与并发治理
- 事件流与增量数据处理
- Node 后端的服务化实践
- 内存与性能剖析
- 并行计算与 Worker
- TypeScript 专题 — TypeScript 专题
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。