MongoDBNotes

第 12 章:Change Stream

zjc 于 2026-01-12 发布

这是《MongoDB 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Change Stream 让应用可以订阅集合、数据库或整个实例的数据变更。它基于副本集 oplog,天然适合缓存失效、搜索索引同步、审计联动和事件驱动架构。

12.1 基本原理

application
  -> open change stream cursor
  -> tail oplog
  -> receive change events
  -> resume from resume token

要求:

  1. 使用副本集或分片集群;
  2. 用户有 findchangeStream 权限;
  3. oplog 窗口足够长;
  4. 应用保存 resume token;
  5. 断开后可恢复。

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

更稳的顺序是:

  1. 读取事件;
  2. 在下游事务中处理事件并保存新 token;
  3. 处理失败重试;
  4. 死信记录;
  5. 恢复时从已提交 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 容量与故障

关注:

  1. oplog 大小;
  2. 事件消费速率;
  3. resume token 是否过期;
  4. 下游写入吞吐;
  5. 重复事件比例;
  6. 全量重建窗口。

处理积压:

  1. 提高消费并行度;
  2. 批量写下游;
  3. 减少无关集合监听;
  4. 拆分消费者;
  5. 全量重建。

本章小结

Change Stream 是构建增量同步和事件联动的基础能力,依赖副本集 oplog 和可靠的 resume token 机制。设计时必须把 token 保存、事件幂等、全量初始化和消费积压治理纳入方案。

思考题

  1. Change Stream 为什么依赖副本集?
  2. Resume Token 应该如何保存?
  3. fullDocument: updateLookup 有什么并发语义?
  4. 下游索引重建时如何组合全量和增量?
  5. Resume Token 过期怎么办?