ストリーミング マテリアライズドビューは、Atlas Stream Processing ストリーム プロセッサーが更新し続けるコレクションです。これを実現するため、ストリーム プロセッサーはソース コレクションへの変更を読み取り、その変更の影響を計算し、その結果をビューに適用します。ビューは、手動またはスケジュールされた更新なしでソース データの現在の状態を反映します。
ストリーミングマテリアライズドビューは、ダッシュボード、実行中の合計、開始アイテムのカウントなど、ソースデータの変更よりも計算結果の読み取りが多いワークロードに適しています。データが到着するとビューが更新されるため、読み取りにより、低レイテンシで現在の結果が返されます。
オンデマンド マテリアライズドビューとの比較
MongoDB は、マテリアライズドビューに対する2つのアプローチをサポートしています。
ストリーミングマテリアライズドビューは、継続的に実行されるAtlas Stream Processingストリームプロセッサの結果を保存します。
どの方法でも、計算結果をディスクに保存し、ビューから直接読み取りを行います。ビューの更新方法と時期が異なります。
特性 | オンデマンドのマテリアライズドビュー | ストリーミング マテリアライズドビュー |
|---|---|---|
trigger の更新 | 手動またはスケジュールされたもの | 変更駆動型、継続的 |
レイテンシ | 数分から数日 | 1 秒未満~数秒 |
データの新鮮度 | 点インタイム スナップショット | 永久に同期される |
コンピュート モデル | 全体の結果を再計算します | 増分効果を計算します |
Best fit | バッチするレポート作成、定期的集計 | ライブ ダッシュボード、運用分析 |
ストリーミングマテリアライズドビュープロセッサの特性
ストリーミング マテリアライズドビューの維持は、標準の Atlas Stream Processing 集計ステージから構成されるパスワードなしです。このパスワードなしを実装するストリーム プロセッサには、次の特徴があります。
変更ストリーム ソースを読み取ります。
fullDocumentとfullDocumentBeforeChangeがrequiredに設定された$sourceステージは、パイプラインが変更前後の各ドキュメントの状態を比較できるように、ソース コレクションから読み取ります。イベントごとに符号付きデルタを計算します。
$addFieldsステージでは$switch式を使用して、変更が計算された結果にどのように影響するかに応じて、各挿入、アップデート、または削除に正の値または負の値を割り当てることができます。グループの結果はウィンドウになります。ストリーム プロセッサは無制限ストリームで動作するため、すべての ステージはウィンドウステージ
$group内で実行する必要があります。ウィンドウ$groupは、ウィンドウ間隔にわたる各キーのデルタを合計します。ウィンドウ間隔はビューの更新間隔としても機能するため、ビューの更新度が決まります。設定できる最小の間隔は 1 ミリ秒です。結果を追加的に適用します。
$mergeステージでは、whenMatchedパイプラインを使用して、各ウィンドウの結果を置き換えではなく ビューの実行中合計に追加できます。Sink ステージで終了します。ストリーム プロセッサパイプラインはSink ステージで終了する必要があります。 Atlasコレクションに書込むには、
$mergeを使用します。
増分集計のみがストリーミングマテリアライズドビューに変換されます。完全なコレクションスキャンが必要な集計は対象外です。
考慮する必要がある動作の違い
ストリーミングのマテリアライズドビューは、ダウンストリームのコンシューマーに影響を与える方法でバッチ集計とは異なります。
ストリーミングマテリアライズドビューの作成
次のチュートリアルでは、購入方法別の完了した販売数を維持するストリーミングマテリアライズドビューを作成します。ストリーム プロセッサは、sample_supplies.sales コレクションの変更ストリームを読み取り、sample_supplies.sales_by_channel コレクションに書き込み (write)ます。
注意
ストリーム プロセッサは、 Apache Kafkaトピックからの読み取りと、 AWS S3 のApache Ops Manager テーブルへの書き込みも可能です。詳しくは、 「 Apache Kafkaブロック 」と「 $iceberg 集計ステージ 」を参照してください。
前提条件
Stream Processingプロセッサーを作成する前に、ソースデータを保持するクラスターへの Atlas 接続を含む Stream Processingワークスペース が必要です。接続を追加するには、接続の管理 を参照してください。
このセクションのコマンドをクラスターに対して実行します。接続するには、「mongosh を使用してクラスターに接続する」を参照してください。
ソース コレクションを準備し、ビューをシードするには、次の手順に従ってください。
サンプル データセットをロードします。
この手順では、付属のsales sample_issues データセットの コレクションを使用します。サンプルデータをロードする方法については、「 Atlas デプロイへのサンプル データのインポート 」を参照してください。
ビューに現在の結果をシードします。
クラスターに対して次のバッチする集計を実行し、ビューに現在のカウントを移入します。
db.sales.aggregate([ { $match: { status: "completed" } }, { $group: { _id: "$purchaseMethod", active_count: { $sum: 1 } } }, { $merge: { into: "sales_by_channel", whenMatched: "replace", whenNotMatched: "insert" } } ])
シードされたカウントを確認するには、ビューをクエリします。sales_by_channel コレクションには、購入方法ごとに 1 つのドキュメントが含まれています。
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585 }, { _id: 'In store', active_count: 2819 } ]
このシード集計自体がオンデマンドのマテリアライズドビューです。これらは補完的な関係です。オンデマンドのマテリアライズドビューをバッチする集計で初期化し、その後、同じコレクションを最新に保持するストリーム プロセッサを起動できます。オンデマンドビューはストリーミングマテリアライズドビューになります。
注意
パイプラインでウィンドウステージを使用しない場合は、代わりに initialSync を使用してビューにシードすることができます。この場合、ストリームプロセッサーはまず、ソース コレクション内の存在するすべてのドキュメントを挿入イベントとして取り込み、次に新しい変更イベントを処理します。$source オプション initialSync について学ぶには、MongoDB コレクション変更ストリーム。を参照してください。
手順
Atlas UIまたは を選択して、ストリームmongosh プロセッサを作成します。
ビューが最新の状態に保たれていることを確認します
このセクションのコマンドをクラスターに対して実行します。接続するには、「mongosh を使用してクラスターに接続する」を参照してください。
プロセッサーを開始すると、sales コレクションへの変更は数秒以内に sales_by_channel を更新します。これを確認するには、新しい完了したオンライン販売を挿入します。
db.sales.insertOne({ saleDate: new Date(), purchaseMethod: "Online", status: "completed", items: [], customer: {}, couponUsed: false })
プロセッサーは Online カウントを増分し、ウィンドウ境界を lastWindowStart にレコードします。
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1586, lastWindowStart: ISODate('2026-07-23T15:18:01.000Z') }, { _id: 'In store', active_count: 2819 } ]
カスタマーがその販売を返品すると、Online カウントはシード値に戻ります。
db.sales.updateOne( { purchaseMethod: "Online", status: "completed" }, { $set: { status: "returned" } } )
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585, lastWindowStart: ISODate('2026-07-23T15:18:13.000Z') }, { _id: 'In store', active_count: 2819 } ]
例
これらの例では、 Atlas Stream Processing例リポジトリのqueue_stats sv の例について説明します。優先順位レベルごとに 1 つのドキュメントで、オープンなサポート チケットの実行中数を保持する コレクションを維持します。ストリーム プロセッサはsupport_tickets 変更ストリームを読み取り、チケットが開かれ、解決され、削除されるにつれて各カウントを更新します。
パイプラインは、チケットが開いたときに +1、チケットが解決または削除されたときに -1 だけチケットの優先順位のカウントを調整します。集計には5つのステージがあります。
$sourceステージは、事前および書き込みのイメージを使用してsupport_tickets変更ストリームを読み取ります。$addFieldsステージでは、$switchのoperationTypeを使用して を計算し、_delta_priorityをグループキーとして抽出します。ステージでは、デルタが 0
$matchのイベントが削除されます。$tumblingWindowステージは、1 秒の各ウィンドウ内で優先順位によってデルタを合計します。$mergeステージでは、リプレイでのダブルカウントを回避するために、lastWindowStartを高値として使用して、各ウィンドウのデルタを実行中の合計に追加します。
[ { "$source": { "connectionName": "<connection-name>", "db": "support", "coll": "support_tickets", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }, { "$addFields": { "_delta": { "$switch": { "branches": [ { "case": { "$and": [ { "$eq": ["$operationType", "insert"] }, { "$eq": ["$fullDocument.status", "open"] } ] }, "then": 1 }, { "case": { "$and": [ { "$eq": ["$operationType", "update"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] }, { "$eq": ["$fullDocument.status", "resolved"] } ] }, "then": -1 }, { "case": { "$and": [ { "$eq": ["$operationType", "delete"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] } ] }, "then": -1 } ], "default": 0 } }, "_priority": { "$ifNull": [ "$fullDocument.priority", "$fullDocumentBeforeChange.priority" ] } } }, { "$match": { "_delta": { "$ne": 0 } } }, { "$tumblingWindow": { "boundary": "processingTime", "interval": { "size": 1, "unit": "second" }, "pipeline": [ { "$group": { "_id": "$_priority", "open_count": { "$sum": "$_delta" }, "windowStart": { "$first": { "$meta": "stream.window.start" } } } } ] } }, { "$merge": { "into": { "connectionName": "<connection-name>", "db": "support", "coll": "queue_stats" }, "whenMatched": [ { "$set": { "open_count": { "$cond": [ { "$gt": ["$$new.windowStart", { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }] }, { "$add": ["$open_count", "$$new.open_count"] }, "$open_count" ] }, "lastWindowStart": { "$max": [ { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }, "$$new.windowStart" ] } } } ], "whenNotMatched": "insert" } } ]
queue_stats コレクション内の各ドキュメントは、次のようになります。
{ _id: "P1", open_count: <num>, lastWindowStart: <timestamp> }
1 つのイベントで 2 つのグループ キーを調整する必要がある場合は、スカラー _delta を _adjustments 配列に置き換え、配列を調整ごとに 1 つのドキュメントに分散させます。単一キー パイプラインを次のように変更します。
$addFieldsステージを置き換えて、各$switch_adjustmentsブランチが、影響を受けたグループキーごとに 1 つの要素を含む 配列を返すようにします。エスカレーションは 2 つの要素を返します。
単一キー パイプラインのステージ 2 と 3 を次のように置き換えます。
// Stage 2 (replacement): Compute an _adjustments array. { $addFields: { _adjustments: { $switch: { branches: [ { case: { $and: [ { $eq: ["$operationType", "insert"] }, { $eq: ["$fullDocument.status", "open"] } ]}, then: [{ _priority: "$fullDocument.priority", _delta: 1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] }, { $eq: ["$fullDocument.status", "resolved"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocument.status", "open"] }, { $ne: ["$fullDocument.priority", "$fullDocumentBeforeChange.priority"] } ]}, then: [ { _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }, { _priority: "$fullDocument.priority", _delta: 1 } ] }, { case: { $and: [ { $eq: ["$operationType", "delete"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] } ], default: [] } } } }, // Stage 3 (replacement): Fan out into one document per adjustment, // then lift the adjustment fields back to the top level. { $unwind: "$_adjustments" }, { $set: { _priority: "$_adjustments._priority", _delta: "$_adjustments._delta" } }
詳細情報
このガイドのステージと概念についてさらに学ぶには、次のリソースを参照してください。
ストリーム プロセッサーの作成、開始、停止、モニターについては、ストリーム プロセッサーの開発と管理を参照してください。
Atlas Stream Processing がサポートする集計ステージについて学ぶには、集計パイプライン ステージを参照してください。
Stream Processing が読み取り可能なソースについて学ぶには、
$sourceステージ (Stream Processing)。 を参照してください。ウィンドウステージについて学ぶには、Stream Processor Windows. を参照してください。
オンデマンドのマテリアライズドビューの詳細については、「オンデマンドのマテリアライズドビュー」を参照してください。