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

如何在Golang中实现含$divide与$subtract的MongoDB聚合管道

在Golang中实现MongoDB聚合Pipeline

需求:基于MongoDB 4.2版本,编写对应指定原生聚合查询的Golang Pipeline,已完成$match和$unwind阶段,需实现$divide、$subtract等操作以匹配原生查询结果。


MongoDB原生聚合查询

db.getCollection("db").aggregate(
    [
        { 
            "$match" : { 
                "attendanceDate" : "07/26/2022"
            }
        }, 
        { 
            "$unwind" : "$history"
        }, 
        { 
            "$set" : { 
                "timeDiff" : { 
                    "$divide" : [
                        { 
                            "$subtract" : [
                                "$history.endTime", 
                                "$history.startTime"
                            ]
                        }, 
                        60000.0
                    ]
                }
            }
        }, 
        { 
            "$group" : { 
                "_id" : { 
                    "status" : "$history.status", 
                    "displayName" : "$displayName"
                }, 
                "duration" : { 
                    "$sum" : "$timeDiff"
                }
            }
        }, 
        { 
            "$group" : { 
                "_id" : "$_id.displayName", 
                "durations" : { 
                    "$push" : { 
                        "key" : "$_id.status", 
                        "value" : "$duration"
                    }
                }
            }
        }
    ]
)

MongoDB样例文档

{ 
    "_id" : ObjectId("62e01543666e8a64c2aeec56"), 
    "attendanceDate" : "07/26/2022", 
    "displayName" : "John, Doe", 
    "signInDate" : ISODate("2022-07-26T16:24:35.488+0000"), 
    "currentStatus" : "Other", 
    "currentStatusTime" : ISODate("2022-07-26T16:37:54.890+0000"), 
    "history" : [
        {
            "status" : "Other", 
            "startTime" : ISODate("2022-07-26T16:37:54.890+0000")
        }, 
        {
            "status" : "In", 
            "startTime" : ISODate("2022-07-26T16:33:00.655+0000"), 
            "endTime" : ISODate("2022-07-26T16:37:54.890+0000")
        }, 
        {
            "status" : "Training", 
            "startTime" : ISODate("2022-07-26T16:32:01.337+0000"), 
            "endTime" : ISODate("2022-07-26T16:33:00.657+0000")
        }, 
        {
            "status" : "In", 
            "startTime" : ISODate("2022-07-26T16:31:00.764+0000"), 
            "endTime" : ISODate("2022-07-26T16:32:01.338+0000")
        }, 
        {
            "status" : "Lunch", 
            "startTime" : ISODate("2022-07-26T16:30:01.025+0000"), 
            "endTime" : ISODate("2022-07-26T16:31:00.765+0000")
        }, 
        {
            "status" : "In", 
            "startTime" : ISODate("2022-07-26T16:27:33.789+0000"), 
            "endTime" : ISODate("2022-07-26T16:30:01.026+0000")
        }, 
        {
            "status" : "Break", 
            "startTime" : ISODate("2022-07-26T16:25:38.492+0000"), 
            "endTime" : ISODate("2022-07-26T16:27:33.789+0000")
        }, 
        {
            "status" : "In", 
            "startTime" : ISODate("2022-07-26T16:24:41.753+0000"), 
            "endTime" : ISODate("2022-07-26T16:25:38.493+0000")
        }
    ]
}

聚合查询预期结果

{ 
    "_id" : "John, Doe", 
    "durations" : [
        {
            "key" : "Other", 
            "value" : NumberInt(0)
        }, 
        {
            "key" : "Lunch", 
            "value" : 0.9956666666666667
        }, 
        {
            "key" : "In", 
            "value" : 9.3131
        }, 
        {
            "key" : "Training", 
            "value" : 0.9886666666666667
        }, 
        {
            "key" : "Break", 
            "value" : 1.9216166666666668
        }
    ]
}

初始Golang代码

type Attendance struct {
    ID                primitive.ObjectID `json:"_id" bson:"_id"`
    Date              time.Time          `json:"signInDate" bson:"signInDate"`
    DisplayName       string             `json:"displayName" bson:"displayName"`
    CurrentStatus     string             `json:"currentStatus,omitempty" bson:"currentStatus,omitempty"`
    CurrentStatusTime time.Time          `json:"currentStatusTime,omitempty" bson:"currentStatusTime,omitempty"`
    History           []AttendanceHistoryItem `json:"history" bson:"history"`
}

type AttendanceHistoryItem struct {
    Status       string    `json:"status,omitempty" bson:"status,omitempty"`
    StartTime    time.Time `json:"startTime,omitempty" bson:"startTime,omitempty"`
    EndTime      time.Time `json:"endTime,omitempty" bson:"endTime,omitempty"`
}

func (r *repo) Find() ([]domain.Attendance, error) {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    var attendances []domain.Attendance

    pipeline := []bson.M{
        {
            "$match": bson.M{
                "attendanceDate": "07/26/2022",
            },
        },
        {
            "$unwind": "$history",
        },
    }

    cur, err := r.Db.Collection("db").Aggregate(ctx, pipeline)
    defer cur.Close(ctx)
    if err != nil {
        return attendances, err
    }

    for cur.Next(ctx) {
        var attendance domain.Attendance

        err := cur.Decode(&attendance)
        if err != nil {
            return attendances, err
        }
    
        attendances = append(attendances, attendance)
    }

    if err = cur.Err(); err != nil {
        return attendances, err
    }

    return attendances, nil
}

调试通过的修改后Golang Pipeline

pipeline := []bson.M{
    {
        "$match": bson.M{
            "attendanceDate": "07/26/2022",
        },
    },
    {
        "$unwind": "$history",
    },
    {
        "$set": bson.M{
            "timeDiff": bson.M{
                "$subtract": bson.A{
                    "$history.endTime",
                    "$history.startTime",
                },
            },
        },
    },
    {
        "$set": bson.M{
            "timeDiff": bson.M{
                "$divide": bson.A{
                    "$timeDiff",
                    60000.0,
                },
            },
        },
    },
    {
        "$group": bson.M{
            "_id": bson.M{
                "status":      "$history.status",
                "displayName": "$displayName",
            },
            "duration": bson.M{
                "$sum": "$timeDiff",
            },
        },
    },
    {
        "$group": bson.M{
            "_id": "$_id.displayName",
            "durations": bson.M{
                "$push": bson.M{
                    "key":   "$_id.status",
                    "value": "$duration",
                },
            },
        },
    },
}

关键说明

  1. 时间差计算:通过$subtract计算endTime与startTime的差值(MongoDB中日期相减返回毫秒数),再用$divide除以60000转换为分钟数。
  2. 分阶段处理:这里拆分为两个$set阶段分步计算,也可以合并为一个$set阶段嵌套$divide和$subtract,效果一致。
  3. 分组聚合:第一个$group按用户+状态分组,累计各状态的总时长;第二个$group按用户分组,将各状态的时长整理为数组格式。

内容的提问来源于stack exchange,提问作者San G

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 19:07:37