变更流要求副本集或分片集群,单节点不支持;需显式指定replicaSet、传full_document="updateLookup"、用resume_after恢复监听,并注意pipeline字段匹配规则。

变更流需要副本集或分片集群
单节点 MongoDB 实例不支持变更流,直接调用 watch() 会抛出 PymongoError: cannot open $changeStream on non-replica set。必须先将单机升级为副本集(哪怕只有 1 个成员),或连接到已启用副本集的集群。
本地快速验证可用以下命令启动最小副本集:
mongod --replSet rs0 --dbpath /data/db --port 27017
然后在 mongo shell 中执行:
rs.initiate({ _id: "rs0", members: [{ _id: 0, host: "localhost:27017" }] })- 生产环境务必使用至少 3 个成员的副本集,避免脑裂和主节点单点故障
- 连接时需在 URI 中显式指定 replica set 名称,例如:
mongodb://localhost:27017/?replicaSet=rs0 - 如果使用 Atlas,确保集群类型是 “Replica Set” 或 “Sharded Cluster”,并勾选 “Enable Change Streams”
watch() 调用必须带 pipeline 和 full_document
默认调用 collection.watch() 只返回变更事件的元信息(如 _id、operationType),不包含文档内容。要获取完整文档,必须传入 full_document="updateLookup" 参数;否则对 update 类型事件只能看到 updateDescription,看不到改了什么字段。
立即学习“Python免费学习笔记(深入)”;
常见 pipeline 示例:
pipeline = [
{ "$match": { "operationType": { "$in": ["insert", "update", "delete"] } } },
{ "$addFields": { "eventTime": "$clusterTime" } }
]-
full_document="updateLookup"会让 driver 自动查一次最新文档,但会增加一次读请求,注意性能影响 - 若只关心某几个字段变化,可在
$match中加"updateDescription.updatedFields"过滤,但字段名需用双引号包裹(MongoDB 的 BSON 语法) - 不要在 pipeline 里写
$project去删字段——变更流事件结构固定,删了关键字段(如_id、operationType)会导致解析失败
监听循环必须处理断连与 ResumeToken
网络抖动或主节点切换时,watch() 迭代器会抛出 PyMongoError 或 StopIteration,不能简单用 except Exception 吞掉就重连——那样会丢失断连期间的变更。
正确做法是捕获异常后,用上一次成功事件的 _id(即 resume token)重启监听:
resume_token = None
while True:
try:
with collection.watch(pipeline, full_document="updateLookup", resume_after=resume_token) as stream:
for change in stream:
print(change)
resume_token = change["_id"]
except PyMongoError as e:
if isinstance(e, InvalidOperation):
# 如 resume_after 无效,可设为 None 重新开始(慎用)
resume_token = None
time.sleep(1)-
resume_after是最轻量的恢复方式,比start_after更安全,但 token 有效期默认 30 分钟(受 oplog 大小限制) - 首次运行时
resume_after=None,后续每次取change["_id"],不是change.get("clusterTime") - 别把 resume token 存在内存里——进程重启就丢了;生产环境建议存到 Redis 或本地文件,并加时间戳防过期
聚合管道中 $match 的字段路径容易写错
变更流事件结构嵌套较深,比如判断是否更新了 status 字段,不能写 {"updateDescription.updatedFields.status": {"$exists": True}} ——因为 updatedFields 是一个对象,其 key 就是字段名,不是嵌套路径。
正确写法是:
{ "$match": { "updateDescription.updatedFields.status": { "$exists": true } } }但注意:MongoDB 的 true 是小写,Python 里用 True 会被自动转成 true,没问题;而字段名中的点号 . 在这里表示“对象的键”,不是路径分隔符。
- 所有变更事件顶层字段:
operationType、fullDocument、updateDescription、documentKey -
updateDescription.updatedFields是一个字典,key 是被修改的字段名(如"name"),value 是新值;所以匹配字段存在,就是查这个 dict 的 key - 想匹配数组元素变化?得用
"updateDescription.updatedFields.array.0.field"这种形式,但前提是数组索引确定——动态索引无法用 pipeline 表达,得靠应用层过滤
变更流不是“实时消息队列”,它依赖 oplog 的滚动窗口和网络稳定性。resume token 过期、pipeline 语法错误、权限不足(比如没授予 changeStream 角色)这三类问题占实际排障的八成以上,检查时优先盯住这三个点。


















