MongoDB Change Streams 订阅变更
MongoDB Change Streams:实时数据变更订阅完全指南
在构建现代应用时,能够实时感知数据库变化对上至数据同步、下至触发业务逻辑都至关重要。MongoDB 的 Change Streams 正是为此而生——它允许应用程序像订阅消息队列一样,监听集合、数据库乃至整个部署的变更事件,且无需任何额外的轮询。
本教程将从零开始,带你掌握 MongoDB Change Streams 的核心概念、操作方法以及生产环境下的注意事项。
什么是 Change Streams?
Change Streams 是 MongoDB 从 3.6 版本(副本集与分片集群)开始引入的一项功能。它基于 oplog(操作日志)实现,能够将数据库中的变更以事件流的形式推送给订阅者。你完全不用关心底层 oplog 的复杂机制,只需要使用 MongoDB 驱动提供的高级 API。
简单来说,它就像数据库的“观察者模式”:
- 有人插入了一条文档 → 你立刻收到
insert事件。 - 有字段被修改 → 你收到包含
updateDescription的update事件。 - 文档被删除 → 你收到
delete事件。
必备前提
在开始编写代码之前,请确认你的环境满足以下条件:
- 部署类型:Change Streams 只能工作在 副本集 或 分片集群 上。单机
mongod进程无法使用。 - 存储引擎:必须使用 WiredTiger 存储引擎,MongoDB 3.6 以上已默认启用。
- 权限:用于监听变更的用户需要拥有
changeStream权限。 - 数据库:若监听的是集合,数据库需已存在;监听数据库或集群级别则无此限制。
基本用法(Node.js 示例)
我们将使用官方的 Node.js 驱动(版本 3.6+)进行演示。其他语言的驱动(Python、Java、C# 等)思路完全一致,只是 API 名称略有不同。
1. 连接并获取监听目标
const { MongoClient } = require('mongodb');
async function main() {
const uri = 'mongodb://localhost:27017/?replicaSet=myReplSet';
const client = new MongoClient(uri);
try {
await client.connect();
const database = client.db('shop');
const collection = database.collection('orders');
// 打开 Change Stream,监听 orders 集合的所有变更
const changeStream = collection.watch();
// 处理事件
changeStream.on('change', (changeEvent) => {
console.log('收到变更事件:', JSON.stringify(changeEvent, null, 2));
});
// 可选:监听错误
changeStream.on('error', (err) => {
console.error('Change Stream 出错:', err);
});
// 程序不退出
console.log('开始监听 shop.orders 的变更...');
await new Promise(resolve => setTimeout(resolve, 60000)); // 演示保持监听
} finally {
await client.close();
}
}
main().catch(console.dir);
运行该脚本后,在 shop.orders 中执行任意增删改操作,控制台都会打印出完整的变更事件。
2. 变更事件结构解析
一个典型的 insert 事件如下:
{
"_id": {
"_data": "8264..." // 恢复令牌 Resumable Token
},
"operationType": "insert",
"clusterTime": { "$timestamp": 1710000000 },
"fullDocument": {
"_id": 1,
"customer": "Alice",
"amount": 99.9
},
"ns": {
"db": "shop",
"coll": "orders"
},
"documentKey": { "_id": 1 }
}
重要字段说明:
operationType:区分insert、update、delete、replace、drop等。fullDocument:仅insert和replace默认包含完整文档。对于update事件,需要通过配置fullDocument: 'updateLookup'来获取更新后的完整文档。updateDescription:update事件特有,包含updatedFields和removedFields,方便精准了解哪部分被修改。ns:发生变更的数据库和集合。_id._data:至关重要的恢复令牌,用于故障恢复。
高级过滤与配置
watch() 方法支持传入一个 pipeline 数组,其中可以放置聚合管道中的 $match 阶段进行过滤。这比在应用层过滤更高效,因为 MongoDB 只推送你关心的变更。
// 只监听插入操作,且金额大于 100 的订单
const pipeline = [
{ $match: { 'operationType': 'insert' } },
{ $match: { 'fullDocument.amount': { $gt: 100 } } }
];
const changeStream = collection.watch(pipeline);
watch 还支持选项对象:
fullDocument: 'default'(默认):update 事件只返回updateDescription。fullDocument: 'updateLookup':update 事件额外返回当前最新版完整文档(通过_id后置查询获取,有轻微性能开销)。fullDocumentBeforeChange: 'whenAvailable'(MongoDB 6.0+):在 replace、update、delete 事件中携带变更前的完整文档,非常适合审计场景。startAtOperationTime: 从指定的时间戳开始监听,用于断点续传。
恢复与容错:利用 Resumable Token
Change Stream 可能因为网络抖动、应用重启或被服务端主动关闭(如运维操作)而中断。此时必须使用保存的恢复令牌重新开启流,避免丢失事件。
let resumeToken;
changeStream.on('change', (change) => {
// 每次收到事件都更新 token
resumeToken = change._id;
// 处理业务逻辑...
});
changeStream.on('close', () => {
console.log('Change Stream 关闭,尝试使用 token 恢复...');
restartStream(resumeToken);
});
async function restartStream(token) {
const newStream = collection.watch([], {
resumeAfter: token // 从上次中断的下一个事件开始
});
// 将 newStream 重新绑定到原有的事件监听逻辑中
}
最佳实践:将 resumeToken 持久化(如存入 Redis 或文件),即使应用完全崩溃,重启后仍能从上次位置继续。
监听级别:集合 vs 数据库 vs 集群
- 集合级别
collection.watch():仅限于该集合的变更。 - 数据库级别
database.watch():该库下任意集合的变更,事件中包含ns字段区分来源集合。 - 集群级别
client.watch():整个 MongoDB 集群(所有库)的变更,通常需要更高的权限。
数据库级别监听常用于需要捕捉新建集合后立即监听其变更的场景(例如 capture create 事件)。
使用场景示例
- 实时数据同步:将 MongoDB 变更近乎实时地同步到 Elasticsearch、Datalake 或另一个 MongoDB。
- 事件驱动架构:用户下单后,订单服务使用 Change Stream 触发通知服务发送短信,无需修改订单模块代码。
- 审计日志:结合
fullDocumentBeforeChange记录每一次字段级别的变动历史。 - 缓存失效:当某份文档被修改时,立即通过流事件通知缓存层使其失效。
注意事项与限制
- oplog 窗口限制:Change Stream 依赖 oplog,如果恢复令牌对应的 oplog 记录已过期(被滚动删除),则无法恢复。必须设置足够大的 oplog 以保证容灾能力。
- 性能考量:每个 Change Stream 都会在服务端占用资源。避免同时打开数量过多的流。对于高吞吐写入,可考虑用一个流驱动多个消费者(扇出模式)。
- 事务中的变更:在分片集群上,事务内的变更会等到事务提交后才出现在 Change Stream 中,且
clusterTime反映的是提交时间。 - 删除集合:当监听的集合被删除 (
drop) 时,Change Stream 会关闭。若需要继续监听同名集合,程序应在收到drop事件后重新调用watch()。 - MongoDB Atlas:Atlas 提供了完全托管的 Change Streams,并支持通过 Eventbridge 等集成到无服务器架构中。
快速上手检验清单
- 确认部署为副本集或分片集群。
- 连接字符串中指定
replicaSet参数。 - 为你的应用账户授予
changeStream和对应集合的find权限。 - 在代码中实现恢复逻辑,并测试中断场景。
- 根据业务选择合适的
fullDocument选项和过滤管道。
总结
MongoDB Change Streams 将数据库变更从被动查询转为主动推送,极大地简化了实时数据处理系统的架构。只需几行代码,你就能构建出可靠的变更订阅管道。记住善用过滤、妥善处理容错恢复,它将成为你数据工具箱中最可靠的一环。