MongoDB Change Streams 订阅变更

FreeGuideOnline 最新 2026-07-11

MongoDB Change Streams:实时数据变更订阅完全指南

在构建现代应用时,能够实时感知数据库变化对上至数据同步、下至触发业务逻辑都至关重要。MongoDB 的 Change Streams 正是为此而生——它允许应用程序像订阅消息队列一样,监听集合、数据库乃至整个部署的变更事件,且无需任何额外的轮询。

本教程将从零开始,带你掌握 MongoDB Change Streams 的核心概念、操作方法以及生产环境下的注意事项。

什么是 Change Streams?

Change Streams 是 MongoDB 从 3.6 版本(副本集与分片集群)开始引入的一项功能。它基于 oplog(操作日志)实现,能够将数据库中的变更以事件流的形式推送给订阅者。你完全不用关心底层 oplog 的复杂机制,只需要使用 MongoDB 驱动提供的高级 API。

简单来说,它就像数据库的“观察者模式”:

  • 有人插入了一条文档 → 你立刻收到 insert 事件。
  • 有字段被修改 → 你收到包含 updateDescriptionupdate 事件。
  • 文档被删除 → 你收到 delete 事件。

必备前提

在开始编写代码之前,请确认你的环境满足以下条件:

  1. 部署类型:Change Streams 只能工作在 副本集分片集群 上。单机 mongod 进程无法使用。
  2. 存储引擎:必须使用 WiredTiger 存储引擎,MongoDB 3.6 以上已默认启用。
  3. 权限:用于监听变更的用户需要拥有 changeStream 权限。
  4. 数据库:若监听的是集合,数据库需已存在;监听数据库或集群级别则无此限制。

基本用法(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:区分 insertupdatedeletereplacedrop 等。
  • fullDocument:仅 insertreplace 默认包含完整文档。对于 update 事件,需要通过配置 fullDocument: 'updateLookup' 来获取更新后的完整文档。
  • updateDescriptionupdate 事件特有,包含 updatedFieldsremovedFields,方便精准了解哪部分被修改。
  • 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 记录每一次字段级别的变动历史。
  • 缓存失效:当某份文档被修改时,立即通过流事件通知缓存层使其失效。

注意事项与限制

  1. oplog 窗口限制:Change Stream 依赖 oplog,如果恢复令牌对应的 oplog 记录已过期(被滚动删除),则无法恢复。必须设置足够大的 oplog 以保证容灾能力。
  2. 性能考量:每个 Change Stream 都会在服务端占用资源。避免同时打开数量过多的流。对于高吞吐写入,可考虑用一个流驱动多个消费者(扇出模式)。
  3. 事务中的变更:在分片集群上,事务内的变更会等到事务提交后才出现在 Change Stream 中,且 clusterTime 反映的是提交时间。
  4. 删除集合:当监听的集合被删除 (drop) 时,Change Stream 会关闭。若需要继续监听同名集合,程序应在收到 drop 事件后重新调用 watch()
  5. MongoDB Atlas:Atlas 提供了完全托管的 Change Streams,并支持通过 Eventbridge 等集成到无服务器架构中。

快速上手检验清单

  • 确认部署为副本集或分片集群。
  • 连接字符串中指定 replicaSet 参数。
  • 为你的应用账户授予 changeStream 和对应集合的 find 权限。
  • 在代码中实现恢复逻辑,并测试中断场景。
  • 根据业务选择合适的 fullDocument 选项和过滤管道。

总结

MongoDB Change Streams 将数据库变更从被动查询转为主动推送,极大地简化了实时数据处理系统的架构。只需几行代码,你就能构建出可靠的变更订阅管道。记住善用过滤、妥善处理容错恢复,它将成为你数据工具箱中最可靠的一环。