Change Streams是什么:从Oplog到实时事件流
要理解Change Streams,得先从MongoDB的底层机制说起。MongoDB副本集内部维护着一个操作日志,也就是Oplog,主节点上所有写入操作都会按顺序记录在这份日志中,从节点正是通过消费Oplog来保持数据一致的。Change Streams本质上就是对Oplog做了一层封装,把原本面向复制机制的底层日志,转换成了面向应用层的、结构化的变更事件流。
在Change Streams出现之前,如果应用想感知数据库的数据变化,通常只能靠轮询,也就是定时扫描表里是否有新增或修改的记录。这种方式不仅延迟高、对数据库压力大,而且很难捕捉到删除操作。有了Change Streams之后,应用只需打开一个游标监听指定集合,任何插入、更新、删除、替换事件都会实时推送过来,事件中带有完整的文档快照、变更字段描述、命名空间信息以及操作时间戳,信息量非常充分。
一个最基础的监听代码如下,这里以Node.js驱动为例:
const { MongoClient } = require("mongodb");
const client = new MongoClient("mongodb://127.0.0.1:27017");
async function watchOrders() {
await client.connect();
const collection = client.db("shop").collection("orders");
// 监听orders集合的变更事件
const changeStream = collection.watch([], { fullDocument: "updateLookup" });
changeStream.on("change", (event) => {
console.log("操作类型:", event.operationType);
console.log("文档数据:", event.fullDocument);
});
}
watchOrders().catch(console.error);需要注意的是,Change Streams要求MongoDB以副本集模式运行,单机standalone模式不支持。本地开发可以用单节点副本集来模拟,生产环境一般本来就是副本集或分片集群,所以门槛并不高。此外监听器还可以针对整个数据库甚至整个集群开启,通过管道聚合对事件做前置过滤和加工,把不关心的事件在数据库端就过滤掉,减少网络传输和应用端处理压力。
典型应用场景:事件驱动架构中的落地方案
场景一:缓存失效与最终一致性。这是最常见也最实用的场景。当订单数据在MongoDB中更新后,Redis里缓存的旧数据如果不及时失效,用户就会看到过期信息。传统做法是在业务代码里写数据库之后手动删缓存,但这段逻辑容易遗漏,一旦某个新的写入路径忘了处理,脏数据问题就出现了。用Change Streams可以把缓存失效逻辑从业务代码中完全剥离出来,由独立的监听服务统一处理:
changeStream.on("change", async (event) => {
const docId = event.documentKey._id;
// 数据变更后删除对应缓存,下次读取时重建
await redis.del(`order:${docId}`);
});场景二:跨服务数据同步与只读副本构建。在微服务架构中,订单服务、报表服务、搜索服务往往需要同一份数据的不同视图。传统方案要么让各服务直接读订单库(造成强耦合),要么引入消息队列做双向同步(链路复杂)。Change Streams可以作为轻量级的同步通道,把变更事件实时写入Elasticsearch构建搜索索引,或者同步到分析库供报表使用,各消费方互不影响,彼此解耦。
场景三:实时通知与业务触发。比如电商系统里订单状态一旦变为已支付,就需要触发发货流程、给用户推送消息、给财务系统记账。通过监听订单集合的update事件,判断更新后的status字段值,就能以事件方式驱动下游流程,业务代码不需要到处埋点调用通知接口。此外,Change Streams也非常适合做审计日志,所有写入操作的来源、内容、时间都有据可查,满足合规要求。
生产环境关键细节:断点续传、过滤与可靠性设计
在真实项目中直接监听变更流是远远不够的,首先要解决的是断点续传问题。服务重启后,如果从当前时刻开始监听,重启期间的变更就全部丢失了。Change Streams的每个事件都带有resumeToken(恢复令牌),消费方应该把成功处理的事件令牌持久化下来,重启后从上次的令牌位置继续消费:
async function watchWithResume(tokens) {
const lastToken = await tokens.get("orders");
const options = {
fullDocument: "updateLookup",
// 从上次记录的resumeToken继续消费
resumeAfter: lastToken || undefined
};
const stream = collection.watch([], options);
stream.on("change", async (event) => {
await processEvent(event);
// 处理成功后再保存令牌,保证至少一次语义
await tokens.save("orders", event._id);
});
}其次要善用聚合管道做事件过滤。与其把所有事件都拉到应用端再判断,不如在watch的第一个参数里传入匹配管道,让数据库端先过滤。比如只关心订单状态从待支付变为已支付的事件,可以直接用$match配合$expr做服务端过滤,事件量可能从每秒上千条降到几十条,资源消耗差别巨大。
还有一个容易被忽视的细节是fullDocument的取值策略。默认情况下update事件只包含变更的字段增量,而updateLookup选项会在事件产生时额外查一次当前完整文档。这带来一个微妙问题:查询到的文档是事件发生后的最新状态,如果同一文档被连续更新多次,拿到的快照可能已经超前于当前事件。对一致性要求高的场景,MongoDB 6.0之后支持fullDocumentBeforeChange,可以拿到变更前的完整文档,配合前置镜像实现更严谨的处理逻辑。
与消息队列方案的对比。Change Streams的天然优势是不需要额外维护Kafka或RabbitMQ这类中间件,数据在库里的每一次变更都会被捕捉,包括直接通过mongoshell执行的运维操作,事件不会因为应用层遗漏而丢失。但它的不足也很明显:没有消息队列那种消费者组横向扩容能力,多个消费者监听同一个Change Stream会各自收到全部事件,需要自己实现分片分发;同时Oplog有容量上限,消费滞后太久令牌可能失效。因此比较务实的做法是,中小规模场景直接用Change Streams驱动事件架构,大规模高频事件场景则可以用Change Streams作为采集端,把事件再投递到Kafka,兼得两者的长处。
MongoDBChange Streams事件驱动架构修改时间:2026-09-06 02:33:26