变更流 Change Streams
数据变了,怎么通知其他系统?传统做法是轮询「更新时间大于上次检查点」的记录,既有延迟又浪费资源,还会漏掉删除操作。Change Streams 把数据库变更变成一条可订阅的实时流,是 MongoDB 做 CDC(Change Data Capture)的官方方案。
1. 为什么需要变更流
1.1 轮询方案的四个缺陷
// 传统轮询
setInterval(async () => {
const docs = await db.posts.find({ updatedAt: { $gt: lastCheck } }).toArray()
lastCheck = new Date()
await syncToES(docs)
}, 5000)| 缺陷 | 说明 |
|---|---|
| 有延迟 | 轮询间隔就是延迟下限 |
| 有开销 | 没有变更时也在查询 |
| 漏删除 | 文档删了就查不到了 |
| 丢变更 | 同一文档在两次轮询之间改了 3 次,只能看到最终态 |
| 时钟依赖 | 依赖 updatedAt 准确,应用忘记维护就失效 |
1.2 变更流的定位
应用写入
│
▼
┌─────────────┐
│ primary │
│ │
│ oplog ─────┼──── 复制给 secondary
│ │ │
└────┼────────┘
│
▼ Change Streams 订阅 oplog(经过封装和过滤)
┌──────────────┐
│ 下游消费者 │ → 同步到 ES / 刷缓存 / 发通知 / 数据仓库
└──────────────┘它本质上是对 oplog 的封装订阅,但提供了 oplog 裸读不具备的能力:跨主从切换的断点续传、有序保证、以及 majority 级别的读一致性(只推送已被多数节点确认的变更,不会推送后来被回滚的数据)。
Change Streams 依赖 oplog,standalone 部署无法使用。这也是第 2 章要求本地用单成员复制集启动的原因之一。
2. 基本用法
2.1 三个订阅层级
// 集合级
const cs = db.posts.watch()
// 数据库级(该库所有集合)
const cs = db.watch()
// 部署级(所有库,需要 admin 权限)
const cs = db.getMongo().watch()while (cs.hasNext()) {
printjson(cs.next())
}2.2 变更事件的结构
插入一条帖子:
db.posts.insertOne({ _id: 500, title: "新帖子", authorId: "alice", status: "draft" })收到的事件:
{
"_id": { "_data": "82665F1A2B000000012B022C0100296E5A100..." },
"operationType": "insert",
"clusterTime": Timestamp({ t: 1717495851, i: 1 }),
"wallTime": ISODate("2024-06-04T09:30:51.234Z"),
"ns": { "db": "community", "coll": "posts" },
"documentKey": { "_id": 500 },
"fullDocument": {
"_id": 500, "title": "新帖子", "authorId": "alice", "status": "draft"
}
}更新一个字段:
db.posts.updateOne({ _id: 500 }, { $set: { status: "published" }, $inc: { "stats.views": 1 } }){
"_id": { "_data": "82665F1A2C..." },
"operationType": "update",
"clusterTime": Timestamp({ t: 1717495852, i: 1 }),
"ns": { "db": "community", "coll": "posts" },
"documentKey": { "_id": 500 },
"updateDescription": {
"updatedFields": { "status": "published", "stats.views": 1 },
"removedFields": [],
"truncatedArrays": []
}
}注意 update 事件默认不带完整文档,只有变更描述。删除事件更极端,只有 documentKey:
{
"operationType": "delete",
"ns": { "db": "community", "coll": "posts" },
"documentKey": { "_id": 500 }
}2.3 operationType 全集
| 类型 | 触发 | 是否带 fullDocument |
|---|---|---|
insert | 插入 | 是 |
update | 更新操作符 | 否(需选项) |
replace | replaceOne | 是 |
delete | 删除 | 否,只有 documentKey |
drop | 集合被删 | — |
rename | 集合改名 | — |
dropDatabase | 库被删 | — |
invalidate | 流失效(drop/rename 后) | — |
集合被 drop 或 rename 时,会先收到一个 drop / rename 事件,紧接着一个 invalidate,然后变更流关闭。此时用旧的 resume token 恢复会失败,必须用 startAfter(4.2+)才能跨过 invalidate 继续。生产消费者必须处理这个情况。
3. 常用选项
3.1 fullDocument
db.posts.watch([], { fullDocument: "updateLookup" })| 值 | 行为 |
|---|---|
default | update 事件不带完整文档 |
updateLookup | 变更后额外查一次当前文档 |
whenAvailable | 有变更前后镜像就带上(6.0+) |
required | 必须有镜像,否则报错(6.0+) |
updateLookup 最常用,但要理解它的语义:它查的是「查询那一刻的文档」,不是「变更后立即的文档」。如果同一文档被连续改了 3 次,三个事件的 fullDocument 可能都是最终态。
3.2 前像与后像(6.0+)
需要精确的「变更前」快照时,开启集合级的 pre/post image:
db.runCommand({ collMod: "posts", changeStreamPreAndPostImages: { enabled: true } })
db.posts.watch([], {
fullDocumentBeforeChange: "whenAvailable",
fullDocument: "whenAvailable"
}){
"operationType": "update",
"fullDocumentBeforeChange": { "_id": 500, "status": "draft", "stats": { "views": 0 } },
"fullDocument": { "_id": 500, "status": "published", "stats": { "views": 1 } }
}这让「审计日志」「字段级 diff」成为可能。代价是 MongoDB 要额外存储镜像,占用空间和写入开销。
镜像存在一个特殊集合里,默认由 expireAfterSeconds 控制过期。如果消费者落后太多,取不到镜像时 required 模式会报错、whenAvailable 模式会返回空。消费延迟监控在开启前后像之后变得更重要。
3.3 batchSize 与 maxAwaitTimeMS
db.posts.watch([], {
batchSize: 100, // 每批最多返回多少个事件
maxAwaitTimeMS: 1000 // 没有新事件时最多等多久(影响空闲时的响应延迟)
})4. 过滤管道
变更流可以接一个聚合管道,在服务端过滤,减少网络传输:
db.posts.watch([
{ $match: {
operationType: { $in: ["insert", "update"] },
"fullDocument.status": "published"
} },
{ $project: {
_id: 1, // resume token 必须保留!
operationType: 1,
documentKey: 1,
"fullDocument.title": 1,
"fullDocument.author": 1,
"updateDescription.updatedFields": 1
} }
], { fullDocument: "updateLookup" })4.1 只关心特定字段的变更
扩展引用模式的同步场景:只有用户名或头像变了才需要同步到 posts。
db.users.watch([
{ $match: {
operationType: "update",
$or: [
{ "updateDescription.updatedFields.username": { $exists: true } },
{ "updateDescription.updatedFields.profile.avatar": { $exists: true } }
]
} }
])变更事件的 _id 就是 resume token,它是断点续传的唯一依据。管道里如果写了 { $project: { operationType: 1 } }(没带 _id),MongoDB 会直接报错,因为流将无法恢复。
4.2 支持的阶段
变更流管道只允许这些阶段:$match、$project、$addFields、$set、$replaceRoot、$replaceWith、$redact、$unset。不允许 $group、$sort、$lookup——因为它们需要看到全部数据,而流是无限的。
5. Resume Token 与断点续传
5.1 为什么必须持久化 token
消费者进程重启、网络断开、主从切换——任何一次中断,都需要知道「上次处理到哪」。resume token 就是这个位置标记。
let resumeToken = loadTokenFromStorage() // 从 Redis/文件/集合里读
const options = { fullDocument: "updateLookup" }
if (resumeToken) options.resumeAfter = resumeToken
const cs = db.posts.watch([], options)
while (cs.hasNext()) {
const change = cs.next()
processChange(change) // 业务处理
resumeToken = change._id
saveTokenToStorage(resumeToken) // 处理成功后再保存
}5.2 保存 token 的顺序决定语义
先保存 token,再处理业务:
处理失败 → token 已推进 → 事件丢失 → at-most-once
先处理业务,再保存 token:
保存失败 → token 未推进 → 事件重放 → at-least-once ✔ 推荐选 at-least-once,然后让处理逻辑幂等。这是所有消息系统的标准做法。
5.3 resumeAfter vs startAfter vs startAtOperationTime
| 选项 | 语义 | 适用 |
|---|---|---|
resumeAfter | 从该 token 之后继续 | 正常续传 |
startAfter | 同上,但能跨过 invalidate | 集合被 drop/rename 后恢复 |
startAtOperationTime | 从某个 clusterTime 开始 | 没有 token 时按时间点启动 |
// 从 10 分钟前开始重放
db.posts.watch([], {
startAtOperationTime: Timestamp({ t: Math.floor(Date.now() / 1000) - 600, i: 1 })
})resume token 指向 oplog 里的一个位置。如果消费者停机时间超过 oplog 的保留窗口(默认是磁盘的 5%,可能只有几小时),那个位置已经被覆盖,恢复时会报 ChangeStreamHistoryLost。此时只能全量重同步。监控消费延迟和 oplog 窗口,是运行 CDC 的必修课。
// 查看 oplog 窗口有多长
db.getSiblingDB("local").oplog.rs.stats().maxSize / 1024 / 1024 / 1024 // GB
rs.printReplicationInfo()configured oplog size: 51200MB
log length start to end: 86400secs (24hrs)
oplog first event time: Mon Jun 03 2024 09:00:00
oplog last event time: Tue Jun 04 2024 09:00:00
now: Tue Jun 04 2024 09:00:0524 小时的窗口意味着:消费者最多能停机 24 小时。
6. 与 oplog 的关系
6.1 oplog 是什么
oplog 是 local.oplog.rs 这个定长集合,记录了所有改变数据的操作,供 secondary 拉取重放。
db.getSiblingDB("local").oplog.rs.find().sort({ $natural: -1 }).limit(1){
"op": "u",
"ns": "community.posts",
"ui": UUID("..."),
"o": { "$v": 2, "diff": { "u": { "status": "published" } } },
"o2": { "_id": 500 },
"ts": Timestamp({ t: 1717495852, i: 1 }),
"t": NumberLong(1),
"wall": ISODate("2024-06-04T09:30:52.001Z")
}6.2 oplog 的幂等性
oplog 里的操作被改写成幂等形式:$inc: { views: 1 } 在 oplog 里会变成 $set: { views: 1201 }。这样重放多次结果一致,是复制正确性的基础。
这也解释了一个现象:为什么 updateMany 影响 10 万条文档会产生 10 万条 oplog——因为每个文档的最终结果都要单独记录。
6.3 变更流相比裸读 oplog 的优势
| 维度 | 裸读 oplog | Change Streams |
|---|---|---|
| 权限 | 需要 local 库读权限 | 集合级权限即可 |
| 格式 | 内部格式,版本间会变 | 稳定的公开 API |
| 一致性 | 可能读到会被回滚的数据 | 只推 majority 已确认的 |
| 主从切换 | 需要自己处理 | token 自动跨节点 |
| 分片 | 要连每个分片自己合并 | mongos 统一有序推送 |
| 过滤 | 客户端过滤 | 服务端管道过滤 |
结论很明确:不要自己读 oplog,除非你在写 MongoDB 的复制组件。
7. CDC 实战:三个场景
7.1 同步扩展引用字段
第 12 章遗留的问题:用户改名后,怎么更新所有帖子里冗余的 author.username。
const cs = db.users.watch([
{ $match: {
operationType: "update",
"updateDescription.updatedFields.username": { $exists: true }
} }
], { fullDocument: "updateLookup" })
while (cs.hasNext()) {
const change = cs.next()
const userId = change.documentKey._id
const newName = change.fullDocument.username
db.posts.updateMany(
{ "author._id": userId },
{ $set: { "author.username": newName } }
)
db.comments.updateMany(
{ "author._id": userId },
{ $set: { "author.username": newName } }
)
saveToken(change._id)
}这就是「最终一致性」的具体实现:改名后几百毫秒内,所有冗余副本都会更新。
7.2 同步到 Elasticsearch
const cs = db.posts.watch([
{ $match: { "ns.coll": "posts" } }
], { fullDocument: "updateLookup" })
while (cs.hasNext()) {
const c = cs.next()
const id = c.documentKey._id.toString()
switch (c.operationType) {
case "insert":
case "replace":
case "update":
if (c.fullDocument) {
await es.index({ index: "posts", id, document: {
title: c.fullDocument.title,
body: c.fullDocument.body,
tags: c.fullDocument.tags
} })
}
break
case "delete":
await es.delete({ index: "posts", id }).catch(() => {})
break
}
saveToken(c._id)
}注意 delete 事件只有 documentKey,所以 ES 的文档 id 必须和 MongoDB 的 _id 一致,否则删除时无法定位。
7.3 缓存失效
db.posts.watch([
{ $match: { operationType: { $in: ["update", "replace", "delete"] } } }
]).forEach(c => {
redis.del(`post:${c.documentKey._id}`)
saveToken(c._id)
})8. 生产环境注意事项
| 问题 | 应对 |
|---|---|
| 单消费者是瓶颈 | 按 documentKey._id 哈希分片,多个消费者各消费一部分 |
| 顺序保证 | 变更流全局有序,但并行消费后就无序了;同一文档的事件要路由到同一消费者 |
| 消费延迟 | 用 clusterTime 和当前时间的差值做监控指标并告警 |
| oplog 窗口 | 至少要能覆盖「最长的计划停机时间 + 报警响应时间」 |
| 主从切换 | 驱动会自动重连并用 token 恢复,但要确认 token 已持久化 |
| 大批量写入 | updateMany 影响 100 万文档会产生 100 万个变更事件,可能压垮消费者 |
一次 db.posts.updateMany({}, { $set: { v: 1 } }) 影响 1000 万文档,变更流会推送 1000 万个事件。如果下游是 ES,很可能直接被打挂。做大批量数据订正时,要么临时停掉消费者,要么分批 + 限速执行。
一、在一个 mongosh 窗口开 db.posts.watch(),另一个窗口做增删改,观察三种 operationType 的事件结构差异;二、开启 fullDocument: "updateLookup" 后重做一遍,对比 update 事件的变化;三、写一个带 resume token 持久化的消费者(token 存到一个 _cdcState 集合里),中途 kill 掉再重启,验证不丢事件;四、写一个只监听「用户名或头像变更」的过滤管道;五、用 rs.printReplicationInfo() 查看你的 oplog 窗口有多长,计算消费者最多能停机多久;六、实现第 7.1 节的用户名同步,并思考如果同一用户在 1 秒内改名 3 次会发生什么。
小结
- Change Streams 是对 oplog 的封装订阅,解决轮询的延迟、开销、漏删除三大问题
- 只推送 majority 已确认的变更,不会推送后来被回滚的数据
- update 事件默认只有变更描述,需要完整文档要开
updateLookup;delete 只有documentKey - 服务端过滤管道能大幅减少传输量,但绝不能 project 掉
_id - resume token 必须持久化,先处理业务再保存 token,配合幂等实现 at-least-once
- token 有效期取决于 oplog 窗口,超期只能全量重同步,必须监控
invalidate事件会关闭流,恢复要用startAfter- 不要裸读 oplog,Change Streams 在权限、稳定性、分片支持上全面更优
- 下一章讲复制集,理解 oplog 和高可用的底层机制 →