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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 18:02:36