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

sp. 进程()(mongosh方法)

sp.process()

在当前 Stream Processing Workspace 上创建临时 Stream Processor

只有在连接到流处理工作区时才能调用此命令。

此命令需要mongosh版本 ≥ 2.0 。

Atlas Stream Processing 工作区支持此方法。

sp.process() 方法使用的语法如下:

sp.process(
[
<pipeline>
],
{
<options>
}
)

sp.process() 采用这些字段:

字段
类型
必要性
说明

pipeline

阵列

必需

流聚合管道是流处理器应用于流媒体数据的管道。

options

对象

Optional

为流处理器定义可选设置的对象。

options.limit

整型

Optional

要返回终端的最大文档数。sp.process() 返回此数量的文档后结束会话。

options.tier

字符串

Optional

要分配给流处理器的层级。层级名称不区分大小写。如果您省略此字段,则流处理器会使用 Stream Processing Workspace 的默认层级。必须是以下值之一:

  • SP2

  • SP5

  • SP10

  • SP30

  • SP50

要学习;了解更多信息,请参阅层级。

sp.process() 在当前 Stream Processing Workspace 上创建一个临时的未命名 Stream Processor,并立即将其初始化。此 Stream Processor 仅在运行时持续存在。如果终止临时 Stream Processor,则必须重新创建才能使用。

临时流处理器在发生以下任何一种情况时停止:

  • 处理器返回由 options.limit 设置的文档数。

  • 10 minutes elapse.

  • 您停止处理器。

运行sp.process() 的用户必须具有atlasAdmin角色。

以下示例创建一个暂时流处理器,该处理器从 sample_stream_solar 连接摄取数据。处理器排除 device_id 字段值为 device_8 的所有文档,并将其余文档传递给持续时间为 10 秒的 滚动窗口。每个窗口将接收到的文档分组,然后返回每个组的各种统计信息。流处理器会将这些记录合并到 solar_db.solar_coll 上的 mongodb1 连接。

sp.process(
[
{
$source: {
connectionName: 'sample_stream_solar',
timeField: {
$dateFromString: {
dateString: '$timestamp'
}
}
}
},
{
$match: {
$expr: {
$ne: [
"$device_id",
"device_8"
]
}
}
},
{
$tumblingWindow: {
interval: {
size: Int32(10),
unit: "second"
},
"pipeline": [
{
$group: {
"_id": { "device_id": "$device_id" },
"max_temp": { $max: "$obs.temp" },
"max_watts": { $max: "$obs.watts" },
"min_watts": { $min: "$obs.watts" },
"avg_watts": { $avg: "$obs.watts" },
"median_watts": {
$median: {
input: "$obs.watts",
method: "approximate"
}
}
}
}
]
}
},
{
$merge: {
into: {
connectionName: "mongodb1",
db: "solar_db",
coll: "solar_coll"
},
on: ["_id"]
}
}
]
)