如何调整PyMongo聚合查询以适配物化视图合并需求?
解决方案
一、调整聚合查询直接生成目标格式
可以通过两次分组+$arrayToObject操作直接得到你要的输出格式,修改后的聚合查询如下:
from pymongo import MongoClient, SON client = MongoClient() db = client.your_database pipeline = [ # 过滤无告警记录 {"$match": {"sevr": {"$ne": "NO_ALARM"}}}, # 第一次分组:统计每个name+sevr组合的数量 {"$group": { "_id": {"name": "$name", "sevr": "$sevr"}, "count": {"$sum": 1} }}, # 第二次分组:按name聚合,将sevr和count转为键值对数组 {"$group": { "_id": "$_id.name", "count_entries": {"$push": {"k": "$_id.sevr", "v": "$count"}} }}, # 将键值对数组转为嵌套对象,生成目标count结构 {"$addFields": { "count": {"$arrayToObject": "$count_entries"} }}, # 移除临时字段 {"$project": {"count_entries": 0}}, # 按需排序(可根据实际需求调整排序规则) {"$sort": SON([("_id", -1)])} ] result = db.events.history.aggregate(pipeline) for doc in result: print(doc)
该查询的逻辑:
- 先过滤掉
sevr为NO_ALARM的无效记录; - 第一次分组统计每个
name与sevr组合的出现次数; - 第二次分组按
name聚合,把sevr和对应数量组装成{k: sevr值, v: 统计数}的数组; - 用
$arrayToObject将数组转为嵌套对象,得到目标格式的count字段; - 清理临时字段并完成排序。
输出结果会直接匹配你的期望格式:
[ {"_id": "name_1", "count": {"SEVR_1": 100, "SEVR_2": 50}}, {"_id": "name_2", "count": {"SEVR_1": 75, "SEVR_2": 30}} ]
二、无需副本集的变更检测方法
MongoDB官方的Change Streams依赖副本集或分片集群,无副本集场景下只能通过轮询检测实现近似的自动变更感知,具体方案如下:
方法1:基于时间字段轮询
如果你的文档包含createdAt/updatedAt这类时间字段,可以记录每次同步物化视图的最后时间戳,定时查询该时间戳之后的新增/修改文档:
import time from datetime import datetime # 初始化最后同步时间(首次同步设为很早的时间) last_sync_time = datetime(1970, 1, 1) while True: # 查询最后同步时间之后的有效变更记录 updated_docs = db.events.history.find({ "updatedAt": {"$gt": last_sync_time}, "sevr": {"$ne": "NO_ALARM"} }) # 执行物化视图更新逻辑(替换为你的实际更新代码) # update_materialized_view(updated_docs) # 更新最后同步时间为当前时间 last_sync_time = datetime.now() # 每隔60秒轮询一次(可根据实时性需求调整间隔) time.sleep(60)
方法2:基于ObjectId的时间特性轮询
如果文档_id是MongoDB默认的ObjectId,它内置了文档创建时间戳,可利用这一点检测新增文档:
import time from bson.objectid import ObjectId # 初始化最后处理的ObjectId(首次同步用一个早期ID) last_id = ObjectId("000000000000000000000000") while True: # 查询比last_id更新的有效文档 updated_docs = db.events.history.find({ "_id": {"$gt": last_id}, "sevr": {"$ne": "NO_ALARM"} }).sort("_id", 1) # 执行物化视图更新逻辑 # update_materialized_view(updated_docs) # 更新last_id为本次查询的最后一个文档ID last_doc = updated_docs.sort("_id", -1).limit(1).next() if last_doc: last_id = last_doc["_id"] time.sleep(60)
注意事项
- 轮询间隔需平衡实时性与资源消耗,间隔越短对数据库压力越大;
- 方法1需确保文档修改时
updatedAt字段被同步更新,方法2仅能检测新增文档,无法感知已有文档的修改; - 可结合定时任务工具(如Linux cron、Python APScheduler)替代无限循环,实现更稳定的周期性同步。
内容的提问来源于stack exchange,提问作者Fonty
相关产品推荐
相关产品推荐

