MongoDB Change Streams 是一种官方提供的变更监听机制,它把数据库中的每次写入操作转换成可订阅的事件流。与轮询查询或手动解析 oplog 相比,Change Streams 延迟更低,也更容易处理断线恢复。事件流基于副本集的 oplog 构建,因此天然具备顺序性和持久性。理解事件类型、游标选项和 resume token 的用法,是可靠使用该特性的关键。

一、Change Streams 的工作机制与事件类型
Change Streams 的核心原理并不复杂。MongoDB 副本集中的每个数据节点都会记录 oplog(操作日志),用于主从同步。主节点执行写操作后,会先生成 oplog 条目,再完成数据变更。Change Streams 正是从这些 oplog 条目中提取变更信息,转换成面向应用的事件。因此,只要写入成功,事件就一定会出现,而且顺序与写入顺序一致。
事件类型主要包括 insert、update、delete、replace 和 invalidate。insert 事件表示新文档插入,documentKey 和 fullDocument 字段携带数据;update 事件包含 updateDescription,说明哪些字段发生了变化;delete 事件只给出 documentKey,因为文档已被删除;replace 事件在整篇文档被替换时出现;invalidate 事件则常见于集合被删除或重命名等导致游标失效的场景。对于分片集群,事件可能来自多个分片,但 MongoDB 会尽量保证全局有序。
每个事件都带有一个 _id 字段,叫做 resume token。它标识了该事件在 oplog 中的位置。应用如果因为网络抖动或进程重启导致游标断开,可以拿着上次保存的 resume token 重新打开监听,从断点继续消费,不会丢失中间的事件。这个特性比传统的 tail -f 解析日志可靠得多。
const { MongoClient } = require('mongodb');
async function openChangeStream() {
const client = new MongoClient('mongodb://localhost:27017');
await client.connect();
const db = client.db('shop');
const orders = db.collection('orders');
// 打开集合级别 Change Stream,返回游标
const changeStream = orders.watch([], {
fullDocument: 'updateLookup'
});
changeStream.on('change', (change) => {
console.log('收到事件:', change.operationType);
console.log('文档键:', change.documentKey);
// 将 resume token 保存到 Redis 或文件中,便于恢复
saveResumeToken(change._id);
});
changeStream.on('error', (error) => {
console.error('监听出错:', error);
});
}
二、实战:连接、过滤与断线恢复
先看如何用 Node.js 驱动打开一个集合级别的 Change Stream。上面的示例中 watch 方法接收两个参数:pipeline 和 options。pipeline 是一个聚合管道,用来过滤事件;options 可以配置 fullDocument、startAtOperationTime 等。fullDocument 设置为 updateLookup 时,update 事件会带上更新后的完整文档,默认情况下 update 事件不包含完整文档。
如果只想监听 insert 和 update,而不关心 delete,可以使用 $match 过滤。例如:
const pipeline = [
{ $match: { operationType: { $in: ['insert', 'update'] } } }
];
const changeStream = orders.watch(pipeline, {
fullDocument: 'updateLookup'
});
这会让游标只返回 insert 和 update 事件,其他操作被直接过滤,降低了网络和计算开销。管道里还可以根据字段内容做过滤,例如只关注 total 大于 200 的订单变更:
const pipeline = [
{ $match: { 'fullDocument.total': { $gt: 200 } } }
];
断线恢复是生产环境必须处理的问题。变更流游标可能因为网络闪断、主从切换或 oplog 过期而失效。合理的做法是监听 error 或 close 事件,捕获最后收到的 resume token,然后重新打开游标。以下是一个带自动重连的简化示例:
const { MongoClient } = require('mongodb');
let resumeToken = null;
let client;
async function startListening() {
client = new MongoClient('mongodb://localhost:27017');
await client.connect();
const orders = client.db('shop').collection('orders');
const options = { fullDocument: 'updateLookup' };
if (resumeToken) {
options.resumeAfter = resumeToken;
}
const changeStream = orders.watch([], options);
changeStream.on('change', (change) => {
resumeToken = change._id;
console.log('处理事件:', change.operationType);
});
changeStream.on('error', async (error) => {
console.error('连接中断,准备重连:', error);
await startListening();
});
changeStream.on('close', async () => {
console.log('游标关闭,准备重连');
await startListening();
});
}
这个示例将 resume token 保存在内存变量中,实际项目中应该放到 Redis、数据库或文件中。重新打开游标时传入 resumeAfter,MongoDB 会从该 token 之后继续推送事件。需要注意的是,如果 oplog 已经推进得太远,resume token 可能失效,此时需要根据业务情况进行全量补偿。
三、生产环境中的部署要求与性能优化
Change Streams 只能在副本集和分片集群上使用,单机模式的 MongoDB 不支持。这是因为底层依赖 oplog,而只有副本集才会维护 oplog。在开发环境如果想快速体验,可以把单实例配置成单节点副本集,启动参数中加入 --replSet。对于分片集群,由于涉及多个分片和配置服务器,建议使用 MongoDB 3.6 以上版本,并确保各分片之间的时间同步。
oplog 的大小直接决定了断线后还能恢复多长时间的事件。默认 oplog 大小可能只有磁盘容量的 5% 左右,在高写入场景下很快被覆盖。如果事件消费程序宕机超过 oplog 保留窗口,resume token 就会失效。可以通过 replSetResizeOplog 命令扩大 oplog,或者尽早持久化 token 并尽快恢复监听。
性能方面,Change Streams 的监听本身对数据库的压力不大,但每个游标都会占用一定的服务器资源。如果多个服务需要监听同一集合,建议复用游标,或者用消息队列将事件分发出去。事件处理必须幂等,因为网络重连可能导致少量重复事件。update 事件默认只返回变更字段,通过 updateLookup 获取完整文档会增加一次查找,对数据库有一定影响,按需开启。
还需要注意,Change Streams 的游标长时间不消费时,MongoDB 会认为客户端已经不活跃,可能自动关闭游标。应用层应实现心跳或定期检查游标状态。对于删除和重命名集合、删除数据库等 DDL 操作,会触发 invalidate 事件,普通变更流游标将无法继续使用,这时需要重新评估监听范围。
四、Change Streams 与消息队列的配合
Change Streams 解决的是数据库变更的实时捕获问题,但它本身不承担消息持久化和广播职责。企业架构中常见做法是:一个或多个消费者监听 Change Streams,把事件转换成业务消息后投递到 Kafka、RabbitMQ 或 RocketMQ,再由下游服务订阅。这种组合兼顾了数据库变更的实时性和消息队列的解耦能力。
如果直接将 Change Streams 暴露给多个业务方,每个业务方都打开游标会加重数据库连接负担。通过消息队列中转后,数据库只需维护少量游标,下游按需订阅 topic。事件投递时建议保留 resume token 作为消息元数据,方便对账和审计。同时要注意消息的幂等消费,可以以 event._id 作为去重键。
另一种简化方案是只用一个监听服务,把事件写入 Redis Streams。代码逻辑与直接消费类似,只是 change 事件回调中改为发布消息。这种架构在中小规模项目中容易落地,维护成本低。无论采用哪种方式,Change Streams 都提供了稳定、有序的事件源,让实时数据同步、缓存失效、通知推送等功能变得简洁可靠。
MongoDB Change Streams实时监听数据变化修改时间:2026-08-19 11:43:52