使用@aws-sdk/lib-storage+JSONStream将MongoDB JSON流上传S3的报错排查
核心问题说明
你遇到的所有错误主要来自两个共性问题:
- 你所用的
@aws-sdk/lib-storage@3.34.0属于v3早期版本,对Node.js原生可读流、Mongo Cursor返回的流兼容性不足 - 方案3、5中
uploadStreamFile是异步函数,调用后返回的是Promise而非流实例,因此pipe时会报dest.on is not a function错误
正确实现代码
import { MongoClient } from 'mongodb'; import { S3Client } from '@aws-sdk/client-s3'; import { Upload } from '@aws-sdk/lib-storage'; import { PassThrough } from 'stream'; import { env } from '../../../env'; const s3Client = new S3Client({ region: env.AWS_REGION }); export const uploadMongoStreamToS3 = async (connectionString, collectionName) => { let client; try { client = await MongoClient.connect(connectionString); const db = client.db(); // 创建中转流,解决旧版SDK的流兼容问题 const passThrough = new PassThrough(); // 初始化S3分片上传任务 const upload = new Upload({ client: s3Client, params: { Bucket: 'test-bucket', Key: 'extracted-data/benda_mongo.json', Body: passThrough, ContentType: 'application/json' }, }); // 异步触发上传,不要阻塞流的管道配置 const uploadTask = upload.done(); // 获取Mongo查询流,直接序列化文档为JSON字符串 const readStream = db.collection(collectionName) .find({}) .limit(5) .stream({ transform: doc => JSON.stringify(doc) + '\n' }); // 流管道连接 readStream.pipe(passThrough); // 等待上传完成 await uploadTask; } catch (err) { log.error('上传失败', err); throw err.name; } finally { if (client) { await client.close(); } } };
如果你需要用JSONStream输出标准JSON数组格式,替换流配置部分即可:
import JSONStream from 'JSONStream'; // ... 其余代码不变 const readStream = db.collection(collectionName).find({}).limit(5).stream(); readStream.pipe(JSONStream.stringify()).pipe(passThrough);
各失败方案的错误原因
- 方案1:直接将未序列化的Mongo对象流传给S3上传的Body参数,SDK无法识别对象类型的流数据
- 方案2:旧版SDK对
JSONStream返回的转换流兼容性不足,无法识别为合法的可读流 - 方案3、5:
uploadStreamFile是异步函数,调用后返回Promise而非流实例,pipe方法无法接收Promise作为目标参数 - 方案4:旧版SDK无法直接识别Mongo Cursor返回的转换流,需要加一层
PassThrough中转
内容的提问来源于stack exchange,提问作者Oron Bendavid
相关产品推荐
相关产品推荐

