Learn
MongoDB/15-change-streams

变更流 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更新操作符否(需选项)
replacereplaceOne是
delete删除否,只有 documentKey
drop集合被删—
rename集合改名—
dropDatabase库被删—
invalidate流失效(drop/rename 后)—
⚠️invalidate 会终止变更流

集合被 drop 或 rename 时,会先收到一个 drop / rename 事件,紧接着一个 invalidate,然后变更流关闭。此时用旧的 resume token 恢复会失败,必须用 startAfter(4.2+)才能跨过 invalidate 继续。生产消费者必须处理这个情况。

3. 常用选项

3.1 fullDocument

db.posts.watch([], { fullDocument: "updateLookup" })
值行为
defaultupdate 事件不带完整文档
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 } }
      ]
  } }
])
⚠️不要在管道里 project 掉 _id

变更事件的 _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 })
})
⚠️token 过期取决于 oplog 窗口

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:05

24 小时的窗口意味着:消费者最多能停机 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 的优势

维度裸读 oplogChange 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 和高可用的底层机制 →