MongoDB Change Streams 是什么?适合做什么?
简化版
MongoDB Change Streams 可以让应用订阅集合、数据库或集群中的数据变更事件,比如 insert、update、delete。它基于复制机制,常用于缓存更新、搜索索引同步、审计、实时通知和数据管道。使用时要处理断点续传、重复消费、消费延迟和事件顺序等问题。
详细版
典型用法:
const changeStream = db.collection("orders").watch();
for await (const change of changeStream) {
console.log(change.operationType, change.documentKey);
}
常见事件:
insertupdatereplacedeleteinvalidate
Change Streams 适合监听数据库变更,但它不是完整业务消息系统的替代品。核心业务事件仍应由业务服务明确发布,数据库变更监听更适合做同步和观察。
完整版教学
一、Change Streams 解决什么问题
很多系统需要知道数据库发生了变化:
- 订单变化后刷新缓存;
- 商品变化后同步搜索引擎;
- 用户资料变化后同步画像;
- 数据变化后写审计日志;
- 后台大屏实时展示。
以前可能需要轮询数据库。Change Streams 提供了订阅变更的机制,应用可以实时消费变化。
二、Change Streams 依赖复制机制
Change Streams 通常要求部署在支持复制的环境中。它利用 MongoDB 的变更记录能力把数据修改以事件形式暴露给客户端。
这意味着它和单机临时脚本不同,更适合生产复制集或分片集群环境。
三、可以监听不同范围
可以监听集合:
db.collection("orders").watch()
也可以监听数据库或集群级变化,具体取决于驱动和权限配置。
监听范围越大,事件越多,过滤和消费压力也越大。实际项目通常只监听关心的集合,并在管道中尽早过滤。
四、断点续传很重要
消费者可能重启、网络中断或处理失败。Change Streams 提供 resume token 机制,应用可以记录消费位置,恢复后从上次位置继续。
如果不保存断点,服务重启后可能漏事件或重复处理。
工程上通常要:
- 保存 resume token;
- 消费逻辑幂等;
- 处理重复事件;
- 监控消费延迟;
- 对失败事件进入重试或死信流程。
五、Change Streams 不是业务消息的万能替代
数据库变更事件描述的是“数据发生了什么变化”,但业务事件描述的是“业务上发生了什么事情”。
例如订单状态从待支付变成已支付,Change Streams 能看到字段变化,但业务事件可能还包含支付渠道、营销归因、风控结果等上下文。
核心业务链路中,推荐由业务服务明确发布领域事件;Change Streams 更适合同步缓存、索引、审计和数据管道。
六、消费语义要按“至少一次”思路设计
Change Streams 消费端重启、网络断开或处理超时后,可能用 resume token 从某个位置继续读。工程上更稳的假设是:事件可能重复到达,消费者必须幂等;事件也可能因为 token 过旧而无法恢复,需要报警和人工补偿。
MongoDB 变更 -> Change Stream -> 消费服务 -> 缓存/搜索/审计
|
保存 resume token
比如同步搜索索引时,可以用文档 _id 作为幂等键,重复收到同一条 update 事件时覆盖写入搜索引擎,而不是追加一条新记录。这样即使同一事件处理 2 次,也不会产生重复数据。
七、常见误区与追问
| 使用场景 | 是否适合 Change Streams | 原因 |
|---|---|---|
| 缓存失效 | 适合 | 数据变化后刷新或删除缓存 |
| 搜索索引同步 | 适合 | 根据变更更新外部索引 |
| 订单支付领域事件 | 不建议完全替代 | 业务上下文可能不在数据库变更里 |
| 审计观察 | 适合但要补上下文 | 能看到变更,责任信息要业务补充 |
记忆钩子:Change Streams 看见的是“数据怎么变”,业务消息表达的是“业务发生了什么”。二者能配合,但不能简单互相替代。
- 误区:Change Streams 可以完全替代 MQ。 它适合同步和观察,核心业务事件仍应由业务服务明确发布,避免丢失业务语义。
- 误区:只要监听集合就不会漏数据。 消费端还要保存 resume token、监控延迟,并处理 token 过期或 oplog 窗口不足的问题。
- 误区:事件一定只消费一次。 网络重试和恢复可能导致重复事件,消费逻辑必须按幂等设计。
- 追问:resume token 有什么用? 它记录变更流位置,消费者重启后可以从上次位置继续,降低漏处理风险。
- 追问:监听范围越大越好吗? 范围越大事件越多,过滤压力和权限风险越高,实际项目通常只监听关心集合并尽早过滤。
- 追问:消费延迟大怎么办? 要监控 lag、扩容消费者、优化下游写入,必要时暂停非关键同步或进入补偿流程。
八、加强记忆
MongoDB Change Streams 记住“订阅数据变化”。它适合缓存刷新、搜索同步、审计和实时数据管道;生产使用要保存 resume token,消费要幂等,并且不要把它误当成所有业务消息的替代品。