编程 用 TransformStream 把 web streams 拼成管线:TextDecoderStream、SSE 拆行与顺序坑

2026-10-01 00:05:55

用 TransformStream 把 web streams 拼成管线

规范与文档:WHATWG Streams Standard、MDN Streams API、MDN Using readable streams、web.dev Streams 指南、Node.js Web Streams。

为什么需要把流拼起来

fetch 拿到的 response.body 是 ReadableStream。块边界由网络决定,和业务消息边界没有关系:一次 read 可能返回半行 JSON,也可能一次带回三个 SSE 事件外加下一行的开头。

如果直接 getReader() 手写循环,等于自己维护一个状态机:解码多字节字符、缓存不完整的尾段、控制读取节奏(背压)、处理中途出错。这些正是 Streams API 已经标准化的部分。把每一步做成一个 transform stream,用 pipeThrough 串联,背压和错误传播由规范负责,业务代码只处理「完整的一行/一个事件」。

TransformStream 与 pipeThrough / pipeTo

Streams API 有三种对象:ReadableStream(数据源)、WritableStream(数据目的地)、TransformStream(算法,把写入的块变换后推出)。

组合方式是 pipe:

  • readable.pipeTo(writable):直接接到可写流;
  • readable.pipeThrough(transform):把 readable 写进 transform.writable,返回 transform.readable,所以能一直链下去。

管道在连接期间会锁定(lock)流,其它 reader 无法再获取。pipeTo / pipeThrough 的 options 里有 preventClose / preventAbort / preventCancel / signal(AbortSignal),用来调整关闭、中止、取消时另一端的默认行为。背压是端到端传播的:后面的流没就绪,信号向前传,让上游放慢读取。

TransformStream 构造:

new TransformStream([transformer[, writableStrategy[, readableStrategy]]])

transformer 里的三个钩子:

  • start(controller):初始化时调用一次;
  • transform(chunk, controller):每收到一个写入块调用,可以 controller.enqueue() 0 次或多次;不 enqueue 就是透传;
  • flush(controller):所有块写完、可写侧即将关闭时调用,用来补出后缀块。

平台内置了几个常用的 transform:TextEncoderStream / TextDecoderStream(字节与文本互转)、CompressionStream / DecompressionStream(format 参数如 'gzip'、'deflate')、WebSocketStream(把 WebSocket 与流集成,Chrome 目前唯一实现)。

组一条真实管线

以 fetch 流式响应为例:

const events = response.body
.pipeThrough(new TextDecoderStream())   // 字节 -> 文本
.pipeThrough(createLineSplitter())      // 文本 -> 一行一个事件
.pipeThrough(createJsonParser());       // 行 -> 业务对象

for await (const event of events) {
handle(event);
}

TextDecoderStream 负责跨块拼接多字节字符,UTF-8 字符被切在两个块里也不会解出乱码。

按行拆分的 TransformStream,关键在跨块缓冲:

function createLineSplitter() {
let buffer = '';
return new TransformStream({
transform(chunk, controller) {
buffer += chunk;
const lines = buffer.split('\n');
buffer = lines.pop();        // 尾部可能不完整,留到下一块
for (const line of lines) {
controller.enqueue(line);
}
},
flush(controller) {
if (buffer) controller.enqueue(buffer);  // 没有换行结尾的最后一行
},
});
}

空行、SSE 的 data: 前缀、: 注释行,都留给下游的业务 transform 处理。

顺序是固定的:字节先进 TextDecoderStream,再按文本拆分;如果载荷本身是压缩数据,则是「先解压、再解码、再拆行」。顺序颠倒会得到乱码或解析失败。

三个顺序/边界坑

1. 浏览器已经解压过,别再套 DecompressionStream

fetch 遇到 Content-Encoding: gzip/br/deflate 会透明解压,response.body 里拿到的已经是解压后的字节。此时再套一层 DecompressionStream('gzip') 会二次解压并报错。

DecompressionStream 只适用于载荷本身就是压缩数据的场景:比如单独下载一个 .gz 文件,或自定义协议里内嵌了压缩块。

2. SSE / NDJSON 拆行必须跨块缓冲

块边界可能落在任何位置,不能假设「一块就是一行」,一行也可能横跨三块。做法就是上面那段:按 \n 切分,把不完整的尾部留在 buffer,等下一块拼上再切。flush 阶段还要处理没有换行结尾的最后一行。

3. pipeThrough 会锁流,错误沿管道传播

一旦进入管道,原流就被锁定,不能再 getReader();需要保留原始 readable,得先 tee()。管道中任意一环抛错,错误会顺链传播,下游读操作 reject。想要「只断当前环节、不牵连整条链」,用 preventClose / preventAbort / preventCancel 和 AbortSignal 控制兜底。另外 pipeThrough 返回的是新流,原来那个 readable 已经不能用了。

什么时候别用管线

  • 响应很小(几百 KB 以内):res.json() / res.text() 更简单,加管线只增加一层调试成本。
  • 需要两次遍历同一份数据,例如先算总长再渲染进度。流只能消费一次,要么 tee() 出两份,要么先落盘。
  • 后端没有稳定分隔符时,流式解析很难做,不如等整包回来再解析。

小结

管线的价值是把「解码 → 拆帧 → 解析」拆成可组合的 transform,把背压和错误传播交给规范。代价是流会被锁定、顺序敏感、跨块状态要自己维护。动手前先确认:数据是否已经解压、消息边界靠什么切分、这条链上的错误由哪一层兜底。

推荐文章

程序员茄子在线接单