变更流
本教程共 50 篇 · 第 48 篇 · 更新于 2026-07-30 · 约 4 分钟阅读
48. 变更流
本节目标:理解变更流能做什么,会用 watch() 监听集合变化,并知道断点续传和部署前提。
有些业务要在数据一变就立刻响应:刷新缓存、写审计日志、同步到别的系统。轮询数据库太低效。变更流(Change Streams)让应用「订阅」数据变化,像监听消息队列一样实时拿到事件。
本节内容依据 MongoDB 8.3 官方文档「Change Streams」整理,变更流仅在副本集与分片集群可用。
48.1 前提:必须跑在副本集上
变更流底层依赖 oplog(操作日志),所以它只能在副本集(Replica Set)或分片集群上用。单机 mongod 不行。
本地开发想试,可以先把单机转成单节点副本集:
# 启动时为副本集模式
mongod --replSet rs0 --dbpath /data/db
// 另一个终端进 mongosh,初始化副本集
test> rs.initiate()
Warning不少新手在单机
mongod上跑watch()直接报错。记住:变更流要先有副本集。这是官方明确的前提。
48.2 用 watch() 打开监听
在集合、数据库、整个部署上都能开变更流。最常用的是集合级别。
// 监听 users 集合的所有变更
test> const stream = db.users.watch()
test> const change = stream.next()
test> printjson(change)
得到的事件文档含:operationType(insert/update/delete 等)、documentKey(对应 _id)、fullDocument(变更后的文档,视选项而定)、ns(命名空间)。
插入一条数据,就能在流里看到对应的 insert 事件:
test> db.users.insertOne({ name: "小明", age: 28 })
48.3 fullDocument:要不要带完整文档
默认更新事件只给「变了什么」,不带完整新文档。想要整条新文档,用 fullDocument 选项。
// updateLookup:更新事件里附带更新后的完整文档
test> const stream = db.users.watch(
[],
{ fullDocument: "updateLookup" }
)
Note
fullDocument可选值包括default、updateLookup、whenAvailable、required。本地 8.3 用updateLookup最常见,能拿到变更后全貌。
48.4 resumeToken:断点续传
网络断了或程序重启,要接着上次的位置听,不能从头来。每个事件文档的 _id 就是恢复令牌(Resume Token)。
// 从保存的令牌处继续监听
test> const stream = db.users.watch(
[],
{ resumeAfter: savedToken }
)
savedToken 就是上一次事件里的 change._id。把它持久化到别处(如一张表或文件),崩溃后读出来续传。
Tip生产环境务必持久化 resumeToken。否则重启会漏掉中断期间的变更,或从头重放造成重复处理。
48.5 用管道过滤事件
watch() 第一个参数是聚合管道,可只关心某些事件。
// 只监听 update 操作
test> const stream = db.users.watch([
{ $match: { operationType: "update" } }
])
也能用 $match 过滤特定字段变化,减少无效通知。
48.6 适用场景
- 实时同步:把变更推到搜索引擎或数据仓库。
- 审计日志:记录谁改了什么,满足合规。
- 缓存失效:文档一变就清掉对应缓存。
- 事件驱动:订单状态变了,触发后续流程。
Note变更流是基于聚合框架的,能过滤也能转换通知内容,比自己 tail oplog 安全、简单得多。
48.7 小结
变更流用 watch() 实时订阅数据变化,事件里的 _id 是 resumeToken,配合 resumeAfter 断点续传,fullDocument 控制是否带全文档。它必须跑在副本集或分片集群上,是做实时系统的好帮手。