如何在MongoDB中计算数值随时间的变化及小时级增益?
MongoDB:计算单文档数值与1小时后对应分钟平均值的对比
场景与需求
流式数据以亚秒级间隔存储为单文档,需要完成两个核心操作:
- 按1分钟时间窗口分组计算数值平均值;
- 为每个原始文档匹配1小时后对应分钟窗口的平均值,并添加
valueLast(该平均值)和gainPct(增益百分比,公式:(valueLast - value)/value)字段。
举个例子:14:30的文档要匹配15:30的分钟平均值,15:35的文档匹配16:35的平均值。
示例数据
[ { _id: ObjectId("63a318c36ccc42d2330fae5e"), timestamp: ISODate("2022-12-21T14:30:31.172Z"), value: 3.8 }, { _id: ObjectId("63a318c46ccc42d2330fae8d"), timestamp: ISODate("2022-12-21T14:30:32.189Z"), value: 4.0 }, { _id: ObjectId("63a318c36ccc42d2330fae5e"), timestamp: ISODate("2022-12-21T15:30:14.025Z"), value: 5.0 }, { _id: ObjectId("63a318c36ccc42d2330fae5e"), timestamp: ISODate("2022-12-21T15:30:18.025Z"), value: 5.5 } ]
现有代码参考
已实现的1分钟分组平均值计算
{ $group: { _id: { "code": "$code", "year": { "$year": "$timestamp" }, "dayOfYear": { "$dayOfYear": "$timestamp" }, "hour": { "$hour": "$timestamp" }, "minute": { "$minute": "$timestamp" } }, value: { "$avg": "$value" }, timestamp: { "$first": "$timestamp" } } }
不符合需求的小时级聚合代码
{ $group: { _id: { "code": "$code", "year": { "$year": "$timestamp" }, "dayOfYear": { "$dayOfYear": "$timestamp" }, "hour": { "$hour": "$timestamp" } }, value: { "$first": "$value" }, valueLast: { "$last": "$value" }, timestamp: { "$first": "$timestamp" } } }
实现方案
核心思路是先预计算所有1分钟窗口的平均值,再通过关联操作将原始文档与1小时后的对应窗口平均值绑定,最后计算增益。完整聚合管道如下:
[ // 阶段1:预计算所有1分钟窗口的平均值 { $group: { _id: { code: "$code", year: { $year: "$timestamp" }, dayOfYear: { $dayOfYear: "$timestamp" }, hour: { $hour: "$timestamp" }, minute: { $minute: "$timestamp" } }, avgValue: { $avg: "$value" } } }, // 阶段2:将预计算的平均值存入临时集合,方便后续关联 { $merge: { into: "minute_avg_temp", whenMatched: "replace", whenNotMatched: "insert" } }, // 阶段3:切换回原始文档,生成用于匹配1小时后窗口的标识 { $unionWith: { coll: "your_original_collection" } }, { $addFields: { targetAvgId: { $let: { vars: { // 计算当前时间+1小时的时间戳 oneHourLaterTs: { $add: [ { $toLong: "$timestamp" }, 3600000 ] } }, in: { code: "$code", year: { $year: { $toDate: "$$oneHourLaterTs" } }, dayOfYear: { $dayOfYear: { $toDate: "$$oneHourLaterTs" } }, hour: { $hour: { $toDate: "$$oneHourLaterTs" } }, minute: { $minute: { $toDate: "$$oneHourLaterTs" } } } } } } }, // 阶段4:关联临时集合中的1小时后平均值 { $lookup: { from: "minute_avg_temp", localField: "targetAvgId", foreignField: "_id", as: "lastAvgData" } }, // 阶段5:提取平均值并计算增益百分比 { $addFields: { valueLast: { $arrayElemAt: [ "$lastAvgData.avgValue", 0 ] }, gainPct: { $cond: { if: { $ne: [ "$value", 0 ] }, then: { $divide: [ { $subtract: [ { $arrayElemAt: [ "$lastAvgData.avgValue", 0 ] }, "$value" ] }, "$value" ] }, else: null // 避免除以0的情况 } } } }, // 阶段6:过滤掉预计算的平均值文档,只保留原始文档,并清理临时字段 { $match: { _id: { $type: "objectId" } } // 假设原始文档的_id是ObjectId,预计算文档的_id是嵌套对象 }, { $project: { targetAvgId: 0, lastAvgData: 0 } } ]
方案说明
- 临时集合存储平均值:用
$merge将1分钟窗口的平均值存入临时集合,避免重复计算; - 时间窗口匹配:通过给每个原始文档生成「1小时后窗口的唯一标识」,与临时集合中的平均值关联;
- 异常处理:添加了除以0的判断,避免
gainPct出现NaN; - 结果过滤:最后只保留原始文档,排除预计算的平均值数据;
- 注意事项:
$unionWith中的your_original_collection需替换为你的原始数据集合名称,临时集合minute_avg_temp可根据需求调整名称。
期望输出示例
[ { _id: ObjectId("63a318c36ccc42d2330fae5e"), timestamp: ISODate("2022-12-21T14:30:31.172Z"), value: 3.8, valueLast: 5.25, gainPct: 0.382 }, { _id: ObjectId("63a318c46ccc42d2330fae8d"), timestamp: ISODate("2022-12-21T14:30:32.189Z"), value: 4.0, valueLast: 5.25, gainPct: 0.313 } ]
内容的提问来源于stack exchange,提问作者fip
相关产品推荐
相关产品推荐

