Change Streams 基于 oplog 的异步事件模型
MongoDB Change Streams 并不是凭空产生的功能,它建立在复制集的 oplog 之上。oplog 是主节点记录所有写操作的特殊集合,从节点通过拉取 oplog 来保持数据同步。Change Streams 订阅被监视集合对应的 oplog 条目,并把它们转换为可读的事件流,通过游标推送给应用程序。这意味着变更事件是异步产生的,原始写操作已经在主节点提交,事件只是其后续的投影。

与数据库内置触发器相比,这一机制有本质区别。传统触发器(例如关系型数据库中的 BEFORE/AFTER INSERT 触发器)在写入路径上同步执行,与触发它的 SQL 语句共享同一事务边界。如果触发器抛出异常,原始写入会被回滚。而 Change Streams 的消费者运行在数据库外部,它接收到的只是已提交的写操作通知,即使消费者处理失败,也无法撤销源操作。这种异步模型降低了数据库的耦合度,但带来了最终一致性的特性,适用于事件驱动架构,不适合需要强一致性的场景。
开发者可以通过 MongoDB 驱动程序订阅集合、数据库或整个部署的变更。常见的事件类型包括 insert、update、replace、delete、invalidate 等。对于 update 操作,默认事件中不包含更新后的完整文档,需要设置 fullDocument 为 updateLookup 才会回查并返回最新文档。下面是一段 Node.js 监听订单集合的示例代码。
const { MongoClient } = require('mongodb');
async function main() {
const client = new MongoClient('mongodb://localhost:27017');
await client.connect();
const db = client.db('shop');
const orders = db.collection('orders');
const changeStream = orders.watch([], { fullDocument: 'updateLookup' });
changeStream.on('change', function(change) {
console.log(change.operationType, change.fullDocument);
});
await new Promise(function(resolve) { setTimeout(resolve, 60000); });
await client.close();
}
main().catch(console.error);
这段代码创建了一个变更流,并监听 orders 集合上的任何写操作。事件回调中打印操作类型和完整文档。实际运行时,回调会在写入发生后的毫秒到秒级内被触发,具体延迟取决于网络、oplog 读取速度以及消费者的处理能力。需要注意的是,Change Streams 要求 MongoDB 部署为副本集或分片集群,单机模式无法使用该功能。
用 Change Streams 实现类触发器的典型场景
虽然 Change Streams 不能完全替代数据库触发器,但在许多场景下它可以安全地接管原本由触发器承担的职责。常见的场景包括审计日志、跨集合数据同步、缓存失效、搜索引擎索引更新以及向外部系统发送事件通知。与触发器相比,这些任务通常允许一定的延迟,并且不需要回滚原始操作。
一个典型的例子是订单支付后自动扣减库存。传统方案中,可以在订单表上建触发器,在插入已支付订单时同步更新库存表。在 MongoDB 中,我们可以用 Change Streams 监听 orders 集合,并对每个已支付订单执行库存扣减。下面的代码演示了这一过程。
const { MongoClient } = require('mongodb');
async function watchOrdersAndUpdateStock() {
const client = new MongoClient('mongodb://localhost:27017');
await client.connect();
const db = client.db('shop');
const orders = db.collection('orders');
const products = db.collection('products');
const pipeline = [{ $match: { 'fullDocument.status': 'paid' } }];
const changeStream = orders.watch(pipeline);
changeStream.on('change', async function(change) {
if (change.operationType === 'insert') {
const order = change.fullDocument;
for (const item of order.items) {
await products.updateOne(
{ _id: item.productId, stock: { $gte: item.quantity } },
{ $inc: { stock: -item.quantity } }
);
}
}
});
await new Promise(function(resolve) { setTimeout(resolve, 3600000); });
}
watchOrdersAndUpdateStock().catch(console.error);
代码中的 $match 管道把事件流限制为 fullDocument.status 等于 paid 的插入事件。库存扣减使用条件更新,要求当前库存大于等于购买数量,从而避免超卖。这个实现与应用代码处于同一进程,逻辑清晰且易于测试。如果扣减失败,可以通过重试机制处理,而不会影响已经成功的订单写入。传统触发器如果扣减库存失败,订单插入也会失败,这可能导致用户支付成功但订单无法记录的问题。
另一个常见场景是审计日志。我们可以监听多个业务集合,把所有变更写入 audit 集合,记录操作时间、用户、修改前后的文档快照。由于 Change Streams 异步消费,审计写入不会拖慢业务请求。这种场景下,即便审计服务短时不可用,只要配置好恢复点,也不会丢失事件。但需要注意,如果你的审计要求与业务操作强一致,Change Streams 的最终一致性会是一个短板。
Change Streams 与触发器在事务语义上的核心差异
Change Streams 不能改变原始操作的结果,也不能在写入前修改文档。数据库触发器中的 BEFORE 触发器可以在数据落盘前对字段做转换或校验,如果校验不通过可以抛出异常阻止写入。Change Streams 只提供事后通知,无法行使这样的控制权。如果你需要在插入前设置默认值或做复杂校验,应该使用 MongoDB 的 schema validation 或应用层代码,而不是指望 Change Streams。
在事务方面,MongoDB 多文档事务提交后,Change Streams 才会在 oplog 中看到相应的变更。消费者读到的是已提交的数据,但消费者自身执行的后续操作独立于原事务。如果消费者处理失败,源事务不会回滚。因此,涉及资金转账、库存扣减等需要原子性的操作,建议使用 MongoDB 事务,或者在同一个数据库操作中完成多个集合的写入,而不是通过 Change Streams 异步补齐。
另一个差异表现在顺序性。在单个复制集内,Change Streams 的事件顺序与 oplog 顺序一致,可以保证对同一文档的写操作按顺序到达。但在分片集群中,不同分片的事件可能交错,MongoDB 使用 clusterTime 提供一个全局排序依据,消费者可以依据该字段对事件排序。若业务要求全局有序,需要额外处理;传统触发器在单个数据库实例内天然严格有序。
生产环境使用 Change Streams 的恢复与扩展要点
生产环境中使用 Change Streams,首先要解决故障恢复的问题。变更流游标支持 resume token,应用程序应该定期持久化这个 token,并在重启后从上次位置继续消费,以避免事件丢失或重复。下面的代码展示了如何保存和恢复 token。
// 保存 token
const token = changeStream.resumeToken;
await saveResumeToken(token);
// 恢复时
const savedToken = await loadResumeToken();
const options = {};
if (savedToken) {
options.resumeAfter = savedToken;
}
const stream = collection.watch(pipeline, options);
值得注意的是,resume token 只在 oplog 保留窗口内有效。如果消费者停机时间超过 oplog 容量,token 可能失效,此时只能重新打开流并可能丢失一部分事件。为了减少这种风险,应该合理设置 oplog 大小,并尽快处理积压事件。也可以使用 startAtOperationTime 指定一个时间点恢复,但同样受 oplog 保留限制。
扩展性方面,多个消费者可以订阅同一个集合的变更流。每个消费者会独立读取 oplog,互不影响。但要注意,如果多个消费者执行有副作用的操作(如扣库存),需要设计幂等逻辑,避免重复处理导致数据错误。通常可以利用 change 的 _id 字段作为事件的唯一标识,在处理前先检查该事件是否已经处理过,以实现至少一次投递下的幂等消费。
权限配置上,使用 Change Streams 的用户需要具有 find 和 changeStream 权限。在副本集或分片集群上,还需要确保操作在 primary 节点执行(默认从 primary 读取 oplog)。此外,MongoDB 官方建议在连接字符串中启用 retryWrites=true,以提高网络瞬断时的可靠性。对于高性能场景,可以使用多个游标并行消费不同分片的数据,但需要自行合并事件顺序。
总的来说,MongoDB Change Streams 在异步、解耦、事件驱动方面明显优于传统触发器,但它无法提供同步执行、原子回滚和 before 语义。架构选型时,应评估一致性要求、延迟容忍度和故障处理复杂度。将 Change Streams 视为数据库的变更日志订阅机制,而不是触发器的等价物,是更务实的态度。
MongoDB Change Streams数据库触发器变更数据捕获修改时间:2026-09-18 05:00:38