TS 流式 I/O:Node Streams 类型体系、背压与大文件处理

系统讲解 TypeScript 流式 I/O:Node Streams 的 Readable/Writable/Transform/Duplex 类型体系、背压与 pipeline、大文件读写、CSV 与 JSON 流式解析、类型化 Transform、Web Streams 与 Node 的互操作、错误处理与资源释放,以及内存与吞吐调优。

引言

把 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 互操作、错误与资源释放,最后给出内存与吞吐调优。

前置:异步控制、Node 后端、事件与流。


目录


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 的差别

维度pipepipeline
错误传播不传播传播并销毁
资源清理需手动自动
返回值目标流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 决策表

需求做法避免
读大文件createReadStreamreadFile
逐行处理readline手动 split
多步转换pipeline链式 .pipe()
超大 JSONstream-jsonJSON.parse
超大 CSVcsv-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 释放 + 常数内存——只要内存占用随数据量增长,就一定有一处把流退化成了全量读取。


延伸阅读

继续阅读

探索更多技术文章

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

全部文章 返回首页

「typescript」更多文章

  1. TS 缓存策略与类型安全:层次、失效、防护与一致性取舍
  2. TS GraphQL 服务端类型安全:codegen、Resolver 与 DataLoader 实践
  3. TS 边缘框架 Hono:Web 标准、端到端类型安全与多运行时部署