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

ストリーミングマテリアライズドビューをビルドする

ストリーミング マテリアライズドビューは、Atlas Stream Processing ストリーム プロセッサーが更新し続けるコレクションです。これを実現するため、ストリーム プロセッサーはソース コレクションへの変更を読み取り、その変更の影響を計算し、その結果をビューに適用します。ビューは、手動またはスケジュールされた更新なしでソース データの現在の状態を反映します。

ストリーミングマテリアライズドビューは、ダッシュボード、実行中の合計、開始アイテムのカウントなど、ソースデータの変更よりも計算結果の読み取りが多いワークロードに適しています。データが到着するとビューが更新されるため、読み取りにより、低レイテンシで現在の結果が返されます。

MongoDB は、マテリアライズドビューに対する2つのアプローチをサポートしています。

  • オンデマンドのマテリアライズドビューには、手動またはスケジュールで実行した集計パイプラインの結果が保存されます。

  • ストリーミングマテリアライズドビューは、継続的に実行されるAtlas Stream Processingストリームプロセッサの結果を保存します。

どの方法でも、計算結果をディスクに保存し、ビューから直接読み取りを行います。ビューの更新方法と時期が異なります。

特性
オンデマンドのマテリアライズドビュー
ストリーミング マテリアライズドビュー

trigger の更新

手動またはスケジュールされたもの

変更駆動型、継続的

レイテンシ

数分から数日

1 秒未満~数秒

データの新鮮度

点インタイム スナップショット

永久に同期される

コンピュート モデル

全体の結果を再計算します

増分効果を計算します

Best fit

バッチするレポート作成、定期的集計

ライブ ダッシュボード、運用分析

ストリーミング マテリアライズドビューの維持は、標準の Atlas Stream Processing 集計ステージから構成されるパスワードなしです。このパスワードなしを実装するストリーム プロセッサには、次の特徴があります。

  • 変更ストリーム ソースを読み取ります。fullDocumentfullDocumentBeforeChangerequired に設定された $source ステージは、パイプラインが変更前後の各ドキュメントの状態を比較できるように、ソース コレクションから読み取ります。

  • イベントごとに符号付きデルタを計算します。 $addFieldsステージでは$switch 式を使用して、変更が計算された結果にどのように影響するかに応じて、各挿入、アップデート、または削除に正の値または負の値を割り当てることができます。

  • グループの結果はウィンドウになります。ストリーム プロセッサは無制限ストリームで動作するため、すべての ステージはウィンドウステージ$group 内で実行する必要があります。ウィンドウ$group は、ウィンドウ間隔にわたる各キーのデルタを合計します。ウィンドウ間隔はビューの更新間隔としても機能するため、ビューの更新度が決まります。設定できる最小の間隔は 1 ミリ秒です。

  • 結果を追加的に適用します。 $mergeステージでは、whenMatched パイプラインを使用して、各ウィンドウの結果を置き換えではなく ビューの実行中合計に追加できます。

  • Sink ステージで終了します。ストリーム プロセッサパイプラインはSink ステージで終了する必要があります。 Atlasコレクションに書込むには、$merge を使用します。

増分集計のみがストリーミングマテリアライズドビューに変換されます。完全なコレクションスキャンが必要な集計は対象外です。

ストリーミングのマテリアライズドビューは、ダウンストリームのコンシューマーに影響を与える方法でバッチ集計とは異なります。

  • ビューはゼロから始まります。デフォルトでは、プロセッサは既存のドキュメントを読み取らないため、プロセッサを開始する前にビューをシードするか、$source ステージで initialSync を有効にします。

  • ソースが更新を操作する変更のみ。パイプラインが $lookupステージを使用している場合、参照コレクションへのその後の変更によって、ビューがすでに書き込んだドキュメントは更新されません。

次のチュートリアルでは、購入方法別の完了した販売数を維持するストリーミングマテリアライズドビューを作成します。ストリーム プロセッサは、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 を使用してクラスターに接続する」を参照してください。

ソース コレクションを準備し、ビューをシードするには、次の手順に従ってください。

1

この手順では、付属のsales sample_issues データセットの コレクションを使用します。サンプルデータをロードする方法については、「 Atlas デプロイへのサンプル データのインポート 」を参照してください。

2

ストリームプロセッサが完了して返品された販売を検出できるように、status フィールドを追加します。クラスターに対して次のコマンドを実行します。

db.sales.updateMany(
{ status: { $exists: false } },
{ $set: { status: "completed" } }
)
{
acknowledged: true,
insertedId: null,
matchedCount: 5000,
modifiedCount: 5000,
upsertedCount: 0
}
3

ストリームプロセッサーが各変更の影響を計算できるように、ソースコレクションでプレイメージと書き込みイメージを有効にします。クラスターに対して次のコマンドを実行します。

db.getSiblingDB("sample_supplies").runCommand({
collMod: "sales",
changeStreamPreAndPostImages: { enabled: true }
})
{ ok: 1, ... }
4

クラスターに対して次のバッチする集計を実行し、ビューに現在のカウントを移入します。

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つのステージがあります。

  1. $source ステージは、事前および書き込みのイメージを使用して support_tickets 変更ストリームを読み取ります。

  2. $addFieldsステージでは、$switchoperationType を使用して を計算し、_delta _priorityをグループキーとして抽出します。

  3. ステージでは、デルタが 0$match のイベントが削除されます。

  4. $tumblingWindow ステージは、1 秒の各ウィンドウ内で優先順位によってデルタを合計します。

  5. $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 つの調整にファンアウトします。

1 つのイベントで 2 つのグループ キーを調整する必要がある場合は、スカラー _delta_adjustments 配列に置き換え、配列を調整ごとに 1 つのドキュメントに分散させます。単一キー パイプラインを次のように変更します。

  1. $addFieldsステージを置き換えて、各$switch _adjustmentsブランチが、影響を受けたグループキーごとに 1 つの要素を含む 配列を返すようにします。エスカレーションは 2 つの要素を返します。

  2. の後に $unwindステージと ステージを追加します。$set $addFields$unwindステージは、各イベントを調整ごとに 1 つのドキュメントに分割し、空の配列を削除し、 ステージを置き換えます。$match $setステージでは、_adjustments フィールドが最上位の_priority_delta に引き上げられ、残りのステージは変更されません。

単一キー パイプラインのステージ 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"
}
}

このガイドのステージと概念についてさらに学ぶには、次のリソースを参照してください。