定义
$source 阶段在连接注册表中指定要从中流式传输数据的连接。支持以下连接类型:
Apache Kafka 代理
MongoDB collection change stream
MongoDB database change stream
MongoDB 集群变更流
AWS Kinesis数据流
文档数组
Cron安排
语法
Apache Kafka 代理
为了操作来自Apache Kafka代理的流媒体数据, $source阶段具有以下原型形式:
{ "$source": { "connectionName": "<registered-connection>", "topic" : ["<source-topic>", ...], "timeField": { $toDate | $dateFromString: <expression> }, "partitionIdleTimeout": { "size": <duration-number>, "unit": "<duration-unit>" }, "schemaRegistry": { "connectionName": "<schema-registry-name>", }, "config": { "auto_offset_reset": "<start-event>", "group_id": "<group-id>", "keyFormat": "<deserialization-type>", "keyFormatError": "<error-handling>" }, } }
$source 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 | |
|---|---|---|---|---|
| 字符串 | 必需 | 用于标识连接注册表中要从中提取数据的连接的标签。 | |
| 字符串或字符串数组 | 必需 | 从中流式传输消息的一个或多个Apache Kafka主题的名称。如果要流多个主题的消息,请在数组中指定这些主题。 | |
| 文档 | Optional | 为传入消息定义权威时间戳的文档。 如果使用
如果您未声明 | |
| 文档 | Optional | 指定分区在水印计算中被忽略之前允许其空闲的时间长度的文档。 默认情况下,此字段为禁用。要处理因空闲而无进展的分区,请为此字段设置一个值。 | |
| 整型 | Optional | 指定分区空闲超时持续时间的数字。 | |
| 字符串 | Optional | 分区空闲超时持续时间的时间单位。
| |
| 文档 | Optional | ||
| 字符串 | 可选的 | 用于 Avro 反序列化的模式注册表连接的名称。 | |
| 文档 | Optional | 包含可覆盖各种默认值的字段的文档。 | |
| 字符串 | Optional | 指定从Apache Kafka源主题中的哪个事件开始摄取。
默认值为 | |
| 字符串 | Optional | Kafka 消费者组的 ID 用于与流处理器关联。如果省略,Atlas Stream Processing 会将 Stream Processing 工作区关联为以下格式的自动生成 ID: Atlas Stream Processing会自动为所有持久流处理器生成此参数的值。对于使用 sp 定义的临时流处理器。 进程(),仅当您手动定义该参数时才会设立该参数。 | |
| 布尔 | 可选的 | 用于确定Kafka代理分区偏移提交策略的标志。Atlas Stream Processing支持两种提交策略:
对于使用 sp 有关使用Kafka作为 | |
| 字符串 | Optional | ||
| 字符串 | Optional | 如何处理在反序列化 Apache Kafka 密钥数据时遇到的错误。必须是以下值之一:
|
注意
Atlas Stream Processing 要求源数据流中的文档在 json 或 ejson有效。Atlas Stream Processing 会将不满足此要求的文档设置为死信队列 (如果您已配置)。
MongoDB 集合变更流
Atlas 集合变更流允许应用程序访问针对单个集合的实时数据更改。要学习如何对集合打开变更流,请参阅变更流。
使用变更流 $source 时,请为源集群配置至少 24 小时的 oplog 窗口。
为读取变更流,Atlas Stream Processing 会扫描 oplog 集合。因此,您可能会在日志中看到 COLLSCAN 警告。这些警告表示正常行为,并不表示错误。
如果您将 config.fullDocument 或 config.fullDocumentBeforeChange 配置为 required,请在要捕获的任何写入操作发生前,为每个集合启用 changeStreamPreAndPostImages。如果由于您在写入发生时尚未启用该功能,或后像已过期,导致某个事件的后像不可用,流处理器将失败。如需了解如何启用前像和后像,请参阅带文档前像和后像的变更流。
重要
仅当oplog集合仍包含其恢复令牌标识的操作时,您才能恢复变更流。为确保您的应用程序可以恢复其变更流,请将最小 Oplog 窗口配置为比您预期的最长中断更长的时间。
为了操作来自 Atlas collection 的流媒体数据, $source阶段具有以下原型形式:
{ "$source": { "connectionName": "<registered-connection>", "timeField": { $toDate | $dateFromString: <expression> }, "db" : "<source-db>", "coll" : ["<source-coll>",...], "initialSync": { "enable": <boolean>, "parallelism": <integer> }, "readPreference": "<read-preference>", "readPreferenceTags": [ {"<key>": "<value>"}, . . . ], "config": { "startAfter": <start-token> | "startAtOperationTime": <timestamp>, "fullDocument": "<full-doc-condition>", "fullDocumentOnly": <boolean>, "fullDocumentBeforeChange": "<before-change-condition>", "pipeline": [{ "<aggregation-stage>" : { <stage-input>, . . . }, . . . }], "maxAwaitTimeMS": <time-ms>, } } }
$source 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 字符串 | 可选的 | 用于标识连接注册表中要从中提取数据的连接的标签。 |
| 文档 | Optional | 为传入消息定义权威时间戳的文档。 如果使用
如果您未声明 |
| 字符串 | 必需 | 在 |
| 字符串或字符串数组 | 必需 | |
| 文档 | Optional | 包含 Atlas Stream Processing 如果 Atlas Stream Processing修复了流处理器启动时要同步的集合设立。它不会同步您在该点之后创建的集合,如果同步任何一个集合失败,整个流处理器都会失败。 如果启用 |
| 布尔 | 可选的 | 确定是否启用 |
| 整型 | Optional | 确定进程 如果
可设置的最大值取决于 Stream Processing 器的层级。要了解更多信息,请参阅 Atlas Stream Processing 层级选择指南。 每个流处理器都有一个由其层级确定的最大累积并行度值。流处理器的累积并行度计算如下:
其中, 示例,如果您的 如果流处理器超过其层级的最大累积并行度, Atlas Stream Processing会引发错误,并建议您实现预期并行度层级所需的最低处理器层级。您必须将处理器扩展到更高层级,或降低阶段的并行度值,才能解决该错误。要学习,了解更多信息,请参阅Stream Processing。 |
| 字符串 | Optional | 读取偏好(read preference)变更流和 默认值为 |
| 阵列 | Optional | 读取读取偏好(read preference)标签用于变更流和 |
| 文档 | Optional | 包含可覆盖各种默认值的字段的文档。 |
| token | 可选的 | 源开始报告的变更事件。这采用 resume token 的形式。 您只能使用 |
| 时间戳日期 | 可选的 | 源应开始报告的操作时间。 您只能使用 接受 MongoDB 扩展 JSON |
| 字符串 | 可选的 | |
| 布尔 | 可选的 | |
| 字符串 | Optional | 指定change stream源是否应在输出中包含处于原始“更改之前”状态的文档。必须是以下之一:
如果没有为 要将此字段与集合变更流一起使用,您必须在该集合上启用变更流前像和后像。 |
| 文档 | Optional | 指定一个聚合管道,用于传递变更流输出进行进一步处理之前对其进行过滤。该管道必须符合修改变更流输出中描述的参数。 重要提示:每个变更事件都包括 |
| 整型 | Optional | 在返回空批次之前,等待新数据更改以报告给变更流游标的最长时间(以毫秒为单位)。 默认值为 |
MongoDB database change stream
Atlas 数据库变更流允许应用程序访问针对单个数据库的实时数据更改。要学习如何针对数据库打开变更流,请参阅变更流。
使用变更流 $source 时,请为源集群配置至少 24 小时的 oplog 窗口。
为读取变更流,Atlas Stream Processing 会扫描 oplog 集合。因此,您可能会在日志中看到 COLLSCAN 警告。这些警告表示正常行为,并不表示错误。
如果您将 config.fullDocument 或 config.fullDocumentBeforeChange 配置为 required,请在要捕获的任何写入操作发生前,为每个集合启用 changeStreamPreAndPostImages。如果由于您在写入发生时尚未启用该功能,或后像已过期,导致某个事件的后像不可用,流处理器将失败。如需了解如何启用前像和后像,请参阅带文档前像和后像的变更流。
重要
仅当oplog集合仍包含其恢复令牌标识的操作时,您才能恢复变更流。为确保您的应用程序可以恢复其变更流,请将最小 Oplog 窗口配置为比您预期的最长中断更长的时间。
为了操作来自 Atlas 数据库变更流的流数据, $source阶段具有以下原型形式:
{ "$source": { "connectionName": "<registered-connection>", "timeField": { $toDate | $dateFromString: <expression> }, "db" : "<source-db>", "readPreference": "<read-preference>", "readPreferenceTags": [ {"<key>": "<value>"}, . . . ], "config": { "startAfter": <start-token> | "startAtOperationTime": <timestamp>, "fullDocument": "<full-doc-condition>", "fullDocumentOnly": <boolean>, "fullDocumentBeforeChange": "<before-change-condition>", "pipeline": [{ "<aggregation-stage>" : { <stage-input>, . . . }, . . . }] }, } }
$source 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 字符串 | 可选的 | 用于标识连接注册表中要从中提取数据的连接的标签。 |
| 文档 | Optional | 为传入消息定义权威时间戳的文档。 如果使用
如果您未声明 |
| 字符串 | 必需 | 在 |
| 字符串 | Optional | 默认值为 |
| 阵列 | Optional | 读取偏好(read preference)标签 for 变更流操作。 |
| 文档 | Optional | 包含可覆盖各种默认值的字段的文档。 |
| token | 可选的 | 源开始报告的变更事件。这采用 resume token 的形式。 您只能使用 |
| 时间戳日期 | 可选的 | 源应开始报告的操作时间。 您只能使用 接受 MongoDB 扩展 JSON |
| 字符串 | 可选的 | 用于控制变更流源是应返回完整文档还是仅在发生更新时返回变更的设置。 必须是以下之一:
如果没有为 要将 |
| 布尔 | 可选的 | |
| 字符串 | Optional | 指定change stream源是否应在输出中包含处于原始“更改之前”状态的文档。必须是以下之一:
如果没有为 要将此字段与数据库变更流一起使用,您必须对该数据库中的每个集合启用变更流前像和后像。 |
| 文档 | Optional | 指定一个聚合管道,用于过滤源点的变更流输出。该管道必须符合修改变更流输出中描述的参数。 重要提示:每个变更事件都包括 |
| 整型 | Optional | 在返回空批次之前,等待新数据更改以报告给变更流游标的最长时间(以毫秒为单位)。 默认值为 |
MongoDB 集群范围变更流源
使用变更流 $source 时,请为源集群配置至少 24 小时的 oplog 窗口。
为读取变更流,Atlas Stream Processing 会扫描 oplog 集合。因此,您可能会在日志中看到 COLLSCAN 警告。这些警告表示正常行为,并不表示错误。
如果您将 config.fullDocument 或 config.fullDocumentBeforeChange 配置为 required,请在要捕获的任何写入操作发生前,为每个集合启用 changeStreamPreAndPostImages。如果由于您在写入发生时尚未启用该功能,或后像已过期,导致某个事件的后像不可用,流处理器将失败。如需了解如何启用前像和后像,请参阅带文档前像和后像的变更流。
重要
仅当oplog集合仍包含其恢复令牌标识的操作时,您才能恢复变更流。为确保您的应用程序可以恢复其变更流,请将最小 Oplog 窗口配置为比您预期的最长中断更长的时间。
要对来自整个 Atlas 集群变更流的流数据进行操作, $source 阶段具有以下原型形式:
{ "$source": { "connectionName": "<registered-connection>", "timeField": { $toDate | $dateFromString: <expression> }, "readPreference": "<read-preference>", "readPreferenceTags": [ {"<key>": "<value>"}, . . . ], "config": { "startAfter": <start-token> | "startAtOperationTime": <timestamp>, "fullDocument": "<full-doc-condition>", "fullDocumentOnly": <boolean>, "fullDocumentBeforeChange": "<before-change-condition>", "pipeline": [{ "<aggregation-stage>" : { <stage-input>, . . . }, . . . }] }, } }
$source 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 字符串 | 可选的 | 用于标识连接注册表中要从中提取数据的连接的标签。 |
| 文档 | Optional | 为传入消息定义权威时间戳的文档。 如果使用
如果您未声明 |
| 字符串 | Optional | 默认值为 |
| 阵列 | Optional | 读取偏好(read preference)标签 for 变更流操作。 |
| 文档 | Optional | 包含可覆盖各种默认值的字段的文档。 |
| token | 可选的 | 源开始报告的变更事件。这采用 resume token 的形式。 您只能使用 |
| 日期 | 时间戳 | 可选的 | 源应开始报告的操作时间。 您只能使用 接受 MongoDB 扩展 JSON |
| 字符串 | 可选的 | 用于控制变更流源是应返回完整文档还是仅在发生更新时返回变更的设置。 必须是以下之一:
如果没有为 要将 |
| 布尔 | 可选的 | |
| 字符串 | Optional | 指定change stream源是否应在输出中包含处于原始“更改之前”状态的文档。必须是以下之一:
如果没有为 要将此字段与数据库变更流一起使用,您必须对该数据库中的每个集合启用变更流前像和后像。 |
| 文档 | Optional | 指定一个聚合管道,用于过滤源点的变更流输出。该管道必须符合修改变更流输出中描述的参数。 请注意, Atlas Stream Processing希望从每个摄取的变更事件中接收 |
| 整型 | Optional | 在返回空批次之前,等待新数据更改以报告给变更流游标的最长时间(以毫秒为单位)。 默认值为 |
AWS Kinesis Data Stream
Atlas Stream Processing支持创建与 AWS Kinesis Streams 的 Private Link 连接。要学习;了解更多信息,请参阅添加Kinesis Private Link 连接。
为了对来自 AWS Kinesis数据流的数据进行操作,$source 阶段具有以下原型形式:
{ "$source": { "connectionName": "<registered-connection>", "stream": "<stream-name>", "region": "<aws-region>", "timeField": { $toDate | $dateFromString: <expression> }, "tsFieldName": "<field-name>", "shardIdleTimeout": { "size": <duration-number>, "unit": "<duration-unit>" }, "config": { "consumerARN": "<aws-arn>", "initialPosition": <initial-position>, reshardDetectionIntervalSecs: <interval> } } }
$source 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 字符串 | 必需 | 用于标识连接注册表中连接的标签,从中提取数据。 |
| 字符串 | 必需 | ARN 与Kinesis使用方相对应的。您的使用者必须使用增强型扇出。 |
| 字符串 | 必需 | 要从其中传输消息的AWS Kinesis数据流。 |
| 字符串 | 可选的 | 指定流所在的 AWS 区域。Kinesis支持在不同区域具有相同名称的多个数据流。如果您在同一连接内的两个或多个区域中对数据流使用相同的名称,则必须使用此字段来指定要使用的名称和区域的组合。 |
| 文档 | Optional | 为传入消息定义权威时间戳的文档。 如果使用
如果您未声明 |
| 字符串 | Optional | 投影文档中用于时间戳的字段名。使用此字段以覆盖默认时间戳字段名。 |
| 文档 | Optional | 指定分片在水印计算中被忽略之前允许分片的时间长度的文档。 默认下,此字段处于禁用状态。要处理由于空闲而不向前移动的分片,请为此字段设立一个值。 |
| 文档 | Optional | 指定分片空闲超时持续时间的数字。 |
| 文档 | Optional | 分片空闲超时持续时间的时间单位。
|
| 文档 | Optional | 包含可覆盖各种默认值的字段的文档。 |
| 字符串 | Optional | 开始摄取消息的 Kinesis 数据流历史中的位置。必须是以下之一:
默认值为 |
| 日期 | 可选的 | 开始摄取消息的时间戳。当 |
| 整型 | Optional | 为了重新分片而检查流经Kinesis流的数据速率的时间间隔(以秒为单位)。 默认为 |
文档数组
为了对文档数组进行操作, $source阶段具有以下原型形式:
{ "$source": { "timeField": { $toDate | $dateFromString: <expression> }, "documents" : [{source-doc},...] | <expression> } }
$source 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 文档 | Optional | 为传入消息定义权威时间戳的文档。 如果使用
如果您未声明 |
| 阵列 | 可选的 | 用作流媒体数据源的文档数组。该字段的值可以是对象数组,也可以是计算结果为对象数组的表达式。使用 |
Cron 计划
要安排生成文档而不是从连接读取,$source 阶段具有以下原型形式:
{ "$source": { "schedule": "<cron-expression>", "tsFieldName": "<timestamp-field-name>" } }
$source 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 字符串 | 必需 | |
| 字符串 | Optional | 该阶段将计划时间戳投影到的字段的名称。默认为 |
行为
$source必须是它所在的任何管道的第一阶段。 每个管道只能使用一个$source阶段。
对于Kafka $source 阶段, Atlas Stream Processing会并行读取源主题内的多个分区。分区限制由处理器层级决定。如需学习;了解更多信息,请参阅流处理计费参考文档。
对于 cron $source 阶段,每个计划时间都会生成一个带有计划时间戳的空文档。使用管道中的后续阶段来填充文档。
当 schedule表达式同时限制一个月中的某一天和一周中的某一天时, Atlas Stream Processing仅在满足这两个字段的计划时间发出文档。这与 cron 实现不同,后者在任一字段匹配时发出文档。
如果具有 cron $source 阶段的流处理器重新启动,则在未运行时,它会为错过的每个计划时间发出一个文档。在您停止发送后,它不会在错过的计划时间内发出文档。
具有 cron $source 阶段的流处理器在计划时间之间连续运行, Atlas在整个时间段内在处理器层级对其进行计费。不频繁的安排不会降低处理器运行的费用。如需学习;了解更多信息,请参阅流处理计费参考文档。
示例
Kafka 示例
流数据源从不同位置生成详细的天气报告,符合 示例天气数据集的模式。以下聚合分为三个阶段:
$source阶段与Apache Kafka代理建立连接,在名为my_weatherdata的主题中收集这些报告,从而在将每条记录摄取到后续聚合阶段时将其公开。此阶段还会覆盖其投影的时间戳字段的名称,将其设置为ingestionTime。$match阶段会排除dewPoint.value小于或等于5.0的文档,并将dewPoint.value大于5.0的文档传递到下一个阶段。$merge阶段将输出写入sample_weatherstream数据库中名为stream的 Atlas 集合。如果不存在此类数据库或集合,Atlas 会创建它们。
[{ "$source": { "connectionName": "sample_weatherdata", "topic": "my_weatherdata" } }, { "$match": { "dewPoint.value": { "$gt": 5 } } }, { "$merge": { "into": { "connectionName": "weatherStreamOutput", "db": "sample_weatherstream", "coll": "stream" } } }]
要查看生成的 sample_weatherstream.stream 集合中的文档,请连接到您的 Atlas 集群并运行以下命令:
db.getSiblingDB("sample_weatherstream").stream.find()
注意
以上是一个有代表性的示例。流数据不是静态的,每个用户看到的都是不同的文档。
变更流示例
以下聚合会从 cluster0-collection 源引入数据,而该源会连接到已加载示例数据集的 Atlas 集群。要了解如何创建 Stream Processing 工作区并将 Atlas 集群连接添加到连接注册表,请参阅Atlas Stream Processing 入门。此聚合会运行两个阶段以打开变更流并记录对 sample_weatherdata 数据库中 data 集合的更改:
$source阶段会连接到cluster0-collection源,并针对sample_weatherdata数据库中的data集合打开变更流。$merge阶段会将过滤后的变更流文档写入sample_weatherdata数据库中名为data_changes的 Atlas 集合。如果不存在此类集合,Atlas 则会进行创建。
[{ "$source": { "connectionName": "cluster0-connection", "db": "sample_weatherdata", "coll": "data" } }, { "$merge": { "into": { "connectionName": "cluster0-connection", "db": "sample_weatherdata", "coll": "data_changes" } } }]
以下 mongosh 命令会删除 data 文档:
db.getSiblingDB("sample_weatherdata").data.deleteOne( { _id: ObjectId("5553a99ae4b02cf715120e4b") } )
删除 data 文档后,流处理器会将此变更流事件文档写入 sample_weatherdata.data_changes 集合。要查看生成的 sample_weatherdata.data_changes 集合中的文档,请使用 mongosh 连接到您的 Atlas 集群并运行以下命令:
db.getSiblingDB("sample_weatherdata").data_changes.find()
数据库变更流示例
以下聚合从 db-change-stream-connection 源摄取数据,该源连接到加载了 Sample Mflix Dataset Collections 数据集的Atlas 集群。要学习;了解如何创建流处理工作区以及如何将与Atlas 集群的连接添加到连接注册表中,请参阅Atlas Stream Processing入门。此聚合运行两个阶段,以打开针对 sample_mflix 源数据库的变更流,并在 sample_mflix_changes 接收器数据库中记录对 db_changes集合的更改:
$source阶段连接到db-change-stream-connection源并打开针对sample_mflix源数据库的变更流。config.startAtOperationTime字段设置源开始报告的时间。此示例将startAt值设置为从过去一分钟开始。$merge阶段将变更流文档写入sample_mflix_changes接收器数据库中名为db_changes的Atlas集合。
const startAt = new Date(Date.now() - 60 * 1000); const pipeline = [ { $source: { connectionName: "db-change-stream-connection", db: "sample_mflix", config: { startAtOperationTime: startAt } } }, { $merge: { into: { connectionName: "db-change-stream-connection", db: "sample_mflix_changes", coll: "db_changes" } } } ];
对源集群运行以下命令,验证流处理器的行为。要进行连接,请参阅通过mongosh连接到集群。
在 sample_mflix 源数据库的movies集合中插入一个文档,在 comments集合中插入一个文档:
db.getSiblingDB("sample_mflix").movies.insertOne({ title: "The Stream Processor", year: 2026 }) db.getSiblingDB("sample_mflix").comments.insertOne({ name: "Ada Lovelace", text: "A fine film about data in motion." })
插入文档后,流处理器会将每次插入的变更流事件文档写入 sample_mflix_changes.db_changes集合。对集群运行以下命令,查看生成的 sample_mflix_changes.db_changes集合中的文档:
db.getSiblingDB("sample_mflix_changes").db_changes.find( {}, { _id: 0, clusterTime: 1, ns: 1, operationType: 1, fullDocument: 1 } )
每个事件都包括时间戳、完整文档更改、operationType 和 ns字段。 ns字段命名源数据库和集合,以便您可以区分哪个集合产生了每个更改。
集群范围的变更流示例
以下聚合从 cluster-changestream-connection 源摄取数据,该源连接到加载了样本 Mflix 数据集集合数据集的Atlas 集群,并写入到 cluster-changestream-sink-connection 目标,该目标连接到单独的Atlas 集群。要学习;了解如何创建流处理工作区以及如何将与Atlas 集群的连接添加到连接注册表中,请参阅Atlas Stream Processing入门。此聚合运行两个阶段,以在源集群上打开集群范围的变更流,并在目标集群上的 cluster_changes数据库中记录对 events集合的更改。此 $source 阶段省略了 db 和 coll,因此它报告源集群上每个数据库和集合的更改,而不是单个集合或数据库的更改。
$source阶段连接到cluster-changestream-connection源并打开针对整个源集群的变更流。config.startAtOperationTime字段指定流处理器应开始报告在指定时间或之后发生的更改。$merge阶段将变更流文档写入目标集群上cluster_changes数据库中名为events的Atlas集合。如果不存在此类数据库或集合, Atlas会进行创建。
[{ "$source": { "connectionName": "cluster-changestream-connection", "config": { "startAtOperationTime": {"$date": "2024-08-19T18:00:00.000Z"} } } }, { "$merge": { "into": { "connectionName": "cluster-changestream-sink-connection", "db": "cluster_changes", "coll": "events" } } }]
以下 mongosh 命令将文档插入到源集群上 sample_mflix数据库的 movies集合中:
db.getSiblingDB("sample_mflix").movies.insertOne( { _id: ObjectId("66c1a1f1f1f1f1f1f1f1f1f1"), title: "Example Movie" } )
插入文档后,流处理器会将变更流事件文档写入目标集群上的 cluster_changes.events集合。
要查看生成的 cluster_changes.events集合中的文档,请对目标集群运行以下命令。要进行连接,请参阅通过mongosh连接到集群。
db.getSiblingDB("cluster_changes").events.find()
输出文档中的 ns字段显示该更改源自源集群上的 sample_mflix.movies 插入。变更流聚合管道会将这一更改反映到目标集群上的 cluster_changes.events集合中。
文档数组示例
以下聚合使用内联文档数组作为流媒体数据源,其中包含三个位置的天气观测值。该数组使用与 示例天气数据集 相同的模式。此聚合运行三个阶段:
$source阶段将天气观测的内联documents数组定义为流媒体数据源,并使用timeField将每个文档的timestamp字段指定为权威时间戳。$match阶段仅将dewPoint.value大于5.0的文档传递到下一阶段。$merge阶段将输出写入sample_weatherstream数据库中名为stream的 Atlas 集合。如果不存在此类数据库或集合,Atlas 会创建它们。
[{ "$source": { "documents": [ { "location": "New York", "timestamp": ISODate('2024-01-15T08:00:00Z'), "temp": 23.5, "dewPoint": { "value": 6.2 } }, { "location": "Los Angeles", "timestamp": ISODate('2024-01-15T08:05:00Z'), "temp": 18.2, "dewPoint": { "value": 4.8 } }, { "location": "Chicago", "timestamp": ISODate('2024-01-15T08:10:00Z'), "temp": 26.8, "dewPoint": { "value": 7.5 } } ] } }, { "$match": { "dewPoint.value": { "$gt": 5.0 } } }, { "$merge": { "into": { "connectionName": "weatherStreamOutput", "db": "sample_weatherstream", "coll": "stream" } } }]
要查看生成的 sample_weatherstream.stream 集合中的文档,请连接到您的 Atlas 集群并运行以下命令:
db.getSiblingDB("sample_weatherstream").stream.find()
Solar Stream 示例
以下聚合筛选来自示例流媒体数据源sample_stream_solar 的报告。它将结果存档到 solar-cluster-connection 连接。要学习;了解如何创建流处理工作区并添加与Atlas Stream Processing工作区的连接,请参阅开始使用Atlas Stream Processing和添加Atlas Stream Processing连接。此聚合运行三个阶段来过滤来自 sample_stream_solar 源的报告,并将结果写入solarDb数据库中名为 solarColl 的集合:
$source阶段连接到sample_stream_solar源。timeField字段使用$dateFromString将每个传入报告的timestamp字段转换为日期。$match阶段会排除device_id为device_8的文档,并将来自所有其他设备的报告传递到下一阶段。$merge阶段将输出写入solar-cluster-connection集群的solarDb数据库中的solarColl集合。
[{ "$source": { "connectionName": "sample_stream_solar", "timeField": { "$dateFromString": { "dateString": "$timestamp" } } } }, { "$match": { "device_id": { "$ne": "device_8" } } }, { "$merge": { "into": { "connectionName": "solar-cluster-connection", "db": "solarDb", "coll": "solarColl" } } }]
注意
sample_stream_solar 源是一种仅供测试的连接,每秒都会生成示例文档,用于快速原型设计。
针对 solar-cluster-connection集群运行以下命令,验证流处理器的行为。要进行连接,请参阅通过mongosh连接到集群。
要查看生成的 solarDb.solarColl集合中的文档,运行:
db.getSiblingDB("solarDb").solarColl.find()
每个文档都包含太阳能设备报告的 device_id、event_type、group_id、max_watts、obs 和 timestamp 字段。您可以观察到输出中没有文档的 device_id 为 device_8,因为 $match 阶段在 $merge 阶段将剩余文档写入 solarColl 之前排除了这些报告。
注意
以上是一个有代表性的示例。流数据不是静态的,每个用户看到的都是不同的文档。
Cron 计划示例
以下聚合每五分钟生成一个文档,而不是从连接中读取。此聚合运行三个阶段:
$source阶段每五分钟的第二个0发出一个空文档(UTC 时间)。$project阶段使用作业名称标记每个文档,并将计划时间戳从_ts字段复制到runAt字段。$merge阶段将输出写入sample_weatherstream数据库中名为heartbeats的 Atlas 集合。如果不存在此类数据库或集合,Atlas 会创建它们。
[{ "$source": { "schedule": "0 0/5 * * * *" } }, { "$project": { "job": "five-minute-heartbeat", "runAt": "$_ts" } }, { "$merge": { "into": { "connectionName": "weatherStreamOutput", "db": "sample_weatherstream", "coll": "heartbeats" } } }]
要查看生成的 sample_weatherstream.heartbeats 集合中的文档,请连接到您的 Atlas 集群并运行以下命令:
db.getSiblingDB("sample_weatherstream").heartbeats.find()
注意
上述是一个具有代表性的示例。您看到的时间戳取决于您启动流处理器的时间。