MongoDB Change Streams 提供了一种事件驱动型的数据变更监听方案,它让应用能够以订阅的方式实时获取集合中的插入、更新、删除等操作,而不是每隔一段时间去扫描数据表。相比传统轮询,Change Streams 的延迟更低,同时能够保留事件之间的顺序,并为断线恢复提供 resume token 机制。这个能力建立在 MongoDB 的复制日志之上,因此对集群拓扑有一定要求。要把它用在生产环境里,需要理解事件结构、监听器的创建方式以及断点续传的实现方法。

一、从 oplog 到 Change Streams:它到底依赖什么
MongoDB 的副本集通过复制日志 oplog 来同步主节点和从节点的数据。oplog 是一条特殊的 capped collection,每条记录对应一次数据变更操作,比如插入、更新或删除。主节点执行写操作后,会把变更追加到 oplog,从节点则持续拉取这些记录并重放,从而保持数据一致。Change Streams 的本质就是在这套复制机制上提供对外的事件流接口,驱动内部会跟踪 oplog 中的变更,并将其转换成结构化的通知。
虽然都基于 oplog,但 Change Streams 和直接读 oplog 有明显区别。oplog 记录的是内部复制格式,字段和结构可能随版本变化而调整,不适合应用直接消费。Change Streams 把这些细节封装起来,提供稳定的事件模型,同时支持聚合管道过滤、权限控制以及 resume token。因此即使你对 oplog 很熟悉,也建议把 Change Streams 作为应用层的数据变更入口。
从部署条件看,Change Streams 要求数据库必须是副本集或分片集群,单机 mongod 无法使用。这是因为只有存在 oplog 的拓扑才有复制日志可供订阅。同时,多数生产环境还需要开启多数读关注,确保监听到的事件已经提交到大多数节点,避免读到回滚的数据。若使用分片集群,建议在 mongos 上打开 Change Streams,由驱动和 mongos 统一协调各分片的事件顺序。
二、监听集合数据变化:Node.js 与 Java 示例
以 Node.js 驱动为例,创建一个针对 orders 集合的变更监听器非常直接。先建立 MongoClient 连接,然后调用集合的 watch 方法获得 ChangeStream 对象,并通过 change 事件处理变更。下面是完整示例。
const { MongoClient } = require('mongodb');
async function watchOrders() {
const client = new MongoClient('mongodb://localhost:27017/?replicaSet=rs0');
await client.connect();
const db = client.db('shop');
const orders = db.collection('orders');
const stream = orders.watch();
stream.on('change', (event) => {
console.log('操作类型:', event.operationType);
console.log('目标文档:', event.documentKey);
if (event.fullDocument) {
console.log('完整文档:', event.fullDocument);
}
});
}
watchOrders().catch(console.error);在这个例子中,watch 方法没有传任何参数,表示监听该集合上的所有操作。每个变更事件都包含 operationType、documentKey、ns 等字段,其中 documentKey 是被修改文档的 _id,ns 是命名空间。完整文档 fullDocument 只在 insert 和 replace 事件里默认出现,对于 update 事件,如果希望拿到更新后的完整文档,需要在 watch 参数里加上 fullDocument 选项并设为 updateLookup。
Java 开发者也使用类似的 API。通过 MongoCollection 的 watch 方法拿到游标,然后遍历事件。
import com.mongodb.client.*;
import com.mongodb.client.model.changestream.ChangeStreamDocument;
import org.bson.Document;
public class ChangeStreamDemo {
public static void main(String[] args) {
String uri = "mongodb://localhost:27017/?replicaSet=rs0";
try (MongoClient mongoClient = MongoClients.create(uri)) {
MongoDatabase database = mongoClient.getDatabase("shop");
MongoCollection<Document> collection = database.getCollection("orders");
MongoChangeStreamCursor<ChangeStreamDocument<Document>> cursor = collection.watch().cursor();
while (cursor.hasNext()) {
ChangeStreamDocument<Document> event = cursor.next();
System.out.println("操作类型: " + event.getOperationType());
System.out.println("文档键: " + event.getDocumentKey());
}
}
}
}两种语言的监听逻辑本质相同,都是通过游标不断拉取事件。实际项目中,通常会把事件转换成内部消息对象,异步投递到队列或更新缓存,而不是在回调里执行耗时操作。
三、用聚合管道过滤事件:只监听真正关心的变更
如果集合写入频繁,但应用只关心特定类型或字段的变化,直接监听所有事件会浪费不少资源。Change Streams 支持在 watch 方法中传入聚合管道,利用 $match 等阶段对事件进行过滤。比如,只想监听 insert 操作,并且新插入的订单金额大于 100,可以这样写:
const pipeline = [
{
$match: {
$and: [
{ operationType: 'insert' },
{ 'fullDocument.totalAmount': { $gt: 100 } }
]
}
}
];
const stream = orders.watch(pipeline);
stream.on('change', (event) => {
console.log('新增大额订单:', event.fullDocument);
});管道过滤发生在数据库端,不满足条件的事件不会推送到应用,能有效降低网络传输和客户端处理压力。常用的过滤条件包括 operationType、databaseName、namespace、fullDocument 中的字段等。$match 支持大多数 MongoDB 查询运算符,但不允许使用 $out、$merge 这类写阶段。
需要注意的是,更新事件的 fullDocument 默认不存在,若要针对更新后的内容过滤,必须设置 fullDocument 为 updateLookup。updateLookup 会从源集合查找更新后的文档,并填充到 fullDocument,这会在每次 update 事件时额外执行一次查询,增加数据库开销。因此建议仅在需要时启用,并尽量把过滤条件写得精确。
另一个常用技巧是使用 $project 阶段裁剪事件字段。例如只保留 operationType、documentKey 和 fullDocument,避免把 clusterTime、ns 等冗余信息传给下游。合理使用管道可以让事件流更轻量,也更容易维护。
四、断点续传:让监听器在重启后不丢事件
应用重启、网络抖动或临时异常都可能导致 Change Streams 连接中断。随意重新建立监听很容易丢失中断期间的事件。MongoDB 为每个变化事件分配了唯一的 resume token,它就是事件中的 _id 字段。应用只要在处理事件后保存这个 token,就能在下次启动时从断点继续消费。
let resumeToken = null;
const stream = orders.watch();
stream.on('change', (event) => {
resumeToken = event._id;
// 将 resumeToken 保存到 Redis 或本地文件
console.log('事件已处理,resumeToken:', resumeToken);
});
// 重启或重连后
const newStream = orders.watch([], { resumeAfter: resumeToken });resumeAfter 要求 token 仍然有效,即对应事件还在 oplog 范围内。如果 token 太旧,oplog 空间可能已被覆盖,此时会收到错误。为了降低这种情况,一方面可以适当扩大 oplog 大小,另一方面可以实现无 token 时的全量对账逻辑,或将 token 持久化到高可靠的存储中。
生产环境还需要考虑错误恢复。驱动在遇到可恢复错误时会自动尝试重连,但应用应监听 error 事件并记录日志。较新的 MongoDB 驱动支持 resumeAfter 和 startAfter 两个参数,前者从指定事件之后继续,后者从指定事件开始继续。对于大多数场景,使用 event._id 作为保存值并配合 resumeAfter 就足够了。如果应用是多实例部署,务必确保同一集合的多个实例不会重复消费或互相竞争,可借助分布式锁或给每个实例分配不同分片键。
五、性能调优与常见误区
Change Streams 虽然降低了轮询成本,但并不意味着可以无节制地创建监听器。每个打开的 change stream 都会占用数据库连接和游标资源,大量监听时可能触发 mongod 的游标限制。建议对不活跃的监听器设置合理的 maxAwaitTime,避免游标长期挂起。同时,尽量在客户端做异步处理,将事件快速投递到内存队列,由独立线程消费,否则 change stream 的回调耗时过长会拖慢整个游标推进。
另一个常见误区是认为 Change Streams 能完全替代消息队列。Change Streams 主要解决数据库变更感知问题,不负责跨系统投递、重试策略、消费组管理等复杂语义。在需要多个下游系统订阅同一事件流的场景里,比较稳妥的方式是先由一个服务消费 Change Streams,再把事件发布到 Kafka、RabbitMQ 等中间件,由中间件完成分发。
此外,update 事件中的 updateDescription 提供了 updatedFields 和 removedFields,它们能帮助应用只关心被修改的字段,而不必每次都读取完整文档。合理利用 updateDescription 可以进一步减少数据库查询。例如在订单状态变更时,只监听 update 并且 updatedFields.status 等于某个值,从而触发后续通知逻辑。
const pipeline = [
{
$match: {
$and: [
{ operationType: 'update' },
{ 'updateDescription.updatedFields.status': 'shipped' }
]
}
}
];
const stream = orders.watch(pipeline, { fullDocument: 'updateLookup' });
stream.on('change', (event) => {
console.log('订单已发货,单号:', event.fullDocument.trackingNo);
});通过这种精确过滤,可以避免大量无关事件进入处理链路,也能让代码意图更加清晰。
MongoDB Change Streams数据变化监听实时数据同步修改时间:2026-09-25 09:54:33