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

MongoDB 变更流

变更流(Change Streams)允许应用程序访问实时数据变更,从而避免事先手动追踪 oplog 的复杂性和风险。应用程序可使用变更流来订阅针对单个集合、数据库或整个部署的所有数据变更,并立即对它们做出响应。由于变更流采用聚合框架,因此,应用程序还可对特定变更进行过滤,或是随意转换通知。

注意

变更流仅限于数据库事件。Atlas Stream Processing具有扩展功能,包括管理多种数据事件类型,以及使用与Atlas数据库相同的查询API处理复杂数据流。有关更多信息,请参阅Atlas Stream Processing。

从 MongoDB 5.1 开始,我们对变更流进行了优化,提高了资源利用率,并加快了某些聚合管道阶段的处理速度。

变更流可用于副本集分片集群

  • 存储引擎。

    副本集和分片集群必须使用 WiredTiger 存储引擎。变更流也可以用于使用静态加密功能的部署。

  • 副本集协议版本。

    副本集和分片集群必须使用副本集协议版本 1 (pv1)。

  • 读关注“多数”启用。

    Change streams are available regardless of the "majority" read concern support; that is, read concern majority support can be either enabled (default) or disabled to use change streams.

注意

时间序列集合不支持变更流,因为时间序列集合使用优化的存储格式,而不是跟踪文档级别的更改。您不能使用时间序列集合作为Atlas Stream Processing 的源。

提示

有关详细信息,请参阅时间序列集合限制

变更流包含在Stable API V1 中。 但是, Stable APIV 中不包含 showExpandedEvents 选项。1

连接 change stream 可以使用带有+srv连接选项的 DNS 种子列表,也可以在连接字符串中单独列出服务器。

如果驱动程序与变更流失去连接或连接中断,它则会尝试通过集群中具有匹配读取偏好的其他节点与变更流重新建立连接。如果驱动程序未找到具有正确读取偏好的节点,则会引发异常。

有关更多信息,请参阅连接字符串 URI 格式。

可以针对如下情况打开变更流:

目标
说明

集合

可以为单个集合(除 system 集合,adminlocalconfig 数据库中的任何集合)打开变更流游标。

The examples on this page include both MongoDB driver and mongosh examples for opening and working with a change stream cursor for a single collection.

数据库

您可以为单个数据库(不包括 adminlocalconfig 数据库)打开变更流游标,以监视其所有非系统集合的更改。

For the MongoDB driver method, refer to your driver documentation. See also the mongosh method db.watch().

部署

您可以为部署(副本集或分片集群)打开变更流游标,以监控所有数据库(除 adminlocalconfig 外)中对所有非系统集合的变更。

For the MongoDB driver method, refer to your driver documentation. See also the mongosh method Mongo.watch().

如果针对数据库打开的活动变更流的数量超过连接池大小,则可能会出现通知延迟。每个变更流在等待下一个事件时,都会通过 getMore 操作保持一个连接处于打开状态。为避免产生延迟,请确保池大小应大于已打开的变更流数量。有关详情,请参阅 maxPoolSize 设置。

在分片集群上打开变更流时:

  • mongos每个分片上创建单独的变更流。无论变更流是否针对特定的分片密钥范围,都会出现这种行为。

  • mongos 收到更改流结果时,它会对这些结果进行排序和筛选。如有必要,mongos 还会执行 fullDocument 查找。

为获得最佳性能,请在变更流中限制对 $lookup 查询的使用。

要打开变更流:

  • 对于副本集,请从任何承载数据的节点打开变更流。

  • 对于分片集群,从 mongos 打开变更流。

以下示例将为某一集合打开一个变更流,并对游标进行迭代以检索变更流文档。[1]

迭代变更流游标以检索变更事件。有关变更流文档格式,请参阅“变更流事件”。

变更流游标保持打开状态,直到出现以下任一情况:

  • 游标已明确关闭。

  • 发生失效事件;例如删除或重命名集合。

  • 与MongoDB 部署的连接关闭或超时。有关更多信息,请参阅行为。

  • 如果部署是分片集群,则分片删除可能会导致打开的变更流游标关闭。关闭的变更流游标可能无法完全恢复。

注意

未关闭游标的生命周期取决于语言。

[1] 您可以指定 startAtOperationTime 在特定时间点打开游标。如果指定的起点在过去,它必须在 oplog 的时间范围内。

提示

The _id field of the change stream event document act as the resume token. Do not use the pipeline to modify or remove the change stream event's _id field.

从 MongoDB 4.2 开始,如果变更流聚合管道修改了事件的 _id 字段,则变更流会引发异常。

有关变更流响应文档格式的更多信息,请参阅 变更流事件

默认情况下,变更流仅在更新操作期间返回字段的增量。不过,您可以配置变更流以返回已更新文档的最新多数提交版本。

updateLookup 操作从集合中读取由分片键和文档标识符标识的文档。集合由名称标识,并使用处理变更流时存在的集合数据。考虑以下场景:

  • 如果重命名了集合,则不会返回任何文档。

  • 如果重命名了集合,并使用旧名称创建了新集合,则会对新集合执行查找操作。如果找到匹配的文档,则将其返回。

警告

对于涉及快速删除或流量峰值的情况,使用 $match 过滤配置 fullDocument: "updateLookup" 可能会导致“未找到恢复令牌”错误。当文档删除导致 fullDocument字段返回 null 值时,会发生这种情况,因为没有匹配的文档,这会阻止变更流查找恢复令牌。

相反,请使用带有 fullDocumentBeforeChange: "whenAvailable"fullDocument: "whenAvailable" 的前像和后像。请参阅具有文档前映像和后映像的Change Streams部分。

注意

如果在更新操作之后但在查找之前有一个或多个多数提交操作修改了更新的文档,则返回的完整文档可能显著不同于更新操作时的文档。

但是,变更流文档中包含的增量始终正确地描述应用于该变更流事件的被监控集合更改。

如果以下任一条件为真,则更新事件的 fullDocument 字段可能会缺失:

  • 如果文档被删除,或者集合在更新和查找之间被删除。

  • 如果更新更改了该集合分片键中至少一个字段的值。

有关变更流响应文档格式的更多信息,请参阅 变更流事件

在打开游标时将恢复令牌指定为 resumeAfter startAfter,借此恢复变更流。

警告

使用恢复令牌恢复变更流时,请使用与最初生成令牌时相同的管道和选项。如果使用不同的变更流管道或不同的选项,则可能会导致不可预测的行为,对数据一致性影响,或阻止变更流恢复。

您可以在打开游标时将恢复令牌传递给 resumeAfter ,从而在特定事件发生后恢复 change stream。

请参阅恢复令牌以了解有关恢复令牌的更多信息。

重要

  • 如果时间戳位于过去,oplog 必须有足够的历史记录来定位与令牌或时间戳相关的操作。

  • 在某一无效事件(例如,集合删除或重命名)关闭变更流后,您无法使用 resumeAfter 来恢复变更流。请改为使用 startAfter无效事件后启动新的变更流。

You can start a new change stream after a specific event by passing a resume token to startAfter when opening the cursor. Unlike resumeAfter, startAfter can resume notifications after an invalidate event by creating a new change stream.

请参阅恢复令牌以了解有关恢复令牌的更多信息。

重要

  • 如果时间戳位于过去,oplog 必须有足够的历史记录来定位与令牌或时间戳相关的操作。

恢复令牌有两种类型:

  • 事件令牌:标识特定的变更事件。每次发生变更事件时,变更流游标都会生成一个事件令牌。

  • 高水位标记:表示没有关联变更事件的时间点。即使没有发生更改事件,服务器也会定期生成高水位线令牌以指示集群时间已提前。

提示

服务器会定期将高水位线恢复令牌中的时间戳向前推进。在写入不频繁的空闲分片上,对于某些使用案例来说,这种进展可能不会频繁发生。要更频繁地提前高水位线时间戳,您可以使用 appendOplogNote 命令将无操作条目写入空闲分片上的oplog 。

您可以在多个来源中找到恢复令牌:

说明

更改事件通知包含针对 _id 字段的恢复词元:

$changeStream 聚合阶段在 cursor.postBatchResumeToken 字段中包含恢复令牌。

该字段仅在使用 aggregate 命令时显示。

getMore 命令在 cursor.postBatchResumeToken 字段中包含一个恢复令牌。

从 MongoDB 4.2 开始,如果变更流聚合管道修改了事件的 _id 字段,则变更流会引发异常。

提示

MongoDB provides a "snippet", an extension to mongosh, that decodes hex-encoded resume tokens.

You can install and run the resumetoken snippet from mongosh:

snippet install resumetoken
decodeResumeToken('<RESUME TOKEN>')

You can also run resumetoken from the command line (without using mongosh) if npm is installed on your system:

npx mongodb-resumetoken-decoder <RESUME TOKEN>

请参阅以下内容了解详细信息:

更改事件通知包含针对 _id 字段的恢复令牌:

{
"_id": {
"_data": "82635019A0000000012B042C0100296E5A1004AB1154ACACD849A48C61756D70D3B21F463C6F7065726174696F6E54797065003C696E736572740046646F63756D656E744B65790046645F69640064635019A078BE67426D7CF4D2000004"
},
"operationType": "insert",
"clusterTime": Timestamp({ "t": 1666193824, "i": 1 }),
"collectionUUID": new UUID("ab1154ac-acd8-49a4-8c61-756d70d3b21f"),
"wallTime": ISODate("2022-10-19T15:37:04.604Z"),
"fullDocument": {
"_id": ObjectId("635019a078be67426d7cf4d2"'),
"name": "Giovanni Verga"
},
"ns": {
"db": "test",
"coll": "names"
},
"documentKey": {
"_id": ObjectId("635019a078be67426d7cf4d2")
}
}

使用 aggregate 命令时,$changeStream 聚合阶段在 cursor.postBatchResumeToken 字段中包含恢复令牌:

{
"cursor": {
"firstBatch": [],
"postBatchResumeToken": {
"_data": "8263515EAC000000022B0429296E1404"
},
"id": Long("4309380460777152828"),
"ns": "test.names"
},
"ok": 1,
"$clusterTime": {
"clusterTime": Timestamp({ "t": 1666277036, "i": 1 }),
"signature": {
"hash": Binary(Buffer.from("0000000000000000000000000000000000000000", "hex"), 0),
"keyId": Long("0")
}
},
"operationTime": Timestamp({ "t": 1666277036, "i": 1 })
}

getMore 命令还在 cursor.postBatchResumeToken 字段中包含一个恢复令牌:

{
"cursor": {
"nextBatch": [],
"postBatchResumeToken": {
"_data": "8263515979000000022B0429296E1404"
},
"id": Long("7049907285270685005"),
"ns": "test.names"
},
"ok": 1,
"$clusterTime": {
"clusterTime": Timestamp( { "t": 1666275705, "i": 1 } ),
"signature": {
"hash": Binary(Buffer.from("0000000000000000000000000000000000000000", "hex"), 0),
"keyId": Long("0")
}
},
"operationTime": Timestamp({ "t": 1666275705, "i": 1 })
}

变更流对于采用业务依赖型系统的基础设施很有益处,因为数据更改一旦变为持久更改,它就会通知下游系统。例如,在实施提取、转换和加载 (ETL) 服务、跨平台同步、协作功能以及通知服务时,变更流可为开发人员节省时间。

对于在自管理部署上执行身份验证授权的部署:

  • 要打开针对特定集合的变更流,应用程序必须具有对相应集合授予 changeStreamfind 动作的特权。

    { resource: { db: <dbname>, collection: <collection> }, actions: [ "find", "changeStream" ] }
  • 要在单个数据库上打开变更流,应用程序必须具有对数据库中所有非 system 集合授予 changeStreamfind 动作的特权。

    { resource: { db: <dbname>, collection: "" }, actions: [ "find", "changeStream" ] }
  • 要在整个部署中打开变更流,应用程序必须具有对部署中所有数据库的所有非 system 集合授予 changeStreamfind 动作的特权。

    { resource: { db: "", collection: "" }, actions: [ "find", "changeStream" ] }

变更流仅在数据发生更改时通知副本集中的大多数数据承载节点。变更流确保通知仅由大多数已提交且在故障情况下持续存在的更改触发。

例如,考虑一个 3 节点副本集,针对主节点打开了变更流游标。如果客户端发出插入操作,则只有在插入持续到大多数数据承载节点后,变更流才会将数据更改通知应用程序。

如果某个操作与事务相关联,则变更事件文档包括 txnNumberlsid

除非提供了显式排序规则,否则变更流使用 simple 二进制比较。

Starting in MongoDB 5.3, during range migration, change stream events are not generated for updates to orphaned documents.

从 MongoDB 6.0 开始,可使用变更流事件来输出更改前后的文档版本(文档前映像和后映像):

  • 前映像是指被替换、更新或删除之前的文档。已插入的文档没有前映像。

  • 后图像是插入、替换或更新后的文档。 已删除的文档没有后图像。

  • 使用 db.createCollection()createcollMod 为集合启用 changeStreamPreAndPostImages。示例,使用 collMod 命令时:

    db.runCommand( {
    collMod: <collection>,
    changeStreamPreAndPostImages: { enabled: true }
    } )

To check the current changeStreamPreAndPostImages setting for a collection, run db.getCollectionInfos():

db.getCollectionInfos( {
name: "<collection>"
} )[0].options.changeStreamPreAndPostImages
{ changeStreamPreAndPostImages: { enabled: true } }

如果图像属于以下情况,则前像和后像不可用于变更流事件

  • 在文档更新或删除操作时未对集合启用。

  • expireAfterSeconds 中设置的前像和后像保留时间后之后被删除。

    • 以下示例将整个集群上的 expireAfterSeconds 设置为 100 秒:

      use admin
      db.runCommand( {
      setClusterParameter:
      { changeStreamOptions: {
      preAndPostImages: { expireAfterSeconds: 100 }
      } }
      } )

      注意

      The setClusterParameter command is not supported in MongoDB Atlas clusters. For information on Atlas support for all commands, see Unsupported Commands in Atlas Clusters.

    • 以下示例返回当前的 changeStreamOptions 设置,包括 expireAfterSeconds

      db.adminCommand( { getClusterParameter: "changeStreamOptions" } )
    • expireAfterSeconds 设置为 off 可使用默认保留策略:将保留前像和后像,直到从 oplog 中删除对应的变更流事件。

    • 如果变更流事件从 oplog 中删除,则无论 expireAfterSeconds 前映像和后映像保留时间如何,相应的前映像和后映像也会被删除。

其他考量:

  • 启用前像和后像会占用存储空间并增加处理时间。仅在需要时启用前像和后像。

  • 将变更流事件大小限制为小于 16 MiB。要限制事件大小,您可以:

    • 将文档大小限制为 8 MB。如果其他 change stream 事件字段(例如 updateDescription)不是很大,则可以在 change stream 输出中同时请求更新前的文档和更新后的文档。

    • 如果其他变更流事件字段(例如 updateDescription)并不大,则仅请求变更流输出中最多 16 MiB 的文档的后像。

    • 在以下情况下,仅请求变更流输出中最多 16 MiB 的文档的前像:

      • 文档更新仅影响文档结构或内容的一小部分,

      • 不会引起 replace 变更事件。replace 事件始终包含后像。

  • 要请求前图像,请在db.collection.watch()中将fullDocumentBeforeChange设置为requiredwhenAvailable。要请求后图像,您可以使用相同的方法设置fullDocument

  • 前像被写入 config.system.preimages 集合。

    • config.system.preimages 集合可能会变大。要限制集合大小,可如前文所示为前映像设置 expireAfterSeconds 时间。

    • 要监控 config.system.preimages 的大小,请连接到分片集群上的分片节点或副本集上的 mongod 节点。然后,运行以下命令:

      use config
      db.system.preimages.totalSize()
      db.system.preimages.stats()

      注意

      您必须对 config.system.preimages 集合拥有 collStats 权限操作才能运行这些命令。

      要在 MongoDB Atlas 部署上运行这些命令,您必须具有 atlasAdmin 角色。

    • 前像由后台进程异步删除。

重要

向后不兼容的功能

从 MongoDB 6.0 开始,如果您将文档前图像和后图像用于 change stream,则必须使用 collMod 命令为每个集合禁用 changeStreamPreAndPostImages,然后才能降级到早期 MongoDB 版本。

提示

有关变更流输出的完整示例,请参阅使用文档前像和后像的变更流

如果集群的 initialSyncMethod 参数为 fileCopyBased,则对变更流侦听器没有影响。

如果 initialSyncMethodlogical,并且在新同步的节点上打开变更流并从早于逻辑初始同步完成的时间点读取事件,则前映像和后映像可能会丢失。