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
相关产品推荐
相关产品推荐

