Go中MongoDB聚合更新未持久化问题求助
问题根源
你代码里的核心错误是用Aggregate聚合管道执行更新操作——Aggregate是只读的查询操作,只会返回匹配的文档,不会对数据库里的文档做任何修改。所以不管你怎么提交事务,数据库里的status字段都不会变,每次查询返回的数量自然也不会减少。
另外,你在事务里执行只读的Aggregate操作完全没必要,事务应该用于需要原子性的写操作场景。
修复方案
改用UpdateMany方法,结合匹配条件来批量更新符合要求的文档。同时可以把原来聚合里的逻辑转化为UpdateMany的匹配条件,或者用管道式更新(MongoDB 4.2+支持)。
修改后的代码
func CheckShipmentExpiryDates(c *mongo.Client) (int, error) { coll := c.Database(os.Getenv("DATABASE")).Collection("shipments") // 用MongoDB的当前时间替代本地time.Now(),避免机器时间偏差 now := bson.M{"$date": time.Now()} // 构建匹配条件,对应原来聚合里的过滤逻辑 matchStage := bson.M{ "expiration_date": bson.M{ "$exists": true, "$type": 9, // 对应Date类型 "$lt": now, }, "status": bson.M{"$ne": "EXPIRED"}, // 用$expr关联子查询,检查shipment_quotes里没有WON状态的记录 "$expr": bson.M{ "$eq": []interface{}{ bson.M{ "$size": bson.M{ "$lookup": bson.M{ "from": "shipment_quotes", "let": bson.M{"shipmentID": "$_id"}, "pipeline": []bson.M{ {"$match": bson.M{"$expr": bson.M{"$and": []bson.M{ {"$eq": []string{"$shipment_id", "$$shipmentID"}}, {"$eq": []string{"$status", "WON"}}, }}}}, }, "as": "quotes", }, }, }, 0, }, }, } update := bson.M{ "$set": bson.M{ "status": "EXPIRED", "updated_at": now, // 同样用MongoDB时间 }, } // 执行批量更新,单UpdateMany操作本身具备原子性,无需额外事务 result, err := coll.UpdateMany(context.TODO(), matchStage, update) if err != nil { return 0, err } return int(result.ModifiedCount), nil }
关键改动说明
- 替换
Aggregate为UpdateMany:UpdateMany是专门用于批量更新的方法,会直接修改数据库里的文档,返回实际修改的记录数。 - 聚合逻辑转为匹配条件:把原来聚合里的
$lookup和$match逻辑整合到UpdateMany的匹配条件中,用$expr和$size判断关联的shipment_quotes是否为空。 - 使用MongoDB服务器时间:用
bson.M{"$date": time.Now()}让MongoDB使用服务器时间,避免本地机器时间与数据库服务器时间不一致的问题。 - 移除不必要的事务:单
UpdateMany操作本身就是原子性的,不需要额外开启事务,除非你有多个写操作需要原子执行。
额外优化建议
如果你的MongoDB版本支持(4.2+),也可以用管道式更新,逻辑更贴近你原来的聚合思路:
func CheckShipmentExpiryDates(c *mongo.Client) (int, error) { coll := c.Database(os.Getenv("DATABASE")).Collection("shipments") now := bson.M{"$date": time.Now()} updatePipeline := []bson.M{ {"$lookup": bson.M{ "from": "shipment_quotes", "let": bson.M{"shipmentID": "$_id"}, "pipeline": []bson.M{ {"$match": bson.M{"$expr": bson.M{"$and": []bson.M{ {"$eq": []string{"$shipment_id", "$$shipmentID"}}, {"$eq": []string{"$status", "WON"}}, }}}}, }, "as": "quotes", }}, {"$match": bson.M{ "expiration_date": bson.M{"$exists": true, "$type": 9, "$lt": now}, "status": bson.M{"$ne": "EXPIRED"}, "$expr": bson.M{"$eq": []interface{}{bson.M{"$size": "$quotes"}, 0}}, }}, {"$set": bson.M{"status": "EXPIRED", "updated_at": now}}, } result, err := coll.UpdateMany(context.TODO(), bson.M{}, updatePipeline) if err != nil { return 0, err } return int(result.ModifiedCount), nil }
验证方法
修改后执行定时任务,查看MongoDB Compass里的文档status是否变为EXPIRED,同时日志里的expiredShipments数量应该每次执行后减少(直到没有符合条件的文档)。
内容的提问来源于stack exchange,提问作者Georgescu Rares Daniel
相关产品推荐
相关产品推荐

