コレクションでの変更ストリームの開始
コレクションでwatch()を呼び出し、非同期反復を使ってNode.jsアプリケーションでイベントストリームを消費します。
「コレクションでの変更ストリームの開始」はCoddyKit上の無料MongoDB Academyレッスンです。 これはレッスン1/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはMongoDB Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 MongoDB Academyコースには全4レッスンが含まれています。
チェンジストリームとは
チェンジストリームは、MongoDBのコレクション、データベース、またはデプロイメント全体に対するinsert、update、replace、delete、invalidate操作のリアルタイムイベントフィードを提供します。MongoDB 3.6で導入されたチェンジストリームは、レプリカセットのレプリケーションジャーナルであるoplog(operation log)の上に構築されていますが、生のoplog形式を解析する必要はなく、高レベルで再開可能なカーソルAPIを利用できます。
チェンジストリームの前提条件
チェンジストリームにはレプリカセットまたはシャーディングクラスターが必要です。oplogに依存するため、スタンドアロンのMongoDBインスタンスでは動作しません。MongoDB Atlasでは、すべてのクラスター(無料のM0を含む)がレプリカセットであるため、チェンジストリームはそのまま利用できます。ローカル開発では、--replSet rs0を指定してmongodを起動し、mongoshでrs.initiate()を実行してレプリカセットを初期化する必要があります。
// Verify you're on a replica set before using change streams
// In mongosh:
rs.status() // should show replica set members, not an error
// If running locally without a replica set, start one:
// mongod --replSet rs0 --port 27017 --dbpath /data/db
// Then in mongosh: rs.initiate()watch()でチェンジストリームを開く
collection.watch()を呼び出すと、特定のコレクションに対するチェンジストリームを開けます。このメソッドはChangeStreamカーソルオブジェクトを返し、非同期反復、next()メソッド、またはイベントリスナーを使用して反復できます。ストリームは開いたままになり、イベントが発生すると配信します。パイプラインを指定しない空のwatch()呼び出しでは、そのコレクションのすべての変更イベント種別を受信します。
const { MongoClient } = require('mongodb');
const client = new MongoClient(process.env.MONGODB_URI);
async function watchCollection() {
const db = client.db('ecommerce');
const orders = db.collection('orders');
// Open a change stream on the orders collection
const changeStream = orders.watch();
console.log('Watching for changes...');
// The stream is now open and ready to deliver events
return changeStream;
}非同期反復でイベントを消費する
最新のNode.jsでチェンジストリームのイベントを消費する最も読みやすい方法は、for await...ofを使用した非同期反復です。この構文はカーソルのnext()呼び出しを自動的に処理し、イベント間ではループを一時停止します。チェンジストリームが閉じられるかエラーが発生するまで、ループは無期限に実行されます。プロセス終了時にストリームが確実に閉じられるよう、必ずループをtry/finallyで囲んでください。
async function processChanges() {
const changeStream = db.collection('orders').watch();
try {
for await (const change of changeStream) {
console.log('Change received:', change.operationType);
console.log('Document key:', change.documentKey._id);
// Handle the change event
await handleOrderChange(change);
}
} finally {
await changeStream.close();
}
}EventEmitterでイベントを消費する
別の方法として、チェンジストリームをNode.jsのEventEmitterとして使用できます。この方法は、コード内の他の場所でストリームを使用している場合に親しみやすいでしょう。通常のイベントには'change'イベントハンドラーを、接続に関する問題には'error'ハンドラーを登録します。この形式は、awaitループで関数をブロックせずにイベントへ反応したい場合に便利です。
const changeStream = db.collection('products').watch();
changeStream.on('change', (event) => {
console.log('Product changed:', event.operationType, event.documentKey._id);
if (event.operationType === 'update') {
invalidateProductCache(event.documentKey._id);
}
});
changeStream.on('error', (error) => {
console.error('Change stream error:', error);
// Implement reconnect logic or close gracefully
});
// Close when done
// changeStream.close();チェンジストリームのスコープ:コレクション、データベース、クライアント
チェンジストリームは、3つのスコープで開けます。コレクションレベル(1つのコレクションを監視)、データベースレベル(データベース内のすべてのコレクションを監視)、クライアントレベル(デプロイメント内のすべてのデータベースとコレクションを監視)です。スコープが広いほど、イベント量は増加します。ほとんどのアプリケーションでは、必要なイベントだけを受信できるよう、特定のコレクションを監視します。
// Collection-level (most common)
const stream1 = db.collection('orders').watch();
// Database-level — all collections in 'ecommerce'
const stream2 = client.db('ecommerce').watch();
// Client-level — everything in the deployment
const stream3 = client.watch();
// Events at db/client scope include 'ns' field to identify which collection changed
for await (const change of stream2) {
console.log('Changed collection:', change.ns.coll);
}チェンジストリームのオプション:fullDocument
デフォルトでは、更新イベントに含まれるのはドキュメント全体ではなく、変更されたフィールド(更新の説明)だけです。イベントのペイロードに更新後のドキュメント全体が必要な場合は、{ fullDocument: 'updateLookup' }をwatch()に渡します。これにより、MongoDBは更新後にドキュメントを追加で検索し、その内容を変更イベントに含めます。これによってレイテンシが増加し、イベント後に別の読み取りが実行される点に注意してください。
// Receive the full document in update events
const changeStream = db.collection('users').watch(
[], // empty pipeline = all events
{ fullDocument: 'updateLookup' }
);
for await (const change of changeStream) {
if (change.operationType === 'update') {
// change.fullDocument is now the complete updated user document
console.log('Updated user:', change.fullDocument.email);
await syncToSearchIndex(change.fullDocument);
}
}リアルタイムダッシュボードのユースケース
チェンジストリームの典型的なユースケースは、新しい注文が届くたびに表示するライブダッシュボードを実現することです。新しい注文ドキュメントが挿入されるとチェンジストリームがトリガーされ、Node.jsバックエンドはWebSocketまたはServer-Sent Eventsを介して接続中のクライアントに更新をプッシュできます。これによりポーリングが不要になり、クエリを繰り返してデータベースに負荷をかけることなく、真のリアルタイム更新を実現できます。
// Server: push new orders to dashboard clients via WebSocket
async function startOrderWatcher(io) { // io = socket.io instance
const changeStream = db.collection('orders').watch([
{ $match: { operationType: 'insert' } }
]);
for await (const change of changeStream) {
const newOrder = change.fullDocument;
// Broadcast to all connected dashboard clients
io.to('dashboard').emit('newOrder', {
id: newOrder._id,
customer: newOrder.customerId,
amount: newOrder.total,
timestamp: newOrder.createdAt
});
}
}チェンジストリームを使用したイベント駆動マイクロサービス
チェンジストリームは、単純なイベント駆動マイクロサービスパターンにおいて、KafkaやRabbitMQなどのメッセージブローカーの代わりに使用できます。Service AがMongoDBに書き込むと、Service Bがコレクションを監視して変更に反応します。これにより、イベント量が少ない、または中程度の場合に、別個のメッセージバスを運用する複雑さを避けられます。ただし、高スループットのユースケースでは、専用のメッセージブローカーのほうがチェンジストリームより優れた保証と高いスループットを提供します。
// Inventory service reacts to confirmed orders
async function inventoryWatcher() {
const changeStream = db.collection('orders').watch([
{
$match: {
operationType: 'update',
'updateDescription.updatedFields.status': 'confirmed'
}
}
], { fullDocument: 'updateLookup' });
for await (const change of changeStream) {
const order = change.fullDocument;
for (const item of order.items) {
await reserveInventory(item.productId, item.quantity);
}
}
}チェンジストリームを適切に閉じる
チェンジストリームはMongoDBサーバーへの永続的な接続を使用します。プロセスをシャットダウンするときや、ストリームが不要になったときは、必ず閉じてください。changeStream.close()を呼び出すとPromiseが返されます。長時間実行するアプリケーションでは、プロセスシグナル(SIGTERM、SIGINT)を監視し、終了前にストリームとクライアントを閉じて、サーバー側でリソースリークが発生しないようにしてください。
const changeStream = db.collection('events').watch();
// Graceful shutdown
process.on('SIGTERM', async () => {
console.log('Shutting down...');
await changeStream.close();
await client.close();
process.exit(0);
});
// Handle SIGINT (Ctrl+C in development)
process.on('SIGINT', async () => {
await changeStream.close();
await client.close();
process.exit(0);
});チェンジストリームとポーリングの比較
チェンジストリームが登場する前は、アプリケーションはポーリングに頼り、MongoDBに定期的にクエリを実行して新規または更新されたドキュメントを検出していました。ポーリングでは、何も変更されていなくてもクエリが実行されるためリソースを無駄にし、検出の遅延がポーリング間隔と同じになるためレイテンシが発生し、イベント率が高い状況ではスケールしにくくなります。チェンジストリームを使用すれば、ポーリングを完全に排除できます。データが変更されるとアプリケーションに即座に通知され、コレクションが静かなときのリソース消費も最小限に抑えられます。そのため、データ変更への反応が必要な機能では、チェンジストリームが推奨されるパターンです。
理解度チェック
このレッスンで学んだMongoDBとNoSQLデータベースの概念について、理解度を確認しましょう。
レッスンのまとめ
このレッスンでは、チェンジストリームがコレクション、データベース、またはデプロイメントに対するすべてのCRUD操作のリアルタイムイベントフィードを提供すること、レプリカセットが必要で、非同期反復またはEventEmitterを介してイベントを消費すること、そしてfullDocumentオプションによって更新イベントに更新後のドキュメント全体を含められることを学びました。次は、変更イベントドキュメントの構造と、各操作種別の処理方法について説明します。
AI チューターと学ぶ JavaScript — 無料
ブラウザでリアルコードを書いて実行し、24/7 の AI チューターから瞬時にサポートを受け、ウェブまたはアプリで続きから学習できます。
- コース
- 30
- レッスン
- 120
よくある質問
「コレクションでの変更ストリームの開始」レッスンは無料ですか?
はい。「コレクションでの変更ストリームの開始」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、MongoDB Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 MongoDB Academyコースには全4レッスンが含まれています。
「コレクションでの変更ストリームの開始」で何を学びますか?
コレクションでwatch()を呼び出し、非同期反復を使ってNode.jsアプリケーションでイベントストリームを消費します。 ブラウザで直接実行するハンズオンコードでMongoDB Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
MongoDB Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのMongoDB Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン1/4です。
「コレクションでの変更ストリームの開始」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このMongoDB Academyレッスンでコードを書いて実行できますか?
はい。すべてのMongoDB Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- コレクションでの変更ストリームの開始
- 変更イベントドキュメントの構造
- 集約パイプラインによるイベントのフィルタリング
- 中断後の変更ストリームの再開