Node.js:解决Tee流单分支消费引发的内存占用过高问题
我使用S3对象存储数据,会定期对存储的Blob进行更新。为降低内存占用,采用流的方式读取并回写数据,同时应用更新。但部分更新会使Blob体积变大,而S3写入操作必须指定ContentLength,因此需要先获取转换后流的长度。于是对ReadableStream执行tee操作,通过第一个分支计算目标流长度,再对第二个分支转换后写入:
let doc = await db.send(new GetObjectCommand({Bucket: 'things', Key: key})); let stms = Readable.toWeb(doc.Body).tee(); // Calculate the length of the destination stream let newlength = await stms[0] .pipeThrough(CalcStreamLength(updates)) .getReader() .read() ; newlength = newlength.value; //stm[0].close(); // THIS DOES BAD THINGS // Transform the stream and write to destination stms[1] = stms[1].pipeThrough(CalcNewDataset(fields, updates)); stms[1] = Readable.fromWeb(stms[1]); await context.send(new PutObjectCommand({ Bucket: 'things', Key: key, Body: stms[1], ContentLength: newlength, }));
但这导致第一个流被全部缓冲到内存中,内存占用剧增(这正是我要避免的情况)。根据MDN文档说明:“如果仅消费一个分支,整个流内容会被排入内存队列。” 请问是否有可行的解决方法?
1. 改用S3分段上传(Multipart Upload)
S3的分段上传不需要提前指定总长度,你可以把转换后的流拆成固定大小的块(比如5MB)逐个上传,最后合并为完整对象。这种方式完全绕开了预计算长度的需求,而且天然适配流式处理,不会把整个数据加载到内存中。
2. 单次遍历流同时完成计算与缓存(仅适用于小数据)
如果Blob体量不大,可以在一次流遍历过程中,一边计算转换后的总长度,一边把转换后的数据缓存到内存或临时文件。遍历完成后,直接用缓存的数据和计算好的长度执行S3写入。但这种方法本质还是会加载全量数据到内存,只适合小体积场景。
3. 自定义单遍历双逻辑的流处理
放弃原生tee,自己实现流处理逻辑:读取原始流的每个chunk时,同时执行两个操作——计算转换后chunk的长度并累加得到总长度,同时将转换后的chunk传入输出流。这样整个过程只遍历一次原始流,不会出现分支缓冲的问题。示例伪代码如下:
let doc = await db.send(new GetObjectCommand({Bucket: 'things', Key: key})); const originalStream = Readable.toWeb(doc.Body); let totalLength = 0; // 自定义转换流,同时计算长度 const transformedStream = new ReadableStream({ async start(controller) { const reader = originalStream.getReader(); try { while (true) { const { done, value } = await reader.read(); if (done) break; // 应用更新转换chunk const transformedChunk = applyUpdates(value, updates); // 累加转换后的长度 totalLength += transformedChunk.length; // 将转换后的chunk加入输出流 controller.enqueue(transformedChunk); } controller.close(); } catch (err) { controller.error(err); } finally { reader.releaseLock(); } } }); // 转换为Node.js流并上传 const nodeStream = Readable.fromWeb(transformedStream); await context.send(new PutObjectCommand({ Bucket: 'things', Key: key, Body: nodeStream, ContentLength: totalLength, }));
4. 基于更新规则预估长度(仅适用于可预测的更新)
如果你的更新逻辑是可预测的(比如固定给每条记录增加N字节),可以直接读取原始Blob的ContentLength,结合更新规则计算出目标长度,无需遍历流。比如原始长度为L,共M条记录,每条增加50字节,目标长度就是L + 50 * M。但这种方法只适用于更新逻辑完全可控的场景。
内容的提问来源于stack exchange,提问作者Jefferey Cave

