Overview
このガイドでは、変更ストリームを使用してデータに対するリアルタイムの変更を監視する方法を学習できます。 変更ストリームは、アプリケーションがコレクション、データベース、または配置のデータ変更をサブスクライブできる MongoDB Server の機能です。
Tip
Atlas Stream Processing
変更ストリームの代わりに、Atlas Stream Processing を使用してデータのストリームを処理および変換できます。データベースイベントのみを登録する変更ストリームとは異なり、Atlas Stream Processing は複数のデータイベント型を管理し、拡張データプロセシング機能を提供します。この機能の詳細については、 MongoDB AtlasドキュメントのAtlas Stream Processingを参照してください。
サンプル データ
このガイドの例では、Atlas サンプル データセットの sample_restaurants.restaurants コレクションを使用します。無料の MongoDB Atlas cluster を作成し、サンプル データセットをロードする方法については、クイック スタートを参照してください。
このページの例では、次の Restaurant クラス、Address クラス、GradeEntry クラスをモデルとして使用します。
public class Restaurant { public ObjectId Id { get; set; } public string Name { get; set; } [] public string RestaurantId { get; set; } public string Cuisine { get; set; } public Address Address { get; set; } public string Borough { get; set; } public List<GradeEntry> Grades { get; set; } }
public class Address { public string Building { get; set; } [] public double[] Coordinates { get; set; } public string Street { get; set; } [] public string ZipCode { get; set; } }
public class GradeEntry { public DateTime Date { get; set; } public string Grade { get; set; } public float? Score { get; set; } }
注意
restaurantsコレクションのドキュメントは、スニペット ケースの命名規則を使用します。このガイドの例では、ConventionPack を使用してコレクション内のフィールドをパスカル ケースに逆シリアル化し、Restaurantクラスのプロパティにマップします。
カスタム直列化について詳しくは、「カスタム直列化」を参照してください。
変更ストリームを開く
変更ストリームを開くには、 メソッドまたはWatch() WatchAsync()メソッドを呼び出します。メソッドを呼び出す インスタンスによって、変更ストリームがリッスンするイベントの範囲が決まります。 Watch()WatchAsync()次のクラスで メソッドまたは メソッドを呼び出すことができます。
MongoClient: MongoDB 配置のすべての変更を監視Database: データベース内のすべてのコレクションの変更を監視するにはCollection: コレクションの変更をモニターするには
次の例では、 restaurantsコレクションの変更ストリームを開き、変更が発生に応じて出力します。 AsynchronousSynchronous対応するコードを表示するには、 タブまたは タブを選択します。
var database = client.GetDatabase("sample_restaurants"); var collection = database.GetCollection<Restaurant>("restaurants"); // Opens a change streams and print the changes as they're received using var cursor = await collection.WatchAsync(); await cursor.ForEachAsync(change => { Console.WriteLine("Received the following type of change: " + change.BackingDocument); });
var database = client.GetDatabase("sample_restaurants"); var collection = database.GetCollection<Restaurant>("restaurants"); // Opens a change stream and prints the changes as they're received using (var cursor = collection.Watch()) { foreach (var change in cursor.ToEnumerable()) { Console.WriteLine("Received the following type of change: " + change.BackingDocument); } }
変更の監視を開始するには、アプリケーションを実行します。 次に、別のアプリケーションまたは shell で、 restaurantsコレクションを変更します。 "name"の値が"Blarney Castle"であるドキュメントを更新すると、次の変更ストリーム出力が生成されます。
{ "_id" : { "_data" : "..." }, "operationType" : "update", "clusterTime" : Timestamp(...), "wallTime" : ISODate("..."), "ns" : { "db" : "sample_restaurants", "coll" : "restaurants" }, "documentKey" : { "_id" : ObjectId("...") }, "updateDescription" : { "updatedFields" : { "cuisine" : "Irish" }, "removedFields" : [], "truncatedArrays" : [] } }
変更ストリーム出力の変更
変更ストリーム出力を変更するには、 パラメータを メソッドと メソッドに渡します。pipelineWatch()WatchAsync()このパラメーターを使用すると、指定された変更イベントのみを監視できます。 EmptyPipelineDefinitionクラスを使用し、関連する集計ステージ メソッドを追加して、パイプラインを作成します。
pipelineパラメーターでは次の集計ステージを指定できます。
$addFields$changeStreamSplitLargeEvent$match$project$replaceRoot$replaceWith$redact$set$unset
Tip
PipelineDefinitionBuilderクラスを使用して集計パイプラインを構築する方法については、「 ビルダによる操作の集計パイプラインの構築 」を参照してください。
変更ストリーム出力の変更の詳細については、MongoDB Server マニュアルの「 変更ストリーム出力 の変更 」セクションを参照してください。
更新イベントの監視例
次の例では、 pipelineパラメータを使用して、アップデート操作のみを記録する変更ストリームを開きます。 AsynchronousSynchronous対応するコードを表示するには、 タブまたは タブを選択します。
var pipeline = new EmptyPipelineDefinition<ChangeStreamDocument<Restaurant>>() .Match(change => change.OperationType == ChangeStreamOperationType.Update); // Opens a change stream and prints the changes as they're received using (var cursor = await collection.WatchAsync(pipeline)) { await cursor.ForEachAsync(change => { Console.WriteLine("Received the following change: " + change); }); }
var pipeline = new EmptyPipelineDefinition<ChangeStreamDocument<Restaurant>>() .Match(change => change.OperationType == ChangeStreamOperationType.Update); // Opens a change streams and print the changes as they're received using (var cursor = collection.Watch(pipeline)) { foreach (var change in cursor.ToEnumerable()) { Console.WriteLine("Received the following change: " + change); } }
大規模な変更イベントの分裂例
アプリケーションが生成した変更イベントが16 MB を超えるサイズの場合、サーバーはBSONObjectTooLarge エラーを返します。 このエラーを回避するには、$changeStreamSplitLargeEventパイプラインステージを使用してイベントを小さなフラグメントに分裂。 .NET/ C#ドライバー集計API には ChangeStreamSplitLargeEvent() メソッドが含まれており、このメソッドを使用して $changeStreamSplitLargeEvent ステージを変更ストリームパイプラインに追加できます。
この例では、 16 MB の制限を超える変更を監視し、変更イベントを分裂にドライバーに指示します。 このコードは、各イベントの変更ドキュメントを出力し、ヘルパーメソッドを呼び出してイベントフラグメントを再アセンブルします。
var pipeline = new EmptyPipelineDefinition<ChangeStreamDocument<Restaurant>>() .ChangeStreamSplitLargeEvent(); using var cursor = await collection.WatchAsync(pipeline); await foreach (var completeEvent in GetNextChangeStreamEventAsync(cursor)) { Console.WriteLine("Received the following change: " + completeEvent.BackingDocument); }
var pipeline = new EmptyPipelineDefinition<ChangeStreamDocument<Restaurant>>() .ChangeStreamSplitLargeEvent(); using var cursor = collection.Watch(pipeline); foreach (var completeEvent in GetNextChangeStreamEvent(cursor.ToEnumerable().GetEnumerator())) { Console.WriteLine("Received the following change: " + completeEvent.BackingDocument); }
注意
前述の例に示すように、変更イベントフラグメントを再アセンブルすることをお勧めしますが、この手順は任意です。 同じロジックを使用して、分裂と完了した変更イベントの両方を監視できます。
上記の例では、GetNextChangeStreamEvent()、GetNextChangeStreamEventAsync()、MergeFragment() メソッドを使用して、変更イベントフラグメントを単一の変更ストリームドキュメントに再アセンブルします。 次のコードは、これらのメソッドを定義します。
// Fetches the next complete change stream event private static async IAsyncEnumerable<ChangeStreamDocument<TDocument>> GetNextChangeStreamEventAsync<TDocument>( IAsyncCursor<ChangeStreamDocument<TDocument>> changeStreamCursor) { var changeStreamEnumerator = GetNextChangeStreamEventFragmentAsync(changeStreamCursor).GetAsyncEnumerator(); while (await changeStreamEnumerator.MoveNextAsync()) { var changeStreamEvent = changeStreamEnumerator.Current; if (changeStreamEvent.SplitEvent != null) { var fragment = changeStreamEvent; while (fragment.SplitEvent.Fragment < fragment.SplitEvent.Of) { await changeStreamEnumerator.MoveNextAsync(); fragment = changeStreamEnumerator.Current; MergeFragment(changeStreamEvent, fragment); } } yield return changeStreamEvent; } } private static async IAsyncEnumerable<ChangeStreamDocument<TDocument>> GetNextChangeStreamEventFragmentAsync<TDocument>( IAsyncCursor<ChangeStreamDocument<TDocument>> changeStreamCursor) { while (await changeStreamCursor.MoveNextAsync()) { foreach (var changeStreamEvent in changeStreamCursor.Current) { yield return changeStreamEvent; } } } // Merges a fragment into the base event private static void MergeFragment<TDocument>( ChangeStreamDocument<TDocument> changeStreamEvent, ChangeStreamDocument<TDocument> fragment) { foreach (var element in fragment.BackingDocument) { if (element.Name != "_id" && element.Name != "splitEvent") { changeStreamEvent.BackingDocument[element.Name] = element.Value; } } }
// Fetches the next complete change stream event private static IEnumerable<ChangeStreamDocument<TDocument>> GetNextChangeStreamEvent<TDocument>( IEnumerator<ChangeStreamDocument<TDocument>> changeStreamEnumerator) { while (changeStreamEnumerator.MoveNext()) { var changeStreamEvent = changeStreamEnumerator.Current; if (changeStreamEvent.SplitEvent != null) { var fragment = changeStreamEvent; while (fragment.SplitEvent.Fragment < fragment.SplitEvent.Of) { changeStreamEnumerator.MoveNext(); fragment = changeStreamEnumerator.Current; MergeFragment(changeStreamEvent, fragment); } } yield return changeStreamEvent; } } // Merges a fragment into the base event private static void MergeFragment<TDocument>( ChangeStreamDocument<TDocument> changeStreamEvent, ChangeStreamDocument<TDocument> fragment) { foreach (var element in fragment.BackingDocument) { if (element.Name != "_id" && element.Name != "splitEvent") { changeStreamEvent.BackingDocument[element.Name] = element.Value; } } }
Tip
大規模な変更イベントの分割の詳細については、 MongoDB Serverマニュアルの $changeStreamSplitLargeEvent を参照してください。
Watch() 動作を変更する
Watch()メソッドとWatchAsync()メソッドは、操作を構成するために使用できるオプションを表す任意のパラメーターを受け入れます。 オプションを指定しない場合、ドライバーは操作をカスタマイズしません。
次の表では、 Watch()とWatchAsync()の動作をカスタマイズするために設定できるオプションについて説明します。
オプション | 説明 |
|---|---|
| ドキュメントに加えられた変更のみを表示するのではなく、変更後に完全なドキュメントを表示するかどうかを指定します。 このオプションの詳細については、「変更前イメージと変更後イメージを含める」を参照してください。 |
| ドキュメントに加えられた変更のみを表示するのではなく、変更前のドキュメント全体を表示するかどうかを指定します。 このオプションの詳細については、「変更前イメージと変更後イメージを含める」を参照してください。 |
|
|
|
|
|
|
| 空のバッチするを返す前に、新しいデータ変更が変更ストリームカーソルに報告されるまでサーバーが待機する最大時間をミリ秒単位で指定します。 デフォルトは 1000 ミリ秒です。 |
| MongoDB Server v 6.0以降、 変更ストリームは、 |
| 変更ストリームが各バッチで返すことができるドキュメントの最大数を指定します。これは |
| 変更ストリームカーソルに使用する 照合 を指定します。 |
| 操作にコメントを付けます。 |
変更前と変更後のイメージを含めます
重要
配置で MongoDB v 6.0以降が使用されている場合にのみ、コレクションで変更前と変更後のイメージを有効にできます。
デフォルトでは、コレクションに対して操作を実行すると、対応する変更イベントにはその操作によって変更されたフィールドのデルタのみが含まれます。 変更前または変更後の完全なドキュメントを表示するには、 ChangeStreamOptionsオブジェクトを作成し、 FullDocumentBeforeChangeまたはFullDocumentオプションを指定します。 次に、 ChangeStreamOptionsオブジェクトをWatch()またはWatchAsync()メソッドに渡します。
変更前のイメージは、変更前のドキュメントの完全なバージョンです。 変更ストリーム イベントに変更前のイメージを含めるには、 FullDocumentBeforeChangeオプションを次のいずれかの値に設定します。
ChangeStreamFullDocumentBeforeChangeOption.WhenAvailable: 変更イベントには、変更前のイメージが利用可能な場合にのみ、 変更イベント 用の変更されたドキュメントの変更前のイメージが含まれます。ChangeStreamFullDocumentBeforeChangeOption.Required: 変更イベントには、変更イベント用に変更されたドキュメントの変更前のイメージが含まれます。 変更前のイメージが利用できない場合、ドライバーはエラーを発生させます。
変更後のイメージとは、変更後のドキュメントの完全なバージョンです。 変更ストリーム イベントに変更後のイメージを含めるには、 FullDocumentオプションを次のいずれかの値に設定します。
ChangeStreamFullDocumentOption.UpdateLookup: 変更イベントには、変更後一定時間の変更されたドキュメント全体のコピーが含まれます。ChangeStreamFullDocumentOption.WhenAvailable: 変更イベントには、変更後のイメージが利用可能な場合にのみ、 変更イベント 用の変更されたドキュメントの変更後のイメージが含まれます。ChangeStreamFullDocumentOption.Required: 変更イベントには、変更イベントの変更されたドキュメントの変更後のイメージが含まれます。 変更後のイメージが利用できない場合、ドライバーはエラーを発生させます。
次の例では、コレクションの変更ストリームを開き、 FullDocumentオプションを指定して更新されたドキュメントの変更後のイメージを含めます。 AsynchronousSynchronous対応するコードを表示するには、 タブまたは タブを選択します。
var pipeline = new EmptyPipelineDefinition<ChangeStreamDocument<Restaurant>>() .Match(change => change.OperationType == ChangeStreamOperationType.Update); var options = new ChangeStreamOptions { FullDocument = ChangeStreamFullDocumentOption.UpdateLookup, }; using var cursor = await collection.WatchAsync(pipeline, options); await cursor.ForEachAsync(change => { Console.WriteLine(change.FullDocument.ToBsonDocument()); });
var pipeline = new EmptyPipelineDefinition<ChangeStreamDocument<Restaurant>>() .Match(change => change.OperationType == ChangeStreamOperationType.Update); var options = new ChangeStreamOptions { FullDocument = ChangeStreamFullDocumentOption.UpdateLookup, }; using (var cursor = collection.Watch(pipeline, options)) { foreach (var change in cursor.ToEnumerable()) { Console.WriteLine(change.FullDocument.ToBsonDocument()); } }
上記のコード例を実行し、 "name"値が"Blarney Castle"であるドキュメントを更新すると、次の変更ストリーム出力が生成されます。
{ "_id" : ObjectId("..."), "name" : "Blarney Castle", "restaurant_id" : "40366356", "cuisine" : "Traditional Irish", "address" : { "building" : "202-24", "coord" : [-73.925044200000002, 40.5595462], "street" : "Rockaway Point Boulevard", "zipcode" : "11697" }, "borough" : "Queens", "grades" : [...] }
変更前と変更後のイメージの詳細については、Change Streams MongoDB Serverマニュアルの「 とドキュメントの変更 前イメージおよび変更後イメージ 」を参照してください。
詳細情報
Change Streams変更ストリームの詳細については、MongoDB Server マニュアルの 「 ストリーム」 を参照してください。
API ドキュメント
このガイドで説明したメソッドや型の詳細については、次の API ドキュメントを参照してください。