这是《MongoDB 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Change Stream 让应用可以订阅集合、数据库或整个实例的数据变更。它基于副本集 oplog,天然适合缓存失效、搜索索引同步、审计联动和事件驱动架构。
12.1 基本原理
application
-> open change stream cursor
-> tail oplog
-> receive change events
-> resume from resume token
要求:
- 使用副本集或分片集群;
- 用户有
find和changeStream权限; - oplog 窗口足够长;
- 应用保存 resume token;
- 断开后可恢复。
12.2 打开 Change Stream
监听集合:
const stream = db.orders.watch()
while (true) {
if (!stream.hasNext()) break;
printjson(stream.next())
}
Node.js 示例:
const changeStream = db.collection("orders").watch();
changeStream.on("change", async (event) => {
console.log(event);
});
changeStream.on("error", (err) => {
console.error(err);
changeStream.close();
});
12.3 事件结构
插入事件示例:
{
"_id": { "_data": "8266CB..." },
"operationType": "insert",
"clusterTime": { "$timestamp": { "t": 1760000000, "i": 1 } },
"fullDocument": {
"_id": "o_10001",
"status": "CREATED"
},
"ns": { "db": "shop", "coll": "orders" },
"documentKey": { "_id": "o_10001" }
}
常见 operationType:
| 类型 | 说明 |
|---|---|
insert |
插入 |
update |
更新 |
replace |
替换 |
delete |
删除 |
invalidate |
集合失效 |
drop |
集合删除 |
12.4 过滤事件
只看支付订单:
const stream = db.orders.watch([
{
$match: {
"fullDocument.status": "PAID",
operationType: { $in: ["insert", "update", "replace"] }
}
}
])
只看特定字段更新:
const stream = db.orders.watch([
{
$match: {
"updateDescription.updatedFields.status": { $exists: true }
}
}
])
服务端过滤能减少网络传输,但复杂业务判断通常仍在消费端完成。
12.5 获取完整文档
更新事件默认不包含完整最新文档:
const stream = db.orders.watch([], {
fullDocument: "updateLookup"
})
选项:
| 值 | 行为 |
|---|---|
default |
插入替换有全文,更新只有差异 |
updateLookup |
查询当前文档全文 |
whenAvailable |
有缓存则返回 |
required |
要求全文 |
并发更新很快时,updateLookup 看到的可能是更新后状态,需要业务幂等。
12.6 Resume Token
事件 _id 就是 resume token:
const event = stream.next();
db.tokens.updateOne(
{ name: "orders-sync" },
{ $set: { token: event._id, updated_at: new Date() } },
{ upsert: true }
)
恢复:
const saved = db.tokens.findOne({ name: "orders-sync" });
const stream = saved
? db.orders.watch([], { resumeAfter: saved.token })
: db.orders.watch();
Token 保存位置应与业务下游状态一致,避免漏事件或大量重复。
12.7 起始位置
从当前开始:
db.orders.watch([], { startAtOperationTime: new Date() })
从某个时间开始:
db.orders.watch([], {
startAtOperationTime: new Date("2026-08-25T00:00:00Z")
})
时间起点受 oplog 窗口限制。重建下游索引时通常需要全量加增量,而不是只靠 Change Stream。
12.8 事件处理模式
可靠处理流程:
read event
-> save resume token
-> apply event idempotently
-> commit downstream
-> advance token
更稳的顺序是:
- 读取事件;
- 在下游事务中处理事件并保存新 token;
- 处理失败重试;
- 死信记录;
- 恢复时从已提交 token 继续。
如果 token 与下游状态分开提交,可能出现重复或漏处理窗口。
12.9 典型应用
缓存失效:
product update
-> change stream
-> delete redis cache
-> next request rebuilds cache
搜索同步:
full import
-> change stream
-> upsert Elasticsearch document
审计联动:
order status change
-> change stream
-> emit event
-> notification service
12.10 容量与故障
关注:
- oplog 大小;
- 事件消费速率;
- resume token 是否过期;
- 下游写入吞吐;
- 重复事件比例;
- 全量重建窗口。
处理积压:
- 提高消费并行度;
- 批量写下游;
- 减少无关集合监听;
- 拆分消费者;
- 全量重建。
本章小结
Change Stream 是构建增量同步和事件联动的基础能力,依赖副本集 oplog 和可靠的 resume token 机制。设计时必须把 token 保存、事件幂等、全量初始化和消费积压治理纳入方案。
思考题
- Change Stream 为什么依赖副本集?
- Resume Token 应该如何保存?
fullDocument: updateLookup有什么并发语义?- 下游索引重建时如何组合全量和增量?
- Resume Token 过期怎么办?