对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
Docs 菜单

使用变更流监控数据

在本指南中,您可以了解如何使用变更流来监控数据库的实时更改。 变更流是 MongoDB Server 的一项功能,允许应用程序订阅集合、数据库或部署上的数据更改。

提示

Atlas Stream Processing

作为变更流的替代方案,您可以使用Atlas Stream Processing来处理和转换数据流。与仅注册数据库事件的变更流不同,Atlas Stream Processing托管多种数据事件类型并提供扩展的数据处理功能。要学习;了解有关此功能的更多信息,请参阅MongoDB Atlas文档中的Atlas Stream Processing

您可以在以下对象上使用 watch() 方法监视 MongoDB 中的更改:

对于每个对象,watch() 方法都会打开一个更改流,以便在更改事件发生时发出更改事件文档。

watch() 方法可以选择采用由聚合阶段数组组成的聚合管道作为第一个参数。聚合阶段筛选和转换变更事件。

在以下代码段中,$match 阶段匹配 runtime 值小于 15 的所有更改事件文档,并过滤掉所有其他文档。

const pipeline = [ { $match: { runtime: { $lt: 15 } } } ];
const changeStream = myColl.watch(pipeline);

watch() 方法将 options 对象作为第二个参数。有关此对象配置设置的更多信息,请参阅本节末尾的链接。

watch() 方法返回 ChangeStream 的实例。您可以通过迭代变更流或监听事件,读取变更流中的事件。

警告

驱动程序不支持在 EventEmitterIterator 模式下同时使用 ChangeStream,并会导致错误。这是为了防止未定义的行为,在这种行为中,驱动程序无法保证哪个使用者先接收文档。

根据您想从变更流中读取事件的方式选择相应的标签页:

从版本 4.12 开始,ChangeStream 对象是异步可遍历对象。通过此更改,您可以使用 for-await 循环从打开的变更流中检索事件:

for await (const change of changeStream) {
console.log("Received change: ", change);
}

您可以调用ChangeStream对象上的方法,例如:

  • hasNext() 检查流中是否有剩余文档

  • next() 请求流中的下一个文档

  • close() 关闭 ChangeStream

您可以通过调用 on()方法将监听函数附加到 ChangeStream 对象。此方法继承自 Javascript EventEmitter类。如下所示,将 string "change" 作为第一个参数传递,将监听函数作为第二个参数传递:

changeStream.on("change", (changeEvent) => { /* your listener function */ });

监听器函数在发出change事件时触发。 您可以在侦听器中指定逻辑,以便在收到更改事件文档时对其进行进程。

您可以通过调用 pause() 停止触发事件或调用 resume() 继续触发事件来控制变更流。

To stop processing change events, call the close() method on the ChangeStream instance. This closes the change stream and frees resources.

changeStream.close();

注意

您可以使用此示例连接到MongoDB实例,并与包含示例数据的数据库交互。要学习;了解有关连接到MongoDB实例和加载示例数据集的更多信息,请参阅《Node.js驱动程序入门》指南。

注意

无Typescript特定功能

以下代码示例使用JavaScript。驱动程序没有与此使用案例相关的 TypeScript 特定功能。

以下示例在 insertDB 数据库中的 haikus 集合上打开一个变更流,并在发生变更事件时打印变更事件:

1// Watch for changes in a collection by using a change stream
2import { MongoClient } from "mongodb";
3
4// Replace the uri string with your MongoDB deployment's connection string.
5const uri = "<connection string uri>";
6
7const client = new MongoClient(uri);
8
9// Declare a variable to hold the change stream
10let changeStream;
11
12// Define an asynchronous function to manage the change stream
13async function run() {
14 try {
15 const database = client.db("insertDB");
16 const haikus = database.collection("haikus");
17
18 // Open a Change Stream on the "haikus" collection
19 changeStream = haikus.watch();
20
21 // Print change events as they occur
22 for await (const change of changeStream) {
23 console.log("Received change:\n", change);
24 }
25 // Close the change stream when done
26 await changeStream.close();
27
28 } finally {
29 // Close the MongoDB client connection
30 await client.close();
31 }
32}
33run().catch(console.dir);

提示

显式资源管理

The Node.js driver natively supports explicit resource management for MongoClient, ClientSession, ChangeStreams, and cursors. This feature is experimental and subject to change. To learn how to use explicit resource management, see the v6.9 Release Notes.

当您运行此代码,然后对 haikus 集合进行更改(例如执行插入或删除操作)时,您可以在终端中看到打印的变更事件文档。

例如,如果将一个文档插入到集合中,上述代码将打印以下输出:

Received change:
{
_id: {
_data: '...'
},
operationType: 'insert',
clusterTime: new Timestamp({ t: 1675800603, i: 31 }),
fullDocument: {
_id: new ObjectId("..."),
...
},
ns: { db: 'insertDB', coll: 'haikus' },
documentKey: { _id: new ObjectId("...") }
}

注意

从更新接收完整文档

默认情况下,包含更新操作信息的更改事件仅返回修改后的字段,而不是完整的更新文档。您可以将更改流配置为也返回文档的最新版本,方法是将选项对象的 fullDocument 字段设置为 "updateLookup",如下所示:

const options = { fullDocument: "updateLookup" };
// This could be any pipeline.
const pipeline = [];
const changeStream = myColl.watch(pipeline, options);

以下示例在 insertDB 数据库中的 haikus 集合上打开一个变更流。让我们创建一个监听函数来接收和打印集合上发生的变更事件。

首先,打开集合上的变更流,然后使用 on() 方法在变更流上定义侦听器。设置侦听器后,通过对集合执行更改来生成变更事件。

要在集合上生成更改事件,我们使用 insertOne() 方法添加新文档。由于 insertOne() 可能会在监听函数注册之前运行,因此我们使用定义为 simulateAsyncPause 的计时器,在插入之前等待 1 秒。

我们还在插入文档后使用 simulateAsyncPause。这为侦听器函数提供充足的时间来接收变更事件,并让侦听器在使用 close() 方法关闭 ChangeStream 实例之前完成自身的执行。

注意

包含计时器的原因

此示例中使用的计时器仅用于演示。它们可以确保有足够的时间来注册监听器并让监听器能在退出之前处理变更事件。

注意

无Typescript特定功能

以下代码示例使用JavaScript。驱动程序没有与此使用案例相关的 TypeScript 特定功能。

1/* Change stream listener */
2
3import { MongoClient } from "mongodb";
4
5// Replace the uri string with your MongoDB deployment's connection string
6const uri = "<connection string uri>";
7
8const client = new MongoClient(uri);
9
10const simulateAsyncPause = () =>
11 new Promise(resolve => {
12 setTimeout(() => resolve(), 1000);
13 });
14
15let changeStream;
16async function run() {
17 try {
18 const database = client.db("insertDB");
19 const haikus = database.collection("haikus");
20
21 // Open a Change Stream on the "haikus" collection
22 changeStream = haikus.watch();
23
24 // Set up a change stream listener when change events are emitted
25 changeStream.on("change", next => {
26 // Print any change event
27 console.log("received a change to the collection: \t", next);
28 });
29
30 // Pause before inserting a document
31 await simulateAsyncPause();
32
33 // Insert a new document into the collection
34 await myColl.insertOne({
35 title: "Record of a Shriveled Datum",
36 content: "No bytes, no problem. Just insert a document, in MongoDB",
37 });
38
39 // Pause before closing the change stream
40 await simulateAsyncPause();
41
42 // Close the change stream and print a message to the console when it is closed
43 await changeStream.close();
44 console.log("closed the change stream");
45 } finally {
46 // Close the database connection on completion or error
47 await client.close();
48 }
49}
50run().catch(console.dir);

请访问以下资源以获取本页提到的类和方法的更多信息: