本节摘要:Change Streams 订阅集合的插入、更新、删除事件,基于 oplog 提供断点续传令牌(resumeToken),是构建缓存失效、消息联动、审计同步的基础设施。本节从一次缓存与库不一致的漏单事故讲它的用法与容错要点。
订单服务用每 10 分钟一次的定时任务全量对账同步缓存。两次同步之间支付的订单,在缓存里最长有 10 分钟不存在——客服系统看不到,用户投诉"付了钱没订单"。轮询的窗口就是不一致的黑洞。改用 Change Streams 后,库一动、缓存即刷新,黑洞消失。
// 订阅订单集合的变更(需连接复制集,单机模式不可用) const cs = db.orders.watch( [{ $match: { operationType: { $in: ["insert", "update"] } } }], { fullDocument: "updateLookup" } // 更新事件附带完整文档 ); while (cs.hasNext()) { const ev = cs.next(); // ev.fullDocument 即最新文档,刷新缓存 cache.set(ev.documentKey._id, ev.fullDocument); // 持久化令牌:进程重启后从这里继续,不丢事件 saveToken(ev._id); // resumeToken } // 重启恢复 const cs2 = db.orders.watch([], { resumeAfter: loadToken() });
⚠️ Change Streams 只能在复制集(含分片集群)上用,因为它读的就是 oplog。还在跑单机 mongod 的服务,先迁移部署形态再谈订阅。
漏单事故的完整过程值得展开。用户支付成功后,订单服务的写入路径是先落库、再由定时任务每 10 分钟把增量刷进缓存,客服系统只读缓存。这个架构在低峰完全无感,问题藏在高峰:同步任务全量扫描最近变更时耗时从 40 秒涨到 11 分钟——超过了轮询周期本身,下一轮任务与上一轮重叠,同步开始跳批次,黑洞从 10 分钟扩大到半小时。客服投诉升级的那天,正是大促把变更速率翻了三倍的日子。
定位过程分两步。第一步对账:抽样一千笔订单,比对库与缓存的创建时间差,画出分布直方图,P50 差 4 分钟、P99 差 31 分钟——数据一出,问题从"偶发投诉"变成"系统性窗口"。第二步选型:团队在"缩短轮询周期"与"换推送"之间摇摆,缩短周期治标且同步任务重叠问题依旧;推送方案动架构但根治。用 Change Streams 改造花了三天:消费端独立部署、令牌持久化到专用集合、消费逻辑幂等覆盖写、令牌过期时降级为增量对账(按 updatedAt 补最近一小时)。上线后实测库到缓存的平均延迟 80 毫秒,P99 不超过 2 秒,投诉归零。
// 生产级消费者的骨架:令牌持久化 + 异常降级 + 心跳 async function consume() { let token = await db.tokens.findOne({ name: "orders-cache" }); const cs = db.orders.watch( [{ $match: { operationType: { $in: ["insert", "update"] } } }], token ? { resumeAfter: token.v, fullDocument: "updateLookup" } : { fullDocument: "updateLookup" }); try { while (await cs.hasNext()) { const ev = await cs.next(); await cache.set(ev.documentKey._id.toString(), ev.fullDocument); await db.tokens.updateOne({ name: "orders-cache" }, { $set: { v: ev._id, at: new Date() } }, { upsert: true }); } } catch (e) { // 令牌过期(oplog 截断)或连接中断:退化为按时间增量对账后重启 await reconcileByUpdatedHour(); } }
Change Streams 不是免费的。消费端跟不上时事件会在内部缓冲堆积,内存上涨最终断连;oplog 压缩期如果消费端长时间停摆,令牌失效后只能全量重建。两个经验数字:单消费者的处理能力在普通业务文档上约每秒五千事件,超过就该按文档键分片开多个消费者;令牌的有效期约等于 oplog 窗口,所以第 6 章讲的 oplog 容量规划在这里直接兑现——消费端可能停摆的最长时间,必须小于 oplog 窗口。另外 watch 的 $match 放在服务端执行,能过滤掉不关心的事件,别把所有事件拉到客户端再丢。
| 场景 | 方案 | 理由 |
|---|---|---|
| 缓存刷新 | Change Streams 加幂等覆盖 | 毫秒级联动,黑洞消失 |
| 审计同步 | Change Streams 加专用账号 | 只读订阅,权限可最小化 |
| 跨系统对账 | 定时全量对账保留为兜底 | 推送为主、对账为辅的双保险 |
| 单机实例 | 先迁移复制集 | 订阅依赖 oplog,无替代 |