14 Transform、Duplex 与数据转换

1. 四种角色

Readable   提供数据
Writable   接收数据
Duplex     同时可读可写,两侧可相对独立
Transform  特殊 Duplex,输出由输入转换得到

PassThrough 是不修改数据的 Transform,适合统计字节、观察进度或分流调试。

2. 文本 Transform

import { Transform } from "node:stream";

const upperCaseTransform = new Transform({
  transform(chunk, encoding, callback) {
    try {
      callback(null, chunk.toString("utf8").toUpperCase());
    } catch (error: unknown) {
      callback(error instanceof Error ? error : new Error("转换失败"));
    }
  },
});

callback 必须且只能调用一次。抛出的异常或传给 callback 的错误会由 Pipeline 捕获。

3. 字符边界问题

不能假设每个 Chunk 恰好结束在一个 UTF-8 字符或一行末尾:

Chunk 1: 你的一部分字节
Chunk 2: 剩余字节

简单对每块 toString() 可能破坏多字节字符。文本 Transform 应使用 StringDecoder,或者通过 setEncoding() 让 Readable 正确处理字符边界;按行解析还要保留最后一段未完成文本。

4. 进度统计

import { PassThrough } from "node:stream";

function createProgressStream(
  onBytes: (totalBytes: number) => void,
): PassThrough {
  let totalBytes = 0;
  const tracker = new PassThrough();

  tracker.on("data", (chunk: Buffer) => {
    totalBytes += chunk.length;
    onBytes(totalBytes);
  });

  return tracker;
}

进度回调必须轻量;不要在每个 Chunk 中同步写数据库。可以按时间或字节阈值节流更新。

5. gzip 实操

import { createGzip } from "node:zlib";

await pipeline(
  createReadStream(inputPath),
  createGzip(),
  createWriteStream(outputPath),
);

.gz 是单数据流压缩,不等同于能保存目录层级的 ZIP。

6. Markdown 不一定能真正流式解析

markdown-it.render() 接收完整字符串。即使输入通过 ReadStream 获得,最终仍要把单个 Markdown 聚合成字符串:

流式读取整个文件
  → 聚合字符串
  → markdown-it.render()

这不会降低单个 Markdown 的完整内存需求。正确策略是限制每个 Markdown 大小,并把图片复制、ZIP 解压和下载等适合的阶段流式化。

7. Object Mode

可以把每一行解析成对象:

new Transform({
  readableObjectMode: true,
  transform(chunk, encoding, callback) {
    callback(null, { text: chunk.toString() });
  },
});

对象模式适合结构化事件,但大量大对象仍会占用大量内存。

练习题

  1. 编写一个正确处理跨 Chunk 行尾的行计数 Transform。
  2. 使用 PassThrough 每处理 1 MiB 输出一次进度。
  3. 解释 gzip 与 ZIP 的差别。
  4. 为什么给 Markdown 文件套上 ReadStream 不等于渲染器已经流式化?