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

Mongo ChangeStreams超大JSON文件高效读写解析及S3上传优化咨询

优化Mongo ChangeStreams大文件读写与S3上传方案

问题核心痛点

通过Mongo ChangeStreams捕获数据库变更后,当前读写逻辑存在以下致命问题:

  • 写入时每次读取整个文件解析为数组,再重新写入,高频调用下IO和内存开销爆炸,500MB文件直接触发Error: Too long to parse
  • 上传时同步读取整个文件到内存,大文件直接超出内存限制

一、本地写入优化:JSON Lines格式+批量流式追加

核心逻辑

放弃传统JSON数组格式,改用JSON Lines(.jsonl):每行存储一个独立JSON对象,无需维护外层数组。写入时只需追加新行,无需读取整个文件;解析时可逐行处理,完全避免加载大文件到内存。同时针对高频调用增加批量写入逻辑,减少IO次数。

修改后代码

const fs = require('fs').promises;
const path = require('path');

// 批量写入缓存:key为文件路径,value为待写入的JSON字符串数组
const writeCache = new Map();
// 批量写入阈值(可根据业务调整)
const BATCH_SIZE = 50;
// 定时强制刷写缓存(防止数据积压)
const FLUSH_INTERVAL = 1000;

// 初始化定时批量刷写
setInterval(async () => {
  for (const [filePath, dataLines] of writeCache.entries()) {
    if (dataLines.length > 0) {
      await appendToJSONLFile(filePath, dataLines);
      writeCache.set(filePath, []);
    }
  }
}, FLUSH_INTERVAL);

async function exportToLocalFile(data, collectionName) {
  const currentDate = new Date();
  const year = currentDate.getFullYear();
  const month = (currentDate.getMonth() + 1).toString().padStart(2, '0');
  const day = currentDate.getDate().toString().padStart(2, '0');
  const filePath = path.join(LOCAL_AUDIT_LOGS_STORAGE_FOLDER, collectionName, `${year}-${month}-${day}.jsonl`);

  // 自动创建目录
  await fs.mkdir(path.dirname(filePath), { recursive: true });
  // 清理旧文件(替换为异步版本)
  await checkAndDeleteDayOlderFiles(collectionName);

  // 将数据加入缓存
  if (!writeCache.has(filePath)) {
    writeCache.set(filePath, []);
  }
  const cacheList = writeCache.get(filePath);
  cacheList.push(JSON.stringify(data));

  // 达到阈值立即写入
  if (cacheList.length >= BATCH_SIZE) {
    await appendToJSONLFile(filePath, cacheList);
    writeCache.set(filePath, []);
  }
}

async function appendToJSONLFile(filePath, dataLines) {
  // 异步追加写入,每行一个JSON对象
  const content = dataLines.join('\n') + '\n';
  await fs.appendFile(filePath, content, 'utf-8');
  console.log(`批量写入${dataLines.length}条数据到: ${filePath}`);
}

// 异步版旧文件清理逻辑
async function checkAndDeleteDayOlderFiles(collectionName) {
  const dirPath = path.join(LOCAL_AUDIT_LOGS_STORAGE_FOLDER, collectionName);
  try {
    const files = await fs.readdir(dirPath);
    for (const file of files) {
      const filePath = path.join(dirPath, file);
      const stats = await fs.stat(filePath);
      // 这里加入你的旧文件判断逻辑(比如超过7天删除)
      // ...
    }
  } catch (err) {
    console.error('清理旧文件失败:', err);
  }
}

二、S3上传优化:流式上传+异步并发遍历

核心逻辑

  • 用fs.createReadStream创建文件流,直接传入S3上传参数,内存占用仅为流缓冲区大小(默认几KB)
  • 改用异步目录遍历,避免同步操作阻塞事件循环
  • 控制并发上传数,防止资源耗尽

修改后代码

const fs = require('fs');
const path = require('path');
const { S3Client, PutObjectCommand } = require('@aws-sdk/client-s3');

const s3 = new S3Client({ /* 你的S3配置 */ });
const Bucket = 'your-bucket-name';
// 并发上传限制(可根据服务器配置调整)
const CONCURRENCY_LIMIT = 5;

async function readAndExportToS3(directory) {
  const files = await fs.promises.readdir(directory);
  const promises = [];

  for (const file of files) {
    const filePath = path.join(directory, file);
    const stats = await fs.promises.stat(filePath);
    
    if (stats.isDirectory()) {
      promises.push(readAndExportToS3(filePath));
    } else {
      promises.push(uploadFileToS3(filePath));
    }

    // 达到并发限制时等待完成,再继续
    if (promises.length >= CONCURRENCY_LIMIT) {
      await Promise.all(promises);
      promises.length = 0;
    }
  }

  // 处理剩余的上传任务
  if (promises.length > 0) {
    await Promise.all(promises);
  }
}

async function uploadFileToS3(filePath) {
  try {
    const key = filePath.replace(/\\/g, '/');
    const finalKey = key.split('db_audit_logs/')[1];
    
    // 创建文件流,直接传入S3
    const fileStream = fs.createReadStream(filePath);

    const uploadParams = {
      Bucket,
      Key: finalKey,
      Body: fileStream,
      ContentType: 'application/jsonl', // 匹配JSON Lines格式
    };

    const command = new PutObjectCommand(uploadParams);
    await s3.send(command);
    console.log('上传成功:', finalKey);
  } catch (err) {
    console.error('上传失败:', filePath, err);
  }
}

额外优化建议

  • 按小时拆分文件:若单日文件仍过大,可按小时生成文件(如2024-05-23-14.jsonl),进一步降低单个文件体积
  • 压缩上传:上传前用zlib对文件流进行gzip压缩,减少带宽和存储成本(需添加ContentEncoding: 'gzip'参数)
  • 错误重试:针对写入和上传异常增加重试逻辑(如使用p-retry库),避免数据丢失
  • 资源监控:添加内存、磁盘占用监控,达到阈值时触发告警

内容的提问来源于stack exchange,提问作者Ajay Dhiman

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 12:37:21