MongoDB: Change Streams:リアルタイム変更監視
Change StreamsはMongoDBのリアルタイム変更通知の基盤—CDC(Change Data Capture)とイベント駆動アプリケーションの中核技術です。
1. 学習内容
- Change Streamsの概念(oplogベース)
- db.collection.watch()リスニング
- パイプラインフィルタリング($match)
- fullDocument: 'updateLookup'
- Mongooseとの統合
- 代表的なユースケース(リアルタイム通知、データ同期)
graph LR
A[アプリケーション1<br/>注文挿入] -->|oplog| DB[(MongoDB<br/>oplog)]
B[アプリケーション2<br/>注文更新] -->|oplog| DB
C[アプリケーション3<br/>注文削除] -->|oplog| DB
DB -->|Change Stream| CS[db.collection.watch]
CS -->|operationType=insert| N1[メール通知<br/>WebSocketプッシュ]
CS -->|operationType=update| N2[Elasticsearch<br/>同期]
CS -->|operationType=delete| N3[Redisキャッシュ<br/>無効化]
style CS fill:#d4edda
style N1 fill:#cce5ff
2. Change Streamsの基本
概念概要: Change StreamsはMongoDB 3.6以降で導入されたリアルタイム変更監視APIで、レプリカセットのoplogを使用して実装されています。ポーリングなしでデータ変更をリアルタイムに検出でき、CDC(Change Data Capture)とイベント駆動アーキテクチャの基盤となります。
OpLogの原理: Oplog(Operation Log)は、レプリカセット内の全ての書き込み操作を記録する特別な固定サイズコレクションです。Primaryは各書き込み操作(insert/update/delete)をOplogに追記し、SecondaryはOplogを追跡(tail)してデータを複製します。Change Streamsは本質的にOplogの構造化されたカプセル化です—生のOplogエントリを読み取り可能な変更イベントに変換し、フィルタリングやレジューム機能などの高度な機能を提供します。
graph TB
subgraph "OpLogの仕組み"
W1[書き込み操作] --> OP[(oplog<br/>固定コレクション)]
OP --> S1[Secondary 1<br/>oplog追跡中]
OP --> S2[Secondary 2<br/>oplog追跡中]
OP --> CS[Change Stream<br/>構造化イベント]
end
subgraph "Change Stream vs ポーリング"
P1[ポーリング: 毎秒クエリ] --> C1[高遅延<br/>高負荷<br/>重複クエリ]
CS --> C2[リアルタイム通知<br/>低負荷<br/>重複なし]
end
style CS fill:#d4edda
style C2 fill:#d4edda
style C1 fill:#f8d7da
Change Streams vs ポーリングの比較:
| 項目 | ポーリング | Change Streams |
|---|---|---|
| 遅延 | 1–60秒(間隔による) | < 100 ms(リアルタイムプッシュ) |
| DB負荷 | 高い(毎回インデックススキャン) | 低い(oplogイベントを受動的に受信) |
| イベント完全性 | イベント消失の可能性(間隔内の複数変更) | 消失なし(oplogは順序保証) |
| レジューム | 手動実装必要 | 組み込みresumeToken |
| リソース使用 | 継続的な接続使用量 + CPU | 長時間接続、低CPU |
変更イベントの構造:
// === コレクション変更を監視 ===
const changeStream = db.products.watch();
changeStream.on('change', (change) => {
console.log('変更を検出:', change);
});
出力:
TEXT 📖 参照専用{ _id: { _data: '...' }, // resumeToken(レジューム用) operationType: 'insert', // 操作タイプ fullDocument: { _id: ..., sku: ..., title: ..., ... }, // 完全なドキュメント ns: { db: 'shopdb', coll: 'products' }, // 名前空間 documentKey: { _id: ObjectId('...') } // ドキュメント主キー }
| operationType | 意味 | fullDocument |
|---|---|---|
insert |
新規ドキュメント挿入 | ✅ 完全なドキュメント |
update |
ドキュメント更新 | ❌ 変更フィールドのみ(updateLookup必要) |
replace |
ドキュメント置換 | ✅ 完全なドキュメント |
delete |
ドキュメント削除 | ❌ なし(ドキュメント削除済み) |
drop |
コレクション削除 | ❌ なし |
rename |
コレクション名変更 | ❌ なし |
invalidate |
無効化(コレクション削除など) | ❌ なし |
ポイント解説:
- Change Streamsはレプリカセットまたはシャードクラスターで実行する必要があります(oplogに依存)。
updateイベントはデフォルトで完全なドキュメントを返しません。fullDocument: 'updateLookup'を設定する必要があります。_idフィールドはresumeTokenで、レジュームに使用します。保存しておけば、再起動後も続きから監視できます。
3. パイプラインフィルタリング
概念説明: Change Streamsは集計パイプラインによるフィルタリングをサポート—$matchや$projectなどのステージをoplogイベントストリームに適用し、関心のあるイベントのみをアプリケーションに渡します。フィルタリングはサーバー側で実行され、ネットワークトラフィックとアプリケーション層の処理負荷を削減します。
動作原理: watch()は集計パイプラインの配列をパラメータとして受け取ります。パイプラインステージはMongoDBサーバー側で実行されます。イベントはまずパイプラインでフィルタリングされ、通過したものだけがクライアントにプッシュされます。$match、$project、$addFields、$replaceRootなどのステージをサポートしています。
フィルタリングポリシー:
| フィルタ対象 | パイプラインステージ | 例 |
|---|---|---|
| 操作タイプ | $match: {operationType} |
挿入のみ監視 |
| ドキュメントフィールド | $match: {'fullDocument.field'} |
特定カテゴリのみ |
| 変更フィールド | $match: {'updateDescription.updatedFields'} |
価格変更のみ |
| 複合条件 | $match: {$or: [...]} |
挿入または価格変更 |
// === 特定の操作でフィルタリング ===
const changeStream = db.products.watch([
{ $match: { operationType: 'insert' } }
]);
// === 特定のフィールドでフィルタリング ===
const changeStream = db.products.watch([
{ $match: { 'fullDocument.category': 'Electronics' } }
]);
// === 価格変更でフィルタリング ===
const changeStream = db.products.watch([
{
$match: {
$or: [
{ operationType: 'insert' },
{ operationType: 'update', 'updateDescription.updatedFields.price': { $exists: true } }
]
}
}
]);
ポイント解説:
- フィルタリングはサーバー側で実行され、不要なネットワークトラフィックを削減します。
$matchのフィールドパスはoplogイベント構造(fullDocument.categoryなど)を使用し、元のドキュメントフィールドではありません。$groupや$limitなどグローバル状態を必要とするステージはサポートされません。
4. fullDocument設定
概念説明: fullDocumentオプションは、Change Streamイベントが完全なドキュメントを含むかどうかを制御します。デフォルトでは、updateイベントは変更されたフィールドのみ(updateDescription)を返し、完全なドキュメントは返しません。fullDocument: 'updateLookup'を設定すると、MongoDBは追加でクエリを実行して現在のドキュメントの完全な内容を取得します。
動作原理:
whenAvailable(デフォルト): insert/replaceイベントは完全なドキュメントを含む、updateイベントは変更フィールドのみ、deleteイベントはドキュメントを含まないupdateLookup: 全イベントが完全なドキュメントを含む(updateイベントは追加で最新ドキュメントをクエリ)、ただし追加の読み取りオーバーヘッドと若干の遅延ありrequired: 完全なドキュメントを返す必要あり、そうでなければエラー
updateLookupのトレードオフ:
| 項目 | whenAvailable | updateLookup |
|---|---|---|
| 完全なドキュメント | 挿入/置換のみ | 全操作タイプ |
| パフォーマンス | ベースライン | 追加クエリオーバーヘッド(+10–20%) |
| データ鮮度 | 変更時点 | クエリ時点(若干の遅延可能性) |
| 使用シーン | 変更フィールドのみ知りたい | 完全なドキュメントが必要な後続処理 |
sequenceDiagram
participant App as アプリケーション
participant DB as MongoDB
participant Doc as ドキュメント
App->>DB: watch([], {fullDocument: 'updateLookup'})
DB->>DB: oplogがupdateイベント生成
Note over DB: updateイベントはデフォルトで変更フィールドのみ含む
DB->>Doc: 現在の完全なドキュメントを検索
Doc-->>DB: 最新ドキュメントを返す
DB-->>App: イベント + fullDocumentをプッシュ
Note over App: イベントが完全なドキュメントを含む
// === updateで完全なドキュメントを返す ===
const changeStream = db.products.watch([], {
fullDocument: 'updateLookup'
});
changeStream.on('change', (change) => {
if (change.operationType === 'update') {
console.log('更新されたドキュメント:', change.fullDocument);
// 完全なドキュメント(デフォルトは変更後の値)
}
});
// === 変更フィールドのみを返す ===
const changeStream = db.products.watch([], {
fullDocument: 'whenAvailable' // デフォルト
});
ポイント解説:
updateLookupは追加のクエリを実行するため、頻繁な更新があるシナリオではDB負荷が増加します。updateLookupはクエリ時点での最新ドキュメントを返すため、変更時点のドキュメントとは若干異なる可能性があります(他の同時変更による)。deleteイベントはupdateLookupを設定してもドキュメントを返しません(ドキュメントが既に存在しないため)
5. Mongoose統合
概念説明: Mongoose 6+はChange Streamsをネイティブにサポートし、Model.watch()で返されます。これはネイティブMongoDBドライバのAPIと完全に一致しますが、モデル層で直接使用することでMongooseの開発規約に合致します。
イベント駆動アーキテクチャ: Change Streamsはイベント駆動アーキテクチャの中核コンポーネントです—データベース変更がイベントソースとなり、キャッシュ更新、検索同期、通知プッシュなどのダウンストリーム副作用を駆動します。
graph TB
subgraph "イベントソース"
DB[(MongoDB<br/>oplog)]
end
subgraph "Change Streamバス"
CS[Model.watch<br/>変更フロー]
end
subgraph "イベントコンシューマー"
N1[メール通知]
N2[Elasticsearch<br/>検索同期]
N3[Redis<br/>キャッシュ無効化]
N4[WebSocket<br/>リアルタイム通知]
N5[監査ログ<br/>コンプライアンス記録]
end
DB --> CS
CS --> N1
CS --> N2
CS --> N3
CS --> N4
CS --> N5
style CS fill:#d4edda
// === Mongoose Change Streams(Mongoose 6+)===
const Product = mongoose.model('Product', productSchema);
// 商品コレクションの変更を監視
const changeStream = Product.watch();
changeStream.on('change', (change) => {
console.log(`${change.operationType}:`, change.fullDocument);
});
// フィルタリング
const filteredStream = Product.watch([
{ $match: { operationType: { $in: ['insert', 'update'] } } }
]);
ポイント解説:
- Mongooseの
watch()はネイティブMongoDB Change Streamを返し、APIは完全に同一です。 - 監視中はDB接続を維持する必要があります。接続が失われるとChange Streamは自動的に停止します。
- 本番環境では、Change Streamマネージャーでラップすることを推奨:自動再接続 + resumeToken永続化
6. 実践的なユースケース
概念概要: Change Streamsには本番環境での3つの典型的なユースケースがあります:リアルタイム通知、データ同期(CDC)、キャッシュ無効化。それぞれがイベント駆動アーキテクチャの典型的な実装です。
(1) リアルタイム通知システム
シナリオ: ShopHub ECでは、新規注文が作成された際、商家の管理画面にリアルタイムでメール通知とWebSocketプッシュを送信する必要があります。
// === 新規注文を監視、通知を送信 ===
const OrderStream = db.orders.watch([
{ $match: { operationType: 'insert' } }
]);
OrderStream.on('change', async (change) => {
const order = change.fullDocument;
// メール送信
await sendEmail(order.userId, '注文確認', `注文 ${order._id} を受領しました`);
// プッシュ通知
await pushNotification(order.userId, {
title: '新規注文',
body: `注文合計: $${order.total}`
});
// WebSocketリアルタイム通知
io.emit('new_order', order);
});
(2) データ同期(CDC)
シナリオ: ShopHubはMongoDBの商品データをリアルタイムでElasticsearchに同期し、全文検索をサポートする必要があります。Change StreamsでゼロレイテンシーのCDCパイプラインを実現します。
CDCアーキテクチャ:
graph LR
A[MongoDB<br/>商品データ] -->|Change Stream| B[CDCワーカー<br/>Node.jsプロセス]
B -->|insert/update| C[Elasticsearch<br/>検索インデックス]
B -->|delete| C
B -->|変更ログ| D[(Redis<br/>resumeToken)]
D -->|再起動復旧| B
style B fill:#d4edda
// === MongoDB変更を監視、Elasticsearchに同期 ===
const ProductStream = db.products.watch();
ProductStream.on('change', async (change) => {
switch (change.operationType) {
case 'insert':
case 'update':
case 'replace':
await elasticsearch.index({
index: 'products',
id: change.documentKey._id.toString(),
body: change.fullDocument
});
break;
case 'delete':
await elasticsearch.delete({
index: 'products',
id: change.documentKey._id.toString()
});
break;
}
});
(3) キャッシュ無効化
シナリオ: ShopHubはRedisで商品詳細ページをキャッシュしています。商品データが変更されたら、自動的にキャッシュをクリアし、古いデータを防止します。
// === 商品変更を監視、Redisキャッシュを無効化 ===
const ProductStream = db.products.watch();
ProductStream.on('change', async (change) => {
const productId = change.documentKey._id.toString();
await redis.del(`product:${productId}`);
console.log(`${productId}のキャッシュをクリア`);
});
▶ サンプル 1:Change Stream中断後のレジューム
// AliceのTechCorpシステム:Change Streamレジューム、再起動後もイベント消失なし
async function startResumableStream() {
// 1. Redisから最後に保存したresumeTokenを取得
let resumeToken = await redis.get('product_stream_token');
let options = { fullDocument: 'updateLookup' };
if (resumeToken) {
options.resumeAfter = JSON.parse(resumeToken);
console.log('保存されたトークンから再開');
}
// 2. 監視開始
const changeStream = db.products.watch([], options);
changeStream.on('change', async (change) => {
// 変更を処理...
console.log(`${change.operationType}: ${change.documentKey._id}`);
// 3. 各イベント後にresumeTokenを保存
await redis.set('product_stream_token', JSON.stringify(change._id));
});
changeStream.on('error', async (err) => {
console.error('ストリームエラー:', err.message);
// 4. エラー後、遅延再接続
setTimeout(startResumableStream, 5000);
});
}
startResumableStream();
7. Change Streamsの制限
概念説明: Change Streamsはoplogとレプリカセットに依存し、明確な制限があります。信頼性の高いリアルタイムシステムを設計するには、これらの制限を理解することが不可欠です。
制限の詳細:
| 制約 | 説明 | 理由 | 軽減策 |
|---|---|---|---|
| レプリカセット必須 | スタンドアロンモードでは未対応 | oplogに依存 | 開発環境では単一ノードレプリカセットを使用 |
| oplogサイズ制限 | デフォルト:ディスク容量の5% | oplogは固定コレクション | oplogSizeを増やすかウィンドウを監視 |
| クラスター横断不可 | 単一クラスター内 | oplogはクラスターをまたがない | Kafkaで複数クラスターを接続 |
| $where未対応 | 一部演算子 | oplogはクエリ詳細を記録しない | fullDocumentフィールドでフィルタ |
| イベント順序 | 単一コレクション内は順序、コレクション間はグローバル順序なし | oplogはコレクション単位でシャード | クロスコレクションの並べ替えはアプリ層で実行 |
| メモリ使用量 | 監視接続ごとのメモリ使用 | サーバーがカーソル状態を維持 | 監視接続数を制限 |
oplogウィンドウの監視: oplogは固定サイズで、古いエントリは上書きされます。Change Streamの消費速度がoplogの上書き速度より遅いと、resumeTokenが期限切れとなり、レジューム機能が動作しません。
graph LR
A[oplog書き込み速度<br/>1000 ops/s] --> B[oplog容量<br/>ディスクの5% ≈ 50GB]
B --> C[oplogウィンドウ<br/>約72時間]
C --> D{消費速度?}
D -->|追従| E[✅ 正常]
D -->|72h以上遅延| F[❌ resumeToken無効<br/>全量再同期が必要]
style E fill:#d4edda
style F fill:#f8d7da
| 監視項目 | コマンド | アラート閾値 |
|---|---|---|
| oplogウィンドウ | rs.printReplicationInfo() |
< 24時間 |
| oplog使用率 | db.oplog.rs.stats() |
> 80% |
| Change Stream接続数 | db.currentOp() |
> 100 |
| イベント遅延 | アプリ層監視 | > 10秒 |
▶ サンプル:Change Streamsリアルタイム注文通知のハンズオンガイド
// シーン:新規注文を監視、リアルタイムでメール送信 + WebSocketプッシュ + Elasticsearch同期
// 1. 監視開始(Node.js)
const { MongoClient } = require('mongodb');
async function startOrderListener() {
const client = new MongoClient('mongodb://localhost:27017/?replicaSet=rs0');
await client.connect();
const orders = client.db('shopdb').collection('orders');
// 注文コレクションの全変更を監視
const changeStream = orders.watch([
{
$match: {
operationType: 'insert', // 挿入のみ監視
'fullDocument.status': 'paid' // 支払い済み注文のみ処理
}
}
]);
changeStream.on('change', async (change) => {
const order = change.fullDocument;
console.log(`新規注文: ${order._id}, 金額: $${order.total}`);
// 1. メール通知送信
await sendEmail(order.userId, {
subject: '注文確認',
body: `注文 ${order._id} を受領しました。合計金額 $${order.total}`
});
// 2. WebSocketリアルタイムプッシュで商家の管理画面に通知
io.to('merchant-dashboard').emit('new_order', {
orderId: order._id,
total: order.total,
items: order.items,
timestamp: order.createdAt
});
// 3. Elasticsearchに同期して検索用に
await elasticsearch.index({
index: 'orders',
id: order._id.toString(),
body: order
});
});
// 4. 商品変更を監視、Elasticsearchを同期更新
const productStream = client.db('shopdb').collection('products').watch();
productStream.on('change', async (change) => {
if (change.operationType === 'delete') {
await elasticsearch.delete({
index: 'products',
id: change.documentKey._id.toString()
});
} else if (change.fullDocument) {
await elasticsearch.index({
index: 'products',
id: change.documentKey._id.toString(),
body: change.fullDocument
});
}
});
console.log('Change Streams監視中...');
}
// 5. レジューム(アプリ再起動後の進捗追跡復旧)
async function resumeAfterRestart() {
const resumeToken = await redis.get('change_stream_resume_token');
const changeStream = orders.watch([], {
resumeAfter: JSON.parse(resumeToken),
fullDocument: 'updateLookup'
});
changeStream.on('change', async (change) => {
// 変更を処理...
// resume tokenを保存
await redis.set('change_stream_resume_token', JSON.stringify(change._id));
});
}
startOrderListener().catch(console.error);
// 2. mongoshで手動テスト
// イベントトリガー:新規注文を挿入
db.orders.insertOne({
userId: 'user_001',
items: [{ sku: 'PHONE-001', qty: 1, price: 599 }],
total: 599,
status: 'paid',
createdAt: new Date()
});
// アプリが即座に通知を受信:メール送信 + WebSocketプッシュ + ESインデックス
出力:
TEXT 📖 参照専用注文挿入後100ms以内にトリガーされ、ポーリングなしで全ての副作用(メール、プッシュ通知、検索同期)を自動的に完了。
▶ サンプル 3:Change Streamsでリアルタイム分析ダッシュボード(難易度 ⭐⭐)
// シーン:ShopHub管理ダッシュボードでライブ注文統計を表示
const mongoose = require('mongoose');
async function startOrderAnalytics() {
const Order = mongoose.model('Order');
// インメモリで日次注文数と売上を追跡(ダッシュボード用)
let stats = {
date: new Date().toISOString().split('T')[0],
orderCount: 0,
totalRevenue: 0,
lastOrderId: null
};
// 挿入操作のみ監視(新規注文)
const changeStream = Order.watch([
{ $match: { operationType: 'insert' } }
], { fullDocument: true });
changeStream.on('change', async (change) => {
const order = change.fullDocument;
// インメモリ統計を更新
stats.orderCount++;
stats.totalRevenue += order.total || 0;
stats.lastOrderId = order._id;
// WebSocketクライアントにブロードキャスト(ダッシュボード)
broadcastToDashboard({
event: 'new_order',
data: {
orderId: order._id,
total: order.total,
stats: stats
}
});
// デモ用にコンソールに出力
console.log(`[LIVE] 注文 #${order.orderNumber} - $${order.total}`);
console.log(` 本日: ${stats.orderCount}件の注文、$${stats.totalRevenue}の売上`);
});
changeStream.on('error', (err) => {
console.error('Change Streamエラー:', err);
// 再接続ロジックをここに追加
});
console.log('注文分析ストリームを開始...');
}
// WebSocketブロードキャストをシミュレート
function broadcastToDashboard(data) {
// 本番環境: io.emit('order_update', data)
}
startOrderAnalytics();
出力:
TEXT 📖 参照専用注文分析ストリームを開始... [LIVE] 注文 #ORD-2026-001 - $599 本日: 1件の注文、$599の売上 [LIVE] 注文 #ORD-2026-002 - $1299 本日: 2件の注文、$1898の売上
❓ よくある質問
resumeAfterで続きから再開できます。resumeTokenを外部(Redisなど)に保存し、再起動時に復元できます。📖 まとめ
- Change Streams:oplogベースのリアルタイム変更通知
- watch()メソッドでリスニング + $matchフィルタリング
- operationType: insert / update / delete / replace / drop
- fullDocument: 'updateLookup'で完全なドキュメントを取得
- 代表的なシーン:リアルタイム通知、CDC、キャッシュ無効化
- レプリカセット必須、oplogサイズ制限
📝 練習問題
- 基礎問題(⭐):
productsコレクションの全変更を監視し、コンソールに出力。 - 基礎問題(⭐): $matchを使用して挿入操作のみフィルタリング。
- 応用問題(⭐⭐): 新規注文をリッスンし、メール通知を送信(モックメール関数を使用)。
- 応用問題(⭐⭐): 商品キャッシュ無効化を実装(Redis同期)。
- チャレンジ(⭐⭐⭐): 完全なCDCシステムを実装(MongoDBとElasticsearch間のリアルタイム同期)。