如何在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", }, }, }, }, }
关键说明
- 时间差计算:通过
$subtract计算endTime与startTime的差值(MongoDB中日期相减返回毫秒数),再用$divide除以60000转换为分钟数。 - 分阶段处理:这里拆分为两个
$set阶段分步计算,也可以合并为一个$set阶段嵌套$divide和$subtract,效果一致。 - 分组聚合:第一个
$group按用户+状态分组,累计各状态的总时长;第二个$group按用户分组,将各状态的时长整理为数组格式。
内容的提问来源于stack exchange,提问作者San G
相关产品推荐
相关产品推荐

