You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用@aws-sdk/lib-storage+JSONStream将MongoDB JSON流上传S3的报错排查

核心问题说明

你遇到的所有错误主要来自两个共性问题:

  1. 你所用的@aws-sdk/lib-storage@3.34.0属于v3早期版本,对Node.js原生可读流、Mongo Cursor返回的流兼容性不足
  2. 方案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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.01 23:36:03