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" } ])
关键优化建议(针对新手)
- 索引优化:给以下字段建索引,能让聚合操作速度再翻倍:
db.secondary.createIndex({ wellId: 1 }) db.primary.createIndex({ wellId: 1 }) db.user_well_subscriptions.createIndex({ userId: 1, subscribedWellIds: 1 }) - 每日自动执行:可以用
crontab(Linux)或任务计划(Windows)定时调用mongosh执行脚本,无需手动操作。示例命令:mongosh "mongodb://your-host:port/your-db" --username your-user --password your-pass --file sync_script.js - 测试先行:先在测试环境用小批量数据验证同步和通知逻辑,确认无误后再部署到生产环境。
内容的提问来源于stack exchange,提问作者amota
相关产品推荐
相关产品推荐

