MongoDB存在则增量更新否则插入:代码生成重复文档问题
排查MongoDB Upsert生成重复文档的问题
看起来你遇到了MongoDB upsert操作意外生成重复文档的问题,我来帮你拆解可能的原因并给出解决办法:
可能的诱因分析
1. 时间戳转换逻辑有偏差
你用moment.unix(timestamp).set('seconds', 0).toDate()处理时间戳,但这里有两个容易踩坑的点:
- 时间戳单位不匹配:
moment.unix()只接受秒级时间戳,如果你的timestamp变量是毫秒级的,转换后的Date会完全偏离预期,每次处理出来的时间都不一样,自然会触发新文档插入。 - 时区/精度差异:即使单位正确,
set('seconds', 0)没有清零毫秒部分,不同请求的时间戳可能因为毫秒差异生成不同的Date对象;另外服务器和客户端时区不一致的话,也会导致转换后的时间不匹配。
2. 缺少唯一复合索引
MongoDB的upsert在高并发场景下,如果查询条件对应的字段没有唯一索引,会出现竞态条件:多个请求同时判定文档不存在,然后各自插入新文档,最终导致重复数据。
3. 字段类型不匹配
查询条件中的itemId如果和数据库中存储的类型不一致(比如数据库是字符串,你传了数字;或者反过来),MongoDB的类型敏感查询会找不到已有文档,从而插入新的重复项。
分步解决方案
第一步:修正时间戳转换逻辑
先确保时间戳处理完全一致:
- 如果
timestamp是毫秒级,把moment.unix(timestamp)改成moment(timestamp),同时清零毫秒部分:// 修正后的时间戳处理(毫秒级场景) const normalizedTimestamp = moment(timestamp) .set('seconds', 0) .milliseconds(0) // 确保毫秒也清零,消除细微差异 .toDate(); - 如果
timestamp是秒级,保留moment.unix(timestamp),但同样要清零毫秒:// 秒级时间戳的处理方式 const normalizedTimestamp = moment.unix(timestamp) .set('seconds', 0) .milliseconds(0) .toDate();
可以在代码里打印normalizedTimestamp,确认相同业务时间的转换结果完全一致。
第二步:创建唯一复合索引
为了从根本上阻止重复插入,同时解决并发竞态问题,给你的集合创建itemId和timestamp的复合唯一索引:
// 在集合初始化时执行一次(或直接在MongoDB Shell中运行) db.yourCollectionName.createIndex( { itemId: 1, timestamp: 1 }, { unique: true } );
这个索引会强制itemId + timestamp的组合唯一,即使有并发请求,MongoDB也会阻止重复插入,同时upsert能正确匹配已有文档进行更新。
第三步:校验字段类型一致性
检查itemId的类型:
- 如果数据库中
itemId是ObjectId,确保代码中传入的是ObjectId类型,而非字符串; - 如果是字符串/数字,确保两端类型完全匹配,避免类型不匹配导致查询失效。
优化后的完整代码示例
// 先标准化时间戳 const normalizedTimestamp = moment(timestamp) // 秒级的话替换成moment.unix(timestamp) .set('seconds', 0) .milliseconds(0) .toDate(); Collection.updateOne( // 推荐用updateOne替代update,行为更明确 { itemId: itemId, timestamp: normalizedTimestamp, }, { '$inc': { 'data.ph1': data.ph1, 'data.ph2': data.ph2, 'data.ph3': data.ph3, 'data.total': data.total }, }, { upsert: true }, (err, resp) => { if(err) { // 捕获唯一索引冲突错误,重试更新 if(err.code === 11000) { Collection.updateOne( { itemId: itemId, timestamp: normalizedTimestamp }, { '$inc': { 'data.ph1': data.ph1, 'data.ph2': data.ph2, 'data.ph3': data.ph3, 'data.total': data.total } }, (retryErr) => retryErr ? reject(retryErr) : resolve('minute value incremented' + gatewayId) ); } else { reject(err); } } else { resolve('minute value incremented' + gatewayId ); } } );
额外提示
- 推荐使用
updateOne替代旧的update方法,它明确保证只更新一个文档,更符合你的业务需求; - 如果日志中出现
E11000 duplicate key error,说明唯一索引已经生效,这时候的重试逻辑可以处理并发场景下的冲突。
内容的提问来源于stack exchange,提问作者Bazinga777
相关产品推荐
相关产品推荐

