MongoDB: Change Streams:リアルタイム変更監視

Change StreamsはMongoDBのリアルタイム変更通知の基盤—CDC(Change Data Capture)とイベント駆動アプリケーションの中核技術です。

1. 学習内容


100%
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エントリを読み取り可能な変更イベントに変換し、フィルタリングやレジューム機能などの高度な機能を提供します。

100%
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

変更イベントの構造:

JAVASCRIPT
// === コレクション変更を監視 ===
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 無効化(コレクション削除など) ❌ なし

ポイント解説:

  1. Change Streamsはレプリカセットまたはシャードクラスターで実行する必要があります(oplogに依存)。
  2. updateイベントはデフォルトで完全なドキュメントを返しません。fullDocument: 'updateLookup'を設定する必要があります。
  3. _idフィールドはresumeTokenで、レジュームに使用します。保存しておけば、再起動後も続きから監視できます。


3. パイプラインフィルタリング

概念説明: Change Streamsは集計パイプラインによるフィルタリングをサポート—$match$projectなどのステージをoplogイベントストリームに適用し、関心のあるイベントのみをアプリケーションに渡します。フィルタリングはサーバー側で実行され、ネットワークトラフィックとアプリケーション層の処理負荷を削減します。

動作原理: watch()は集計パイプラインの配列をパラメータとして受け取ります。パイプラインステージはMongoDBサーバー側で実行されます。イベントはまずパイプラインでフィルタリングされ、通過したものだけがクライアントにプッシュされます。$match$project$addFields$replaceRootなどのステージをサポートしています。

フィルタリングポリシー:

フィルタ対象 パイプラインステージ
操作タイプ $match: {operationType} 挿入のみ監視
ドキュメントフィールド $match: {'fullDocument.field'} 特定カテゴリのみ
変更フィールド $match: {'updateDescription.updatedFields'} 価格変更のみ
複合条件 $match: {$or: [...]} 挿入または価格変更
JAVASCRIPT
// === 特定の操作でフィルタリング ===
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 } }
      ]
    }
  }
]);

ポイント解説:

  1. フィルタリングはサーバー側で実行され、不要なネットワークトラフィックを削減します。
  2. $matchのフィールドパスはoplogイベント構造(fullDocument.categoryなど)を使用し、元のドキュメントフィールドではありません。
  3. $group$limitなどグローバル状態を必要とするステージはサポートされません。


4. fullDocument設定

概念説明: fullDocumentオプションは、Change Streamイベントが完全なドキュメントを含むかどうかを制御します。デフォルトでは、updateイベントは変更されたフィールドのみ(updateDescription)を返し、完全なドキュメントは返しません。fullDocument: 'updateLookup'を設定すると、MongoDBは追加でクエリを実行して現在のドキュメントの完全な内容を取得します。

動作原理:

updateLookupのトレードオフ:

項目 whenAvailable updateLookup
完全なドキュメント 挿入/置換のみ 全操作タイプ
パフォーマンス ベースライン 追加クエリオーバーヘッド(+10–20%)
データ鮮度 変更時点 クエリ時点(若干の遅延可能性)
使用シーン 変更フィールドのみ知りたい 完全なドキュメントが必要な後続処理
100%
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: イベントが完全なドキュメントを含む
JAVASCRIPT
// === 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'  // デフォルト
});

ポイント解説:

  1. updateLookupは追加のクエリを実行するため、頻繁な更新があるシナリオではDB負荷が増加します。
  2. updateLookupはクエリ時点での最新ドキュメントを返すため、変更時点のドキュメントとは若干異なる可能性があります(他の同時変更による)。
  3. deleteイベントはupdateLookupを設定してもドキュメントを返しません(ドキュメントが既に存在しないため)


5. Mongoose統合

概念説明: Mongoose 6+はChange Streamsをネイティブにサポートし、Model.watch()で返されます。これはネイティブMongoDBドライバのAPIと完全に一致しますが、モデル層で直接使用することでMongooseの開発規約に合致します。

イベント駆動アーキテクチャ: Change Streamsはイベント駆動アーキテクチャの中核コンポーネントです—データベース変更がイベントソースとなり、キャッシュ更新、検索同期、通知プッシュなどのダウンストリーム副作用を駆動します。

100%
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
JAVASCRIPT
// === 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'] } } }
]);

ポイント解説:

  1. Mongooseのwatch()はネイティブMongoDB Change Streamを返し、APIは完全に同一です。
  2. 監視中はDB接続を維持する必要があります。接続が失われるとChange Streamは自動的に停止します。
  3. 本番環境では、Change Streamマネージャーでラップすることを推奨:自動再接続 + resumeToken永続化


6. 実践的なユースケース

概念概要: Change Streamsには本番環境での3つの典型的なユースケースがあります:リアルタイム通知、データ同期(CDC)、キャッシュ無効化。それぞれがイベント駆動アーキテクチャの典型的な実装です。

(1) リアルタイム通知システム

シナリオ: ShopHub ECでは、新規注文が作成された際、商家の管理画面にリアルタイムでメール通知とWebSocketプッシュを送信する必要があります。

JAVASCRIPT
// === 新規注文を監視、通知を送信 ===
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アーキテクチャ:

100%
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
JAVASCRIPT
// === 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で商品詳細ページをキャッシュしています。商品データが変更されたら、自動的にキャッシュをクリアし、古いデータを防止します。

JAVASCRIPT
// === 商品変更を監視、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中断後のレジューム

JAVASCRIPT
// 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が期限切れとなり、レジューム機能が動作しません。

100%
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リアルタイム注文通知のハンズオンガイド

JAVASCRIPT
// シーン:新規注文を監視、リアルタイムでメール送信 + 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でリアルタイム分析ダッシュボード(難易度 ⭐⭐)

JAVASCRIPT
// シーン: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の売上

❓ よくある質問

Q Change Streamsはリアルタイムですか?
A ほぼリアルタイムです。Change Streamsはoplogエントリ書き込み直後にトリガーされ、遅延は100ms未満です。
Q Change Streamsでイベントは消失しますか?
A いいえ(oplogが上書きされない限り)。resumeAfterで続きから再開できます。
Q 監視の進捗を永続化するには?
A resumeTokenを外部(Redisなど)に保存し、再起動時に復元できます。

📖 まとめ


📝 練習問題

  1. 基礎問題(⭐): productsコレクションの全変更を監視し、コンソールに出力。
  2. 基礎問題(⭐): $matchを使用して挿入操作のみフィルタリング。
  3. 応用問題(⭐⭐): 新規注文をリッスンし、メール通知を送信(モックメール関数を使用)。
  4. 応用問題(⭐⭐): 商品キャッシュ無効化を実装(Redis同期)。
  5. チャレンジ(⭐⭐⭐): 完全なCDCシステムを実装(MongoDBとElasticsearch間のリアルタイム同期)。
Web-Tutorial.com

Web-Tutorial 技術チーム

複数の開発者によって共同維持されているプログラミングチュートリアルプラットフォーム。各チュートリアルは専門分野の開発者が執筆・レビューしています。正確で信頼性の高いコンテンツを目指しています — 問題を見つけた場合はお知らせください。

100%