定義
sp.process()現在の Stream Processing Workspace にエフェメラル Stream Processor を作成します。
このコマンドは、ストリーム処理ワークスペースに接続しているときにのみ呼び出せます。
このコマンドには、
mongoshバージョン `` 2.0が必要です。
互換性
このメソッドは Atlas Stream Processing ワークスペース でサポートされています。
構文
sp.process()メソッドの構文は次のとおりです。
sp.process( [ <pipeline> ], { <options> } )
コマンドフィールド
sp.process() 次のフィールドを取ります。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| 配列 | 必須 | ストリーム プロセッサーがストリーミング データに適用するストリーム集計パイプライン。 |
| オブジェクト | 任意 | ストリーム プロセッサーの任意設定を定義するオブジェクト。 |
| integer | 任意 | ターミナルに返すドキュメントの最大数。 |
| string | 任意 | Stream Processingに割り当てる階層。階層名は大文字と小文字を区別しません。このフィールドを省略すると、ストリームプロセッサーは Stream Processing Workspace のデフォルト階層を使用します。次のいずれかの値である必要があります。
詳細については、階層を参照してください。 |
動作
sp.process() 現在の Stream Processing ワークスペースに一時的な名前のないストリーム プロセッサを作成し、直ちに初期化します。このストリーム プロセッサは実行中にのみ持続します。一時的なストリーム プロセッサを終了する場合は、再度作成する必要があります。
一時的なストリーム プロセッサーは、次のいずれかが発生すると停止します。
プロセッサーは、
options.limitで設定されたドキュメントの数を返します。10 minutes elapse.
プロセッサを停止します。
アクセス制御
sp.process()を実行中ユーザーには atlasAdmin ロールが必要です。
例
次の例では、sample_stream_solar 接続からデータを取り込む一時的なストリームプロセッサーを作成します。プロセッサーは、device_id フィールドの値が device_8 であるすべてのドキュメントを除外し、残りのドキュメントを 10 秒間のタンブリング ウィンドウに渡します。各ウィンドウは受信したドキュメントをグループ化し、各グループの様々な統計情報を返します。ストリームプロセッサーは、これらのレコードをマージして mongodb1 接続を介して solar_db.solar_coll します。
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"] } } ] )