定义
sp.process()在当前 Stream Processing Workspace 上创建临时 Stream Processor。
只有在连接到流处理工作区时才能调用此命令。
此命令需要
mongosh版本 ≥ 2.0 。
兼容性
Atlas Stream Processing 工作区支持此方法。
语法
sp.process() 方法使用的语法如下:
sp.process( [ <pipeline> ], { <options> } )
命令字段
sp.process() 采用这些字段:
行为
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"] } } ] )