用 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,把背压和错误传播交给规范。代价是流会被锁定、顺序敏感、跨块状态要自己维护。动手前先确认:数据是否已经解压、消息边界靠什么切分、这条链上的错误由哪一层兜底。