MongoDB Academy · レッスン

中断後の変更ストリームの再開

再開トークンを永続化し、最後に処理したイベントから変更ストリームを再開して、少なくとも1回の配信を保証します。

レッスン 4/413 ステップ

「中断後の変更ストリームの再開」はCoddyKit上の無料MongoDB Academyレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはMongoDB Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 MongoDB Academyコースには全4レッスンが含まれています。

問題:クラッシュ後はどうなるか

変更ストリームのコンシューマプロセスは、ネットワーク障害、デプロイ、クラッシュ、または意図的な再起動によって中断されることがあります。最後に処理したイベントから再開する仕組みがなければ、アプリケーションは現在の時点からストリームを再開するため、停止中に発生したすべてのイベントを取りこぼしてしまいます。MongoDBのresume tokenを使うと、oplogの履歴内にある正確な位置からストリームを再開できるため、この問題を解決できます。

Resume Tokenとは

すべての変更イベントドキュメントには、resume tokenとして機能する_idフィールドがあります。これは、oplog内でのイベントの位置を一意に識別する、バイナリ形式の不透明な値です。{ _data: '...' }のような形式で、dataには16進数でエンコードされた内部識別子が格納されています。内容を理解する必要はありません。保存してMongoDBに渡し、正確にその位置から再開できれば十分です。

// A resume token is the _id field of a change event
// Example structure (your values will differ):
// {
//   _id: {
//     _data: '8264CE3B7200000002463C6F5A100...'
//   }
// }

// Store it as-is — do not parse or modify the _data value

Resume Tokenの永続化

各イベントの処理後、処理済みであることを確認する前に、resume tokenを保存先の永続ストア(データベース、ファイル、Redisなど)に保存してください。この「チェックポイント」パターンにより、再起動後も最後に正常に処理できたイベントから処理を再開できます。処理後に保存し、処理前には保存しないでください。処理中にプロセッサがクラッシュした場合に、イベントをスキップするのを防ぐためです。

const CHECKPOINT_COLLECTION = 'changeStreamCheckpoints';
const STREAM_ID = 'orderProcessorV1';

for await (const change of changeStream) {
  // Process the event
  await processOrderChange(change);

  // Save token AFTER successful processing
  await db.collection(CHECKPOINT_COLLECTION).updateOne(
    { streamId: STREAM_ID },
    { $set: { resumeToken: change._id, processedAt: new Date() } },
    { upsert: true }
  );
}

resumeAfterによる再開

保存したトークンから再開するには、resumeAfterオプションを使用してwatch()に渡します。MongoDBは、トークンで識別されたイベントの後からイベントを再配信します。最後に処理したイベントが再配信されることはありません。トークンがまだoplogに残っているイベントを指している場合、MongoDBは次のイベントから配信を開始します。これにより、べき等なイベント処理と組み合わせた場合にExactly-once配信を実現できます。

async function startOrResumeStream() {
  // Load the last saved resume token
  const checkpoint = await db.collection('changeStreamCheckpoints')
    .findOne({ streamId: 'orderProcessorV1' });

  const watchOptions = checkpoint
    ? { resumeAfter: checkpoint.resumeToken }
    : {}; // Start from current position if no token

  const changeStream = db.collection('orders').watch([], watchOptions);

  for await (const change of changeStream) {
    await processOrderChange(change);
    await saveResumeToken(change._id);
  }
}

startAfterとresumeAfterの違い

MongoDBには、resumeAfterとstartAfterという2つの似たオプションがあります。resumeAfterは通常の操作トークンから再開しますが、'invalidate'イベントのトークンが渡されるとエラーになります。startAfterも同様に動作しますが、invalidateイベントの後から再開することもできます。これは、コレクションが削除されて再作成された後にストリームを再開したい場合に便利です。ほとんどの用途では、resumeAfterが適切な選択です。

// resumeAfter — standard use case, does not work after invalidate tokens
db.collection('orders').watch([], { resumeAfter: savedToken });

// startAfter — can also start after an invalidate event
// (use when the collection may have been dropped and recreated)
db.collection('orders').watch([], { startAfter: savedToken });

時間ベースの再開に使うstartAtOperationTime

resume token がないものの、開始したい時点の timestampがわかっている場合は、MongoDB の Timestamp オブジェクトとともに startAtOperationTime を使用できます。これにより、特定のイベントではなく、特定のクラスタ時刻からストリームを開始できます。これは、デプロイの開始以降に発生したすべてのイベントを再処理するなど、イベントの再生が必要な場合に便利です。ただし、その時刻までの履歴が oplog に残っている必要があります。

const { Timestamp } = require('mongodb');

// Replay events since a specific time
const startTime = new Timestamp({ t: Math.floor(Date.now() / 1000) - 3600, i: 1 }); // 1 hour ago

const changeStream = db.collection('orders').watch([], {
  startAtOperationTime: startTime
});

// Events from 1 hour ago will be delivered

Oplog の保持と有効期間

Change stream の再開が機能するのは、token に対応するイベントがまだ oplog に存在する場合だけです。MongoDB の oplog はサイズに上限のある capped collection であり、新しいエントリが追加されると古いエントリが上書きされます。Atlas では、oplog の保持期間を設定できます(通常は 24~72 時間です)。アプリケーションが oplog の保持期間を超えて停止していた場合、token は期限切れの位置を指すことになり、MongoDB は ChangeStreamHistoryLost エラー(コード 286)を返します。この場合、アプリケーションは新しく開始して対処する必要があります。

async function startOrResumeWithFallback() {
  const checkpoint = await loadResumeToken();
  try {
    const options = checkpoint ? { resumeAfter: checkpoint } : {};
    const stream = db.collection('orders').watch([], options);
    for await (const change of stream) {
      await processChange(change);
      await saveResumeToken(change._id);
    }
  } catch (error) {
    if (error.code === 286) { // ChangeStreamHistoryLost
      console.warn('Resume token expired — starting from now and running catch-up scan');
      await deleteResumeToken();
      await catchUpScan(); // full collection scan to catch missed changes
      await startOrResumeWithFallback();
    } else {
      throw error;
    }
  }
}

冪等なイベント処理

Change stream は少なくとも 1 回の配信を保証するため(resume 後に同じイベントが複数回配信される可能性があります)、イベントハンドラーは冪等でなければなりません。つまり、同じイベントを 2 回処理しても、1 回処理した場合と同じ結果になる必要があります。方法としては、変更を適用する前にドキュメントがすでに変更後の状態になっていないか確認する、insert の代わりに upsert を使用する、処理済みイベント ID を「processed events」セットに保存して重複をスキップする、といったものがあります。

async function processOrderChange(change) {
  if (change.operationType === 'insert') {
    const { fullDocument } = change;
    // Idempotent upsert — safe to replay:
    await db.collection('orderSummaries').updateOne(
      { _id: fullDocument._id },       // match by _id
      { $set: { ...summarize(fullDocument) } }, // idempotent set
      { upsert: true }                 // create if not exists
    );
  }
}

ドライバーレベルでの自動再開

MongoDB Node.js driver(およびその他の公式ドライバー)には、一時的なネットワークエラーに対する自動再開機能が含まれています。MongoDB サーバーとの接続が一時的に失われると、ドライバーはコードで処理しなくても、最後に受信したイベントの token から change stream を透過的に再開します。この自動再開により、一般的な中断の多くに対応できます。手動で再開するロジックが必要なのは、プロセスのクラッシュやデプロイなど、アプリケーションレベルで再起動する場合だけです。

// The driver automatically resumes after network blips
// No code needed in your loop for this case:
for await (const change of changeStream) {
  // If the network drops and reconnects, the driver resumes automatically
  // and continues delivering events from where it left off.
  await processChange(change);
  await saveCheckpoint(change._id); // Still save tokens for process restarts
}

本番環境対応の Change Stream ハンドラー

本番環境の change stream コンシューマーでは、resume token の永続化(プロセスの再起動に対応)、oplog の期限切れへの対応(全件スキャンへのフォールバック)、冪等な処理(イベントを安全に再生)、グレースフルシャットダウン(SIGTERM 受信時にストリームを閉じる)を組み合わせる必要があります。この 4 つを組み合わせることで、ネットワーク障害、デプロイ、長時間の停止が発生しても、信頼性の高いイベント処理を実現できます。

class ReliableChangeStreamConsumer {
  constructor(collection, handler) {
    this.collection = collection;
    this.handler = handler;
    this.running = false;
  }

  async start() {
    this.running = true;
    while (this.running) {
      const token = await loadToken();
      const stream = this.collection.watch([], token ? { resumeAfter: token } : {});
      try {
        for await (const change of stream) {
          await this.handler(change);
          await saveToken(change._id);
        }
      } catch (e) {
        if (e.code === 286) { await deleteToken(); continue; } // expired, restart fresh
        if (!this.running) break; // shutting down
        throw e;
      } finally {
        await stream.close();
      }
    }
  }

  stop() { this.running = false; }
}

グレースフルシャットダウンと token のフラッシュ

計画的にシャットダウンする場合(デプロイや再起動など)は、change stream を閉じる前に必ず最新の resume token をフラッシュしてください。アプリケーションがイベントをバッチ単位で処理している場合(パフォーマンス向上のためにバッファリングしている場合)は、個々のイベントごとではなく、各バッチの処理後に token が保存されるようにします。保存するのは最後に受信したイベントの token ではなく、最後に正常に処理されたイベントの token です。この違いは重要です。イベントを受信してその token を保存した後、処理前にクラッシュすると、再開時にそのイベントが失われてしまうためです。

// Correct: save token AFTER processing, not before
for await (const change of changeStream) {
  await processChange(change);        // process first
  await saveToken(change._id);        // token saved only after success
  // If crash happens here, the event was already processed — safe
}

// Wrong: saving token before processing
for await (const change of changeStream) {
  await saveToken(change._id);        // token saved
  await processChange(change);        // if crash here, event is skipped on resume!
}

理解度チェック

このレッスンで学んだ MongoDB & NoSQL Databases の概念について理解度を確認しましょう。

レッスンのまとめ

このレッスンでは、次のことを学びました。resume token は change._id に保存され、各イベントの処理後に永続化する必要があること、最後に処理したイベントからストリームを再開するには、token を resumeAfter として watch() に渡すこと、そしてtoken の期限が切れて oplog から消えた場合は、現在時刻から再開してキャッチアップスキャンを実行することで ChangeStreamHistoryLost(コード 286)に対応することです。次は、全文検索用の Atlas Search インデックスの作成について学びます。

無料で開始

AI チューターと学ぶ JavaScript — 無料

ブラウザでリアルコードを書いて実行し、24/7 の AI チューターから瞬時にサポートを受け、ウェブまたはアプリで続きから学習できます。

コース
30
レッスン
120

よくある質問

「中断後の変更ストリームの再開」レッスンは無料ですか?

はい。「中断後の変更ストリームの再開」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、MongoDB Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 MongoDB Academyコースには全4レッスンが含まれています。

「中断後の変更ストリームの再開」で何を学びますか?

再開トークンを永続化し、最後に処理したイベントから変更ストリームを再開して、少なくとも1回の配信を保証します。 ブラウザで直接実行するハンズオンコードでMongoDB Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

MongoDB Academyを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのMongoDB Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。

「中断後の変更ストリームの再開」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このMongoDB Academyレッスンでコードを書いて実行できますか?

はい。すべてのMongoDB Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. コレクションでの変更ストリームの開始
  2. 変更イベントドキュメントの構造
  3. 集約パイプラインによるイベントのフィルタリング
  4. 中断後の変更ストリームの再開
← MongoDB Academyに戻る