优化AWS DocumentDB中$in查询与批量更新的执行效率
优化AWS DocumentDB中批量更新已处理任务状态的方案
问题背景
现有两个MongoDB兼容集合:
sourcequeuemanualupload:存储待处理/已分配的上传数据,状态为Unworked或Assigned,示例Assigned记录:
{ "_id" : ObjectId("63e0e46a6047d75b9c20d8ec"), "Properties: Name" : "Hangman - Guess Words", "Appstore URL" : "https://itunes.apple.com/app/id1375993101?hl=None", "Region" : "na", "Create Date" : "na", "AHT" : "1", "sourceId" : "63e0e3719b4f812ba5333a31", "type" : "Manual", "uploadTime" : "2023-02-06T11:28:42.533+0000", "status" : "Assigned", "batchId" : "63e0e3719b4f812ba5333a31_746f22e4319b4d81b8ab255f5e653c2c_612023112842" }
queuedata:存储已完成的任务操作记录,状态均为Completed,通过id字段关联sourcequeuemanualupload的_id,示例Completed记录:
{ "_id" : ObjectId("63e0e4b19b4f812ba5333a34"), "templateId" : "63e0e28e9b4f812ba5333a30", "id" : "63e0e46a6047d75b9c20d8ec", "moderator" : "kodaga", "startTime" : "2023-02-06T11:29:46.048Z", "endTime" : "2023-02-06T11:29:52.438Z", "status" : "Completed", "AHT" : NumberLong(6), "userInput" : [ { "question" : "Is the URL leading to the desired store page link?", "response" : "yes" }, { "question" : "Comments, if any.", "response" : "test 1" } ] }
因操作失误,已在queuedata中标记为Completed的任务,未在sourcequeuemanualupload中将对应Assigned状态更新为Completed。当前数据规模:
> db.sourcequeuemanualupload.count() 414781 > db.sourcequeuemanualupload.count({"status":"Assigned"}) 306418 > db.queuedata.count() 298128
原方案先拉取所有Assigned记录的_id数组,再通过$in查询queuedata匹配记录,最后批量更新,因单次拉取数据量过大,导致执行超时甚至会话中断,即使添加索引也无明显改善。
优化方案(适配AWS DocumentDB 4.0.0)
1. 先创建必要索引
创建复合索引加速筛选和关联查询:
// 加速source集合中Assigned状态记录的筛选 db.sourcequeuemanualupload.createIndex({ status: 1, _id: 1 }) // 加速queuedata中id字段的查询 db.queuedata.createIndex({ id: 1 })
2. 分批次分页处理(推荐)
通过游标分页分批次获取queuedata中的关联ID,再批量更新source集合,避免一次性加载大量数据导致内存溢出或超时:
const batchSize = 1000; // 可根据DocumentDB实例性能调整批次大小 let lastProcessedId = null; while (true) { // 分页查询queuedata的id,按id排序实现游标分页 const query = lastProcessedId ? { id: { $gt: lastProcessedId } } : {}; const cursor = db.queuedata.find(query, { id: 1, _id: 0 }).sort({ id: 1 }).limit(batchSize); const targetIds = []; while (cursor.hasNext()) { const doc = cursor.next(); targetIds.push(ObjectId(doc.id)); // 转换为ObjectId匹配source集合的_id类型 lastProcessedId = doc.id; } if (targetIds.length === 0) break; // 批量更新符合条件的记录 const updateResult = db.sourcequeuemanualupload.updateMany( { status: "Assigned", _id: { $in: targetIds } }, { $set: { status: "Completed" } } ); print(`完成批次更新:处理 ${updateResult.modifiedCount} 条记录,当前批次最后ID:${lastProcessedId}`); }
3. 聚合管道关联筛选(可选)
利用DocumentDB支持的$lookup聚合操作,直接关联两个集合筛选出需要更新的记录,再分批次更新:
const batchSize = 1000; // 聚合获取所有已处理但状态仍为Assigned的记录ID const pendingUpdateIds = db.sourcequeuemanualupload.aggregate([ { $match: { status: "Assigned" } }, { $lookup: { from: "queuedata", localField: "_id", foreignField: "id", as: "processedInfo" } }, { $match: { processedInfo: { $ne: [] } } }, // 筛选已存在于queuedata的记录 { $project: { _id: 1 } } ]).toArray().map(item => item._id); // 分批次执行更新 for (let i = 0; i < pendingUpdateIds.length; i += batchSize) { const batchIds = pendingUpdateIds.slice(i, i + batchSize); db.sourcequeuemanualupload.updateMany( { _id: { $in: batchIds }, status: "Assigned" }, { $set: { status: "Completed" } } ); print(`完成第 ${Math.floor(i/batchSize)+1} 批次更新,处理 ${batchIds.length} 条记录`); }
方案说明
原方案低效的核心原因是一次性加载数十万级别的ID数组,导致内存占用过高、查询性能骤降。分批次处理通过控制单次查询和更新的数据量,降低内存压力,同时利用索引加速筛选,大幅提升执行效率。游标分页方式无需维护全局数组,适合大规模数据处理场景。
内容的提问来源于stack exchange,提问作者coding_newbie
相关产品推荐
相关产品推荐

