变更流允许应用程序访问实时数据更改,从而避免事先手动追踪 oplog 的复杂性和风险。应用程序可使用变更流来订阅针对单个集合、数据库或整个部署的所有数据变更,并立即对它们做出响应。由于变更流采用聚合框架,因此,应用程序还可对特定变更进行过滤,或是随意转换通知。
注意
变更流仅限于数据库事件。Atlas Stream Processing具有扩展功能,包括管理多种数据事件类型,以及使用与Atlas数据库相同的查询API处理复杂数据流。有关更多信息,请参阅 Atlas Stream Processing。
从 MongoDB 5.1 开始,我们对变更流进行了优化,提高了资源利用率,并加快了某些聚合管道阶段的处理速度。
可用性
存储引擎。
副本集和分片集群必须使用 WiredTiger 存储引擎。变更流也可以用于使用静态加密功能的部署。
副本集协议版本。
副本集和分片集群必须使用副本集协议版本 1 (
pv1)。读关注“多数”启用。
无论是否支持
"majority"读关注, 变更流 都可用;也就是说,可以启用(默认)或majority禁用 读关注 支持,以使用变更流。
注意
时间序列集合不支持变更流,因为时间序列集合使用优化的存储格式,而不是跟踪文档级别的更改。您不能使用时间序列集合作为Atlas Stream Processing 的源。
提示
有关详细信息,请参阅时间序列集合限制。
稳定的 API 支持
变更流包含在Stable API V1 中。 但是, Stable APIV 中不包含 showExpandedEvents 选项。1
连接
连接 change stream 可以使用带有+srv连接选项的 DNS 种子列表,也可以在连接字符串中单独列出服务器。
如果驱动程序与变更流失去连接或连接中断,它则会尝试通过集群中具有匹配读取偏好的其他节点与变更流重新建立连接。如果驱动程序未找到具有正确读取偏好的节点,则会引发异常。
有关更多信息,请参阅连接字符串 URI 格式。
监视集合、数据库或部署
可以针对如下情况打开变更流:
目标 | 说明 |
|---|---|
集合 | 可以为单个集合(除 本页上的示例包括 MongoDB 驱动程序和 |
数据库 | 您可以为单个数据库(不包括 有关这种 MongoDB 驱动程序方法,请参阅您的驱动程序文档。另请参阅 |
部署 | 您可以为部署(副本集或分片集群)打开变更流游标,以监控所有数据库(除 有关这种 MongoDB 驱动程序方法,请参阅您的驱动程序文档。另请参阅 |
变更流性能考量
如果针对数据库打开的活动变更流的数量超过连接池大小,则可能会出现通知延迟。每个变更流在等待下一个事件时,都会通过 getMore 操作保持一个连接处于打开状态。为避免产生延迟,请确保池大小应大于已打开的变更流数量。有关详情,请参阅 maxPoolSize 设置。
分片集群注意事项
在分片集群上打开变更流时:
mongos在每个分片上创建单独的变更流。无论变更流是否针对特定的分片密钥范围,都会出现这种行为。当
mongos收到更改流结果时,它会对这些结果进行排序和筛选。如有必要,mongos还会执行fullDocument查找。
为获得最佳性能,请在变更流中限制对 $lookup 查询的使用。
打开变更流
要打开变更流:
对于副本集,请从任何承载数据的节点打开变更流。
对于分片集群,从
mongos打开变更流。
以下示例将为某一集合打开一个变更流,并对游标进行迭代以检索变更流文档。[1]
迭代变更流游标以检索变更事件。有关变更流文档格式,请参阅“变更流事件”。
变更流游标保持打开状态,直到出现以下任一情况:
游标已明确关闭。
发生失效事件;例如删除或重命名集合。
如果部署是分片集群,则分片删除可能会导致打开的变更流游标关闭。关闭的变更流游标可能无法完全恢复。
注意
未关闭游标的生命周期取决于语言。
| [1] | 您可以指定 startAtOperationTime 在特定时间点打开游标。如果指定的起点在过去,它必须在 oplog 的时间范围内。 |
修改变更流输出
查找更新操作的完整文档
默认情况下,变更流仅在更新操作期间返回字段的增量。不过,您可以配置变更流以返回已更新文档的最新多数提交版本。
updateLookup 操作从集合中读取由分片键和文档标识符标识的文档。集合由名称标识,并使用处理变更流时存在的集合数据。考虑以下场景:
如果重命名了集合,则不会返回任何文档。
如果重命名了集合,并使用旧名称创建了新集合,则会对新集合执行查找操作。如果找到匹配的文档,则将其返回。
警告
对于涉及快速删除或流量峰值的情况,使用 $match 过滤配置 fullDocument: "updateLookup" 可能会导致“未找到恢复令牌”错误。当文档删除导致 fullDocument字段返回 null 值时,会发生这种情况,因为没有匹配的文档,这会阻止变更流查找恢复令牌。
相反,请使用带有 fullDocumentBeforeChange:
"whenAvailable" 和 fullDocument: "whenAvailable" 的前像和后像。请参阅具有文档前映像和后映像的Change Streams部分。
注意
如果在更新操作之后但在查找之前有一个或多个多数提交操作修改了更新的文档,则返回的完整文档可能显著不同于更新操作时的文档。
但是,变更流文档中包含的增量始终正确地描述应用于该变更流事件的被监控集合更改。
如果以下任一条件为真,则更新事件的 fullDocument 字段可能会缺失:
如果文档被删除,或者集合在更新和查找之间被删除。
如果更新更改了该集合分片键中至少一个字段的值。
有关变更流响应文档格式的更多信息,请参阅 变更流事件。
恢复变更流
在打开游标时将恢复令牌指定为 resumeAfter 或 startAfter,借此恢复变更流。
警告
使用恢复令牌恢复变更流时,请使用与最初生成令牌时相同的管道和选项。如果使用不同的变更流管道或不同的选项,则可能会导致不可预测的行为,对数据一致性影响,或阻止变更流恢复。
resumeAfter 对于 Change Stream
您可以在打开游标时将恢复令牌传递给 resumeAfter ,从而在特定事件发生后恢复 change stream。
请参阅恢复令牌以了解有关恢复令牌的更多信息。
重要
如果时间戳位于过去,oplog 必须有足够的历史记录来定位与令牌或时间戳相关的操作。
在某一无效事件(例如,集合删除或重命名)关闭变更流后,您无法使用
resumeAfter来恢复变更流。请改为使用 startAfter 在无效事件后启动新的变更流。
startAfter 对于 Change Stream
您可在打开游标时将恢复令牌传递给 startAfter,从而在特定事件之后启动新的变更流。与 resumeAfter 不同,startAfter 可在出现无效事件之后通过创建新的变更流来恢复通知。
请参阅恢复令牌以了解有关恢复令牌的更多信息。
重要
如果时间戳位于过去,oplog 必须有足够的历史记录来定位与令牌或时间戳相关的操作。
恢复令牌
恢复令牌有两种类型:
事件令牌:标识特定的变更事件。每次发生变更事件时,变更流游标都会生成一个事件令牌。
高水位标记:表示没有关联变更事件的时间点。即使没有发生更改事件,服务器也会定期生成高水位线令牌以指示集群时间已提前。
提示
服务器会定期将高水位线恢复令牌中的时间戳向前推进。在写入不频繁的空闲分片上,对于某些使用案例来说,这种进展可能不会频繁发生。要更频繁地提前高水位线时间戳,您可以使用 appendOplogNote 命令将无操作条目写入空闲分片上的oplog 。
您可以在多个来源中找到恢复令牌:
源 | 说明 |
|---|---|
更改事件通知包含针对 | |
该字段仅在使用 | |
|
从 MongoDB 4.2 开始,如果变更流聚合管道修改了事件的 _id 字段,则变更流会引发异常。
提示
MongoDB 提供了“代码段”,这是 mongosh 的扩展,用于解码十六进制编码的恢复令牌。
您可以安装并运行 中的恢复令牌 mongosh代码段:
snippet install resumetoken decodeResumeToken('<RESUME TOKEN>')
如果系统上安装了 npm,那么您还可以在命令行中运行 resumetoken (并且不使用 mongosh):
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
使用 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
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) 服务、跨平台同步、协作功能以及通知服务时,变更流可为开发人员节省时间。
访问控制
对于在自管理部署上执行身份验证和授权的部署:
要打开针对特定集合的变更流,应用程序必须具有对相应集合授予
changeStream和find动作的特权。{ resource: { db: <dbname>, collection: <collection> }, actions: [ "find", "changeStream" ] } 要在单个数据库上打开变更流,应用程序必须具有对数据库中所有非
system集合授予changeStream和find动作的特权。{ resource: { db: <dbname>, collection: "" }, actions: [ "find", "changeStream" ] } 要在整个部署中打开变更流,应用程序必须具有对部署中所有数据库的所有非
system集合授予changeStream和find动作的特权。{ resource: { db: "", collection: "" }, actions: [ "find", "changeStream" ] }
事件通知
变更流仅在数据发生更改时通知副本集中的大多数数据承载节点。变更流确保通知仅由大多数已提交且在故障情况下持续存在的更改触发。
例如,考虑一个 3 节点副本集,针对主节点打开了变更流游标。如果客户端发出插入操作,则只有在插入持续到大多数数据承载节点后,变更流才会将数据更改通知应用程序。
如果某个操作与事务相关联,则变更事件文档包括 txnNumber 和 lsid。
排序规则
除非提供了显式排序规则,否则变更流使用 simple 二进制比较。
变更流和孤立文档
附带文档前映像和后映像的变更流
从 MongoDB 6.0 开始,可使用变更流事件来输出更改前后的文档版本(文档前映像和后映像):
前映像是指被替换、更新或删除之前的文档。已插入的文档没有前映像。
后图像是插入、替换或更新后的文档。 已删除的文档没有后图像。
使用
db.createCollection()、create或collMod为集合启用changeStreamPreAndPostImages。示例,使用collMod命令时:db.runCommand( { collMod: <collection>, changeStreamPreAndPostImages: { enabled: true } } )
如果图像属于以下情况,则前像和后像不可用于变更流事件:
在文档更新或删除操作时未对集合启用。
在
expireAfterSeconds中设置的前像和后像保留时间后之后被删除。以下示例将整个集群上的
expireAfterSeconds设置为100秒:use admin db.runCommand( { setClusterParameter: { changeStreamOptions: { preAndPostImages: { expireAfterSeconds: 100 } } } } ) 注意
MongoDB Atlas集群中不支持
setClusterParameter命令。有关AtlasAtlas支持所有命令的信息,请参阅AtlasAtlas中不支持的命令。以下示例返回当前的
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设置为required或whenAvailable。要请求后图像,您可以使用相同的方法设置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 版本。
提示
有关变更流事件和输出,请参阅变更流事件。
要查看集合的变化,请参阅
db.collection.watch()。有关变更流输出的完整示例,请参阅使用文档前像和后像的变更流。
有关变更流输出的完整示例,请参阅使用文档前像和后像的变更流。
如果集群的 initialSyncMethod 参数为 fileCopyBased,则对变更流侦听器没有影响。
如果 initialSyncMethod 为 logical,并且在新同步的节点上打开变更流并从早于逻辑初始同步完成的时间点读取事件,则前映像和后映像可能会丢失。