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

MongoDB集合对比脚本优化求助:提升39k文档处理效率

用MongoDB查询语言(MQL)优化油井数据同步与通知生成流程

直接用MongoDB原生的聚合管道和批量操作替代JS循环,能大幅提升处理效率——毕竟这些操作是在数据库端批量执行,避免了客户端与服务器之间的多次数据往返,处理39k文档的时间能压缩到分钟级。下面分步骤给出具体实现方案,适合新手快速上手:


第一步:批量同步secondary到primary集合(处理新增+属性变更)

用$merge操作可以一次性完成「匹配更新、无匹配插入」的逻辑,无需循环遍历每一条文档。假设油井的唯一标识是wellId(请根据实际字段替换):

db.secondary.aggregate([
  // 只保留需要同步的核心字段,减少数据传输
  {
    $project: {
      wellId: 1,
      statusAbbr: 1,
      operator: 1,
      wellName: 1, // 保留油井名称用于后续生成通知,按需添加
      _id: 0 // 不需要副集合的_id,主集合会自动生成或保留原有_id
    }
  },
  // 同步到主集合
  {
    $merge: {
      into: "primary", // 目标主集合名
      on: "wellId", // 用于匹配文档的唯一键
      whenMatched: [
        // 匹配到已有文档时,仅更新指定变更字段+记录更新时间
        {
          $set: {
            statusAbbr: "$$new.statusAbbr",
            operator: "$$new.operator",
            lastSyncedAt: new Date()
          }
        }
      ],
      whenNotMatched: "insert" // 无匹配时直接插入新文档
    }
  }
])

第二步:批量识别变更/新增文档,生成通知原始数据

通过聚合管道关联主副集合,批量筛选出需要通知的文档,并生成标准化的通知条目:

// 先将筛选出的通知数据存入临时集合(方便后续按用户合并)
db.secondary.aggregate([
  // 关联主集合的对应文档
  {
    $lookup: {
      from: "primary",
      localField: "wellId",
      foreignField: "wellId",
      as: "primaryDoc"
    }
  },
  // 拆分主集合的匹配结果(兼容新增文档的空值情况)
  { $unwind: { path: "$primaryDoc", preserveNullAndEmptyArrays: true } },
  // 判断文档类型:新增/状态变更/运营商变更
  {
    $addFields: {
      isNew: { $eq: ["$primaryDoc", null] },
      statusChanged: { $ne: ["$statusAbbr", "$primaryDoc.statusAbbr"] },
      operatorChanged: { $ne: ["$operator", "$primaryDoc.operator"] }
    }
  },
  // 筛选出需要通知的文档
  {
    $match: {
      $or: [
        { isNew: true },
        { statusChanged: true },
        { operatorChanged: true }
      ]
    }
  },
  // 生成标准化通知内容
  {
    $project: {
      wellId: 1,
      wellName: "$wellName",
      notificationType: {
        $cond: {
          if: "$isNew",
          then: "油井新增",
          else: {
            $cond: {
              if: { $and: ["$statusChanged", "$operatorChanged"] },
              then: "状态+运营商变更",
              else: {
                $cond: { if: "$statusChanged", then: "状态变更", else: "运营商变更" }
              }
            }
          }
        }
      },
      oldStatus: "$primaryDoc.statusAbbr",
      newStatus: "$statusAbbr",
      oldOperator: "$primaryDoc.operator",
      newOperator: "$operator",
      notifyTime: new Date()
    }
  },
  // 将结果存入临时集合,方便后续按用户合并
  { $out: "temp_notification_events" }
])

第三步:按用户合并通知,减少推送次数

假设用户订阅油井的集合为user_well_subscriptions(结构示例:{ userId: "xxx", subscribedWellIds: ["well001", "well002"] }),通过聚合关联订阅数据与临时通知集合,批量按用户合并通知:

db.user_well_subscriptions.aggregate([
  // 拆分用户订阅的油井ID,便于关联通知
  { $unwind: "$subscribedWellIds" },
  // 关联临时通知集合的对应条目
  {
    $lookup: {
      from: "temp_notification_events",
      localField: "subscribedWellIds",
      foreignField: "wellId",
      as: "user_notifications"
    }
  },
  // 过滤掉无通知的用户
  { $match: { user_notifications: { $ne: [] } } },
  // 拆分通知条目
  { $unwind: "$user_notifications" },
  // 按用户ID分组,合并所有通知
  {
    $group: {
      _id: "$userId",
      notifications: { $push: "$user_notifications" },
      totalCount: { $sum: 1 }
    }
  },
  // 生成最终推送格式
  {
    $project: {
      userId: "$_id",
      notifications: 1,
      totalNotifications: "$totalCount",
      pushTime: new Date()
    }
  },
  // 可选:将最终推送数据存入集合,方便后续调用推送接口
  { $out: "user_push_tasks" }
])

关键优化建议(针对新手)

  1. 索引优化:给以下字段建索引,能让聚合操作速度再翻倍:
    db.secondary.createIndex({ wellId: 1 })
    db.primary.createIndex({ wellId: 1 })
    db.user_well_subscriptions.createIndex({ userId: 1, subscribedWellIds: 1 })
    
  2. 每日自动执行:可以用crontab(Linux)或任务计划(Windows)定时调用mongosh执行脚本,无需手动操作。示例命令:
    mongosh "mongodb://your-host:port/your-db" --username your-user --password your-pass --file sync_script.js
    
  3. 测试先行:先在测试环境用小批量数据验证同步和通知逻辑,确认无误后再部署到生产环境。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 18:01:19