使用同一AWS S3读取流同时处理数据并上传至其他桶失败排查
问题:复用S3读取流同时上传和处理数据失败
我正尝试从第三方AWS S3存储桶读取.gz格式的文件,需要处理文件中的数据并将其上传至我方自有S3存储桶。读取文件时,我通过S3.getObject创建读取流,代码如下:
const fileStream = externalS3.getObject({Bucket: "<bucket-name>", Key: "<key>"}).createReadStream();
为提升代码效率,我计划使用同一个fileStream同时处理文件内容并上传至我方S3存储桶,但以下代码无法完成上传操作:
import Stream from "stream"; const uploadStream = fileStream.pipe(new stream.PassThrough()); const readStream = fileStream.pipe(new stream.PassThrough()); await internalS3.upload({Bucket:"<bucket-name>", Key: "<key>", Body: uploadStream}) .on("httpUploadProgress", progress => {console.log(progress)}) .on("error", error => {console.log(error)}) .promise(); readStream.pipe(createGunzip()) .on("error", err =>{console.log(err)}) .pipe(JSONStream.parse()) .on("data", data => {console.log(data)});
然而,仅保留上传逻辑的代码可成功将文件上传至我方S3存储桶:
const uploadStream = fileStream.pipe(new stream.PassThrough()); await internalS3.upload({Bucket:"<bucket-name>", Key: "<key>", Body: uploadStream}) .on("httpUploadProgress", progress => {console.log(progress)}) .on("error", error => {console.log(error)}) .promise();
注:若使用独立的fileStream分别进行上传和数据读取,功能可正常运行,但我需实现同一fileStream的复用。
解答
核心问题
- Node.js可读流的单次消费特性:可读流默认是单向流动的,数据只能被消费一次。当你将同一个
fileStream直接pipe到多个PassThrough时,数据会被第一个活跃的消费者(这里是uploadStream)优先读取,后续的readStream无法获取完整数据。 - 异步操作顺序错误:你先
await上传操作完成,此时fileStream的所有数据已经被uploadStream完全消费,后续再处理readStream时,流已经没有数据可读,自然无法完成解析。
解决方案
要实现同一个源流的复用,需要通过分流器将源流的数据复制到多个处理流中,并且同时启动所有异步处理逻辑,等待全部完成。
修改后的代码示例:
import { PassThrough } from "stream"; import { createGunzip } from "zlib"; import JSONStream from "jsonstream"; const fileStream = externalS3.getObject({Bucket: "<bucket-name>", Key: "<key>"}).createReadStream(); // 创建分流器,源流数据会先流入这个PassThrough const streamSplitter = new PassThrough(); fileStream.pipe(streamSplitter); // 启动上传任务,从分流器分支出上传流 const uploadTask = internalS3.upload({ Bucket: "<bucket-name>", Key: "<key>", Body: streamSplitter.pipe(new PassThrough()) }) .on("httpUploadProgress", progress => console.log(progress)) .on("error", error => console.log(error)) .promise(); // 启动数据解析任务,从分流器分支出解析流 const parseTask = new Promise((resolve, reject) => { streamSplitter.pipe(createGunzip()) .on("error", err => { console.log(err); reject(err); }) .pipe(JSONStream.parse()) .on("data", data => console.log(data)) .on("end", resolve); }); // 等待上传和解析任务全部完成 await Promise.all([uploadTask, parseTask]);
代码说明
- 使用
PassThrough作为分流器:源流的数据会先流入这个分流器,再从分流器分支到上传和解析两个流,确保两个流都能获取完整的源数据。 - 同时启动异步任务:避免先完成一个任务导致流被消费完毕,用
Promise.all等待两个任务都完成,保证逻辑的并行执行。 - 独立分支流:每个处理逻辑都从分流器分支出独立的流,确保各自的处理互不干扰。
内容的提问来源于stack exchange,提问作者Rachit Anand
相关产品推荐
相关产品推荐

