Mongoose cursor.eachAsync内存泄漏及SQS批量发送方案咨询
解决Mongoose迭代内存泄漏问题
问题根源
原代码使用eachAsync并设置parallel:10,当文档量较大时,并行执行的异步任务(SQS发送+DB更新)可能堆积大量未完成的Promise,加上cursor未被正确清理,导致内存占用持续上升。
优化方案:手动迭代+可控并发
改用for await...of配合并发限制工具(如p-limit),既能保证处理效率,又能严格控制内存占用:
步骤1:安装并发控制依赖
npm install p-limit
步骤2:优化迭代代码
import pLimit from 'p-limit'; // 限制并发数为10,可根据实际场景调整 const limit = pLimit(10); async function processDocuments() { const campaignUserCursor = this.campaignUserModel .aggregate([...query]) .allowDiskUse(true) .cursor({ batchSize: 10 }) .exec(); try { const tasks = []; for await (const campaignUser of campaignUserCursor) { // 将每个处理任务加入并发队列 tasks.push(limit(async () => { const queueObj = this.setObj(campaignUser); await this.sqsProducerService.sendMessage(queueObj); await this.broadcastRepository.updateBroadcast(queueObj, { $inc: { total_queue_count: 1 } }); })); } // 等待所有任务完成 await Promise.all(tasks); } finally { // 强制关闭cursor,释放数据库资源 await campaignUserCursor.close(); } }
关键优化点
- 手动迭代+并发限制:用
for await...of遍历cursor,配合p-limit控制同时运行的异步任务数,避免Promise堆积。 - 强制释放资源:在
finally块中关闭cursor,确保无论成功还是出错,都能释放数据库连接资源。 - 透明可控:相比
eachAsync的黑盒实现,手动迭代更易排查内存泄漏点。
使用@ssut/nestjs-sqs发送批量消息
@ssut/nestjs-sqs提供sendMessageBatch方法实现批量发送,需构造符合AWS SQS规范的消息数组:
批量发送基础代码示例
async function sendBatchToSQS(documents: any[]) { // 转换为SQS批量消息格式,每个消息需唯一Id const batchMessages = documents.map((doc, index) => ({ Id: `msg-${index}-${Date.now()}`, // 保证Id全局唯一 MessageBody: JSON.stringify(this.setObj(doc)), // 可选:添加DelaySeconds、MessageAttributes等参数 })); // 调用批量发送方法 await this.sqsProducerService.sendMessageBatch( 'your-queue-name', // 目标队列名称 batchMessages ); }
结合Mongoose迭代的批量优化方案
为减少SQS请求次数,可将cursor返回的文档按批次打包发送:
import pLimit from 'p-limit'; const limit = pLimit(5); // 控制批量任务的并发数 const SQS_BATCH_SIZE = 10; // 每批次最多发送10条消息(AWS SQS限制) async function processAndSendBatch() { const campaignUserCursor = this.campaignUserModel .aggregate([...query]) .allowDiskUse(true) .cursor({ batchSize: SQS_BATCH_SIZE }) .exec(); try { const tasks = []; let currentBatch = []; for await (const campaignUser of campaignUserCursor) { currentBatch.push(campaignUser); // 批次满额时触发批量处理 if (currentBatch.length === SQS_BATCH_SIZE) { tasks.push(limit(async () => { // 批量发送到SQS await sendBatchToSQS(currentBatch); // 批量更新DB计数 const updatePromises = currentBatch.map(doc => { const queueObj = this.setObj(doc); return this.broadcastRepository.updateBroadcast(queueObj, { $inc: { total_queue_count: 1 } }); }); await Promise.all(updatePromises); })); currentBatch = []; } } // 处理剩余的不足一批的文档 if (currentBatch.length > 0) { tasks.push(limit(async () => { await sendBatchToSQS(currentBatch); const updatePromises = currentBatch.map(doc => { const queueObj = this.setObj(doc); return this.broadcastRepository.updateBroadcast(queueObj, { $inc: { total_queue_count: 1 } }); }); await Promise.all(updatePromises); })); } await Promise.all(tasks); } finally { await campaignUserCursor.close(); } }
批量发送注意事项
- 消息Id唯一性:每个批量消息的
Id必须唯一,否则SQS会拒绝整个批次。 - AWS限制:每个批量请求最多包含10条消息,总大小不超过256KB,需根据消息大小调整批次容量。
- 错误处理:批量发送若部分消息失败,SQS会返回失败列表,需根据业务需求添加重试逻辑。
内容的提问来源于stack exchange,提问作者joe
相关产品推荐
相关产品推荐

