AI エージェント向け: ドキュメントインデックスは https://www.mongodb.com/ja-jp/docs/llms.txt で利用できます。すべてのページの markdown バージョンは、いずれかの URL パスに .md を追加することで利用できます。
Docs Menu

SP.process()(mongoshメソッド)

sp.process()

現在の Stream Processing ワークスペースにエフェメラル Stream Processor を作成します。

このメソッドは Atlas Stream Processing ワークスペース でサポートされています。

sp.process()メソッドの構文は次のとおりです。

sp.process(
[
<pipeline>
]
)

sp.process() 次のフィールドを取ります。

フィールド
タイプ
必要性
説明

pipeline

配列

必須

ストリーム プロセッサーがストリーミング データに適用するストリーム集計パイプライン

options

オブジェクト

任意

ストリーム プロセッサーの任意設定を定義するオブジェクト。

options.limit

integer

任意

ターミナルに返すドキュメントの最大数。sp.process() は、この数のドキュメントを返した後、セッションを終了します。

options.tier

string

任意

Stream Processingに割り当てる階層。階層名は大文字と小文字を区別しません。このフィールドを省略すると、Stream ProcessingはStream Processingワークスペースのデフォルト階層を使用します。次のいずれかの値である必要があります。

  • SP2

  • SP5

  • SP10

  • SP30

  • SP50

詳細については、階層を参照してください。

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"]
}
}
]
)