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

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
}

关键改动说明

  1. 替换Aggregate为UpdateMany:UpdateMany是专门用于批量更新的方法,会直接修改数据库里的文档,返回实际修改的记录数。
  2. 聚合逻辑转为匹配条件:把原来聚合里的$lookup和$match逻辑整合到UpdateMany的匹配条件中,用$expr和$size判断关联的shipment_quotes是否为空。
  3. 使用MongoDB服务器时间:用bson.M{"$date": time.Now()}让MongoDB使用服务器时间,避免本地机器时间与数据库服务器时间不一致的问题。
  4. 移除不必要的事务:单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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 21:01:04