CodeGym /コース /ChatGPT Apps /非同期ジョブ: キュー、ワーカー、再送(retry)

非同期ジョブ: キュー、ワーカー、再送(retry)

ChatGPT Apps
レベル 13 , レッスン 3
使用可能

1. なぜ ChatGPT App に非同期ジョブが必要か

世界が理想的なら、どんな MCP ツールでも数百ミリ秒で終わるでしょう。ですが現実で面白い処理はたいてい時間も負荷も重いものです。

  • ユーザーの購買履歴を含む巨大な CSV の解析;
  • 複数の外部 API からの集約(各 API は時にスリープし、時に 503 を返す);
  • 多数の中間ステップを含む高度なレコメンデーション生成;
  • 大きなレポートやプレゼン資料の生成。

これらを1つの同期 tool‑call に押し込もうとすると、3つの問題につき当たります。

まずはタイムアウトです。ChatGPT セッション、HTTP インフラ、MCP クライアント——いずれも「5分後に返答」には向きません。接続を長く保持しすぎるサーバーは、ChatGPT にもユーザーにも「固まっている」ように見えてしまいます。

次に、負荷制御です。100人のユーザーが同時に「年末のギフト超分析」を実行したとき、MCP サーバーが HTTP スレッド内で100の重いタスクを同期に抱えるのは避けたい状況です。スパイクを受け止め、タスクをキューに並べ、複数のワーカーで処理できる層が必要になります。

最後に、UX です。ユーザーが GiftGenius のボタンを押して 40 秒間ただスピナーを眺めるだけ——まるで昔のインターネットバンキングのような体験です。「すぐの応答 + 進捗 + キャンセル可能」というモデルの方がはるかに快適です。

これらの問題は「起動 → キュー → バックグラウンド → イベント」という共通パターンで解決できます。

2. MCP における async‑job の基本アーキテクチャ

GiftGenius を例に取りましょう。新しい重いシナリオが増えたとします: 「友人の購買履歴やソーシャルから嗜好を深掘り分析する」。この処理は数分かかる可能性があるため、次のようにします。

  1. MCP のツール(tool)がモデルからのパラメータを受け取る。
  2. すぐに計算するのではなく、データベースに Job レコードを作成する。
  3. タスクをキューに積む。
  4. ChatGPT に即座に「分析を開始しました。jobId はこれです」と返す。
  5. バックグラウンドのワーカーがキューからタスクを取り出し、重い処理を実行しながら MCP イベント job.progressjob.partial を送信し、 最後に job.completed または job.failed を送る。

アーキテクチャの見た目はおおよそ次のとおりです:

flowchart LR
    subgraph ChatGPT
      U[ユーザー] --> GPT[モデル + ChatGPT UI]
    end

    GPT -->|call_tool analyze_preferences| MCP[MCPサーバー]

    subgraph Backend
      MCP -->|Job作成| DB[(ジョブDB)]
      MCP -->|enqueue| Q[キュー]
      W[ワーカー] -->|ジョブ取得| Q
      W -->|ステータス/進捗更新| DB
      W -->|MCPイベント: job.progress/job.completed| MCP
    end

    MCP -->|SSEイベント| GPT

重要な考え方: MCP サーバーは必ずしもモノリスである必要はありません。多くの場合、内部の非同期基盤に対するファサードとして振る舞います。つまり、tool‑call を受け取り、ジョブを作成し、イベントを送る一方、重い仕事は別プロセスのワーカーが実行します。

3. 非同期ジョブのデータモデル

まずはシンプルな Job モデルから始めましょう。TypeScript と仮の Node/MCP サーバーを使い、実際のスタックにどう組み込むかがすぐに分かるようにします。

メモリ/DB 上の最小モデルは次のようになります:

// openai/jobs/model.ts
export type JobStatus =
  | 'pending'
  | 'in_progress'
  | 'completed'
  | 'failed'
  | 'canceled';

export interface GiftJob {
  id: string;                // jobId
  type: 'deep_gift_analysis';
  status: JobStatus;
  payload: {
    recipientProfile: string;  // プロフィールのテキスト/ID
    budget: number;
  };
  result?: unknown;          // 最終的なレコメンデーション
  error?: string;            // エラー原因
  attempts: number;          // 実行試行回数
  createdAt: Date;
  updatedAt: Date;
}

実プロジェクトでは GiftJob を Postgres、DynamoDB、Firestore などに保存するでしょうが、この講義で重要なのは次のフィールドです。

  • status — タスクの現在の状態。イベントにも UX にも反映されます;
  • attempts — retry のためのカウンタ;
  • error — ログとデバッグ用;
  • payload — ワーカーが処理に使う入力データ。

4. async‑job を作成する MCP ツール

start_deep_analysis というツールを想定します。以前は同期的に全部やっていたかもしれませんが、今はジョブをキューに積み、jobId を返すだけにします。

// openai/tools/startDeepAnalysis.ts
import { v4 as uuid } from 'uuid';
import { createJobAndEnqueue } from '../jobs/queue';

// MCP SDKの疑似型
type StartDeepAnalysisInput = {
  recipientProfile: string;
  budget: number;
};

type StartDeepAnalysisOutput = {
  jobId: string;
  message: string;
};

export async function startDeepAnalysisTool(
  input: StartDeepAnalysisInput
): Promise<StartDeepAnalysisOutput> {
  const jobId = uuid();

  await createJobAndEnqueue({
    id: jobId,
    type: 'deep_gift_analysis',
    status: 'pending',
    payload: {
      recipientProfile: input.recipientProfile,
      budget: input.budget,
    },
    attempts: 0,
    createdAt: new Date(),
    updatedAt: new Date(),
  });

  return {
    jobId,
    message: `詳細分析を開始しました。ジョブID: ${jobId}。進捗が出次第、順次アップデートを送ります。`,
  };
}

ここで重要なのは次の点です。

  • MCP ツールは高速に動作すること(せいぜい DB/キューへの数回のアクセス)。
  • jobId を含む構造化レスポンスを返すこと。ChatGPT はこれを「ユーザーへの説明」に使えますし、GiftGenius のウィジェットは自身の widgetState に保存できます。

このツールの JSON Schema は jobId(文字列)と message(人間が読むためのテキスト)を記述するだけです。モデルはそれがジョブの識別子だと理解し、後続の対話ステップでも参照できます。

5. シンプルなキューとワーカー: 学習用の実装

いまは Redis や RabbitMQ などを導入せず、インメモリの簡易キューで進めます。実運用ではもちろん別サービス(SQS/BullMQ/Cloud Tasks 等)にしますが、ロジックは同じです。

まずはキューのひな型:

// openai/jobs/queue.ts
import type { GiftJob } from './model';

const jobs = new Map<string, GiftJob>();   // メモリ内の「DB」
export const queue: string[] = [];         // IDベースの簡易キュー

export async function createJobAndEnqueue(job: GiftJob) {
  jobs.set(job.id, job);
  queue.push(job.id);
}

export function getJob(id: string): GiftJob | undefined {
  return jobs.get(id);
}

export function updateJob(id: string, patch: Partial<GiftJob>) {
  const job = jobs.get(id);
  if (!job) return;
  const updated: GiftJob = { ...job, ...patch, updatedAt: new Date() };
  jobs.set(id, updated);
}

次に、定期的にキューを覗いてジョブを処理する素朴なワーカー:

// openai/jobs/worker.ts
import { getJob, updateJob } from './queue';
import { emitJobEvent } from './events';

async function processJob(jobId: string) {
  const job = getJob(jobId);
  if (!job) return;

  updateJob(jobId, { status: 'in_progress' });
  await emitJobEvent(jobId, 'job.started', {});

  try {
    // ここで重いビジネスロジックを呼び出す
    const result = await doDeepGiftAnalysis(job.id, job.payload);

    updateJob(jobId, { status: 'completed', result });
    await emitJobEvent(jobId, 'job.completed', { resultSummary: summarize(result) });
  } catch (err) {
    updateJob(jobId, {
      status: 'failed',
      error: (err as Error).message,
    });
    await emitJobEvent(jobId, 'job.failed', { error: 'Internal error' });
  }
}

そして、アプリ起動時などにどこかで動かせる「ループ型」ワーカー:

// openai/jobs/workerLoop.ts
import { queue } from './queue';
import { processJob } from './worker';

export function startWorkerLoop() {
  setInterval(async () => {
    const jobId = queue.shift(); // 本来は競合対策が必要
    if (!jobId) return;

    await processJob(jobId);
  }, 1000); // 1秒ごとにキューをチェック
}

これは学習用の例です。実運用では setInterval の代わりに、メッセージ到着時に自動でワーカーを「起こす」本物のキューを使います。ただし大事なのは、ワーカーが MCP ツールから分離され、バックグラウンドで動き、MCP サーバーとイベントでやり取りしているという構図です。

6. ワーカーから MCP イベントを生成

前の講義で MCP イベントのフォーマット(type、一意の event_idtimestampjob_idpayload)を見ました。 ここでは、ワーカーが emitJobEvent ヘルパーを呼び、その中で MCP サーバーの SSE チャネルへイベントを届ける方法を示します。

シンプルなヘルパーの例:

// openai/jobs/events.ts
import { randomUUID } from 'crypto';
import { sendMcpEvent } from '../mcp/eventBus';

export async function emitJobEvent(
  jobId: string,
  type: 'job.started' | 'job.progress' | 'job.completed' | 'job.failed',
  payload: unknown
) {
  const event = {
    event_id: randomUUID(),
    type,
    job_id: jobId,
    timestamp: new Date().toISOString(),
    payload,
  };

  await sendMcpEvent(event);
}

sendMcpEvent は MCP サーバー側で、 MCP SDK の SSEServerTransport にどう流し込むかを把握しています。たとえばローカルのイベントバスや Redis Pub/Sub を経由します(モジュール 12 で扱いました)。

重要なポイント: ワーカーは ChatGPT と直接は通信しません。ワーカーは MCP サーバーと通信し、MCP サーバーが SSE 接続を維持してクライアントへイベントを中継します。

7. 進捗と部分結果(partial results)をワーカーから送る

いよいよ醍醐味の進捗と部分結果です。GiftGenius の長い分析は次の段階に分けられます。

  • データ収集と正規化;
  • 基本セグメントの構築;
  • 一次的なギフト案の生成;
  • 最終的な再ランク付けと説明文の生成。

各段階で job.progress を、場合によっては job.partial も送って、UI に先行して最初の案を表示させられます。

仮のワーカー:

async function doDeepGiftAnalysis(jobId: string, payload: GiftJob['payload']) {
  await emitJobEvent(jobId, 'job.progress', { step: 1, totalSteps: 4 });

  const normalized = await collectAndNormalizeData(payload);
  await emitJobEvent(jobId, 'job.progress', { step: 2, totalSteps: 4 });

  const roughGifts = await generateInitialGifts(normalized);
  await emitJobEvent(jobId, 'job.partial', { gifts: roughGifts.slice(0, 3) });

  await emitJobEvent(jobId, 'job.progress', { step: 3, totalSteps: 4 });

  const finalGifts = await rerankAndBeautify(roughGifts);
  await emitJobEvent(jobId, 'job.progress', { step: 4, totalSteps: 4 });

  return finalGifts;
}

ウィジェットはイベントを購読して、まずは「まだ調整中」という注釈付きで 3 件の「草案」ギフトを表示し、 job.completed の後にリストを更新し、ローディング指標を外すことができます。 この UX は講義 3 で扱ったパターンにぴったりはまります。

8. ワーカーのための retry ロジック

ここからが一番ヒリつく箇所、エラーと再試行です。

ワーカーが外部の商品の API を叩く際、500429 が時々返るとしましょう。 最初の失敗で諦めるのは不自然ですが、無限に再試行するのもダメです。自分たちや相手のサービスを DDoS してしまいます。

必要なのは、指数的な遅延と試行回数の上限を持つ retry 戦略です。

以降の講義でも使う、エラーの分類から始めましょう。

  • 一時的(transient)— タイムアウト、500503429;
  • 恒久的(permanent)— 不正な入力、存在しないリソース;
  • 致命的(bug)— バグ、TypeError、想定外の例外。

再試行すべきなのは一時的エラーだけです。その他は正直に 'failed' とすべきです。

簡略化してヘルパーを用意します。

// openai/jobs/retry.ts
export function shouldRetry(error: unknown): boolean {
  if (!(error instanceof Error)) return false;
  // 仮に: HTTP 5xx または 429
  return /5\d\d|429/.test(error.message);
}

export function getDelayMs(base: number, attempt: number): number {
  const jitter = Math.random() * 100;   // 小さなジッター
  return base * 2 ** attempt + jitter; // 指数バックオフ
}

次に、attemptsGiftJob に反映するようワーカーを更新します。

// openai/jobs/worker.ts
import { getJob, updateJob } from './queue';
import { emitJobEvent } from './events';
import { shouldRetry, getDelayMs } from './retry';

const MAX_ATTEMPTS = 5;

export async function processJob(jobId: string) {
  const job = getJob(jobId);
  if (!job) return;

  updateJob(jobId, { status: 'in_progress' });

  try {
    const result = await doDeepGiftAnalysis(job.id, job.payload);

    updateJob(jobId, { status: 'completed', result });
    await emitJobEvent(jobId, 'job.completed', {
      resultSummary: summarize(result),
    });
  } catch (err) {
    const attempts = job.attempts + 1;
    const error = err as Error;

    if (attempts <= MAX_ATTEMPTS && shouldRetry(error)) {
      const delay = getDelayMs(1000, attempts); // 1s,2s,4s...

      updateJob(jobId, { attempts, status: 'pending', error: error.message });

      setTimeout(() => {
        // 実際のキューなら、遅延付きで再投入します
        processJob(jobId);
      }, delay);

      await emitJobEvent(jobId, 'job.progress', {
        retry: attempts,
        nextAttemptInMs: delay,
      });
    } else {
      updateJob(jobId, { status: 'failed', error: error.message });
      await emitJobEvent(jobId, 'job.failed', {
        error: '複数回の再試行後も分析を完了できませんでした',
      });
    }
  }
}

ここではいくつか重要な点があります。

第一に、attempts はジョブ自身に保存します。これはログや可観測性の面で便利です(グラフで retry の有無が視覚的に分かります)。

第二に、各 retry ごとに job.progress を送信し、これは何回目の試行かを明示します。モデルはこの情報を使って「ギフトのサーバーが不安定なので再試行中です」といった説明をユーザーに伝えられます。

第三に、最終的には必ず job.completedjob.failed のどちらかが送られることを保証します。「生きているのか死んでいるのか分からない」ジョブは残しません。

キャンセル('canceled')も重要なステータスです。学習用の例では実装しませんが、プロダクションではユーザーの操作(ウィジェットの「キャンセル」ボタン)やタイムアウトで設定されるのが一般的です。 その場合、ワーカーは次にジョブをキューから取るときに status: 'canceled' を見て処理を開始せず、MCP サーバーが最終イベント job.canceled を送ります。

9. 冪等性と retry: 同じ失敗を二度踏まない

retry を導入すると、「同じことを二度やってしまう」リスクがすぐに生まれます。コマース系では致命的(例えば二重課金)ですが、GiftGenius でも悪影響があります(友人へ同じメールを二重送信、内部分析に重複レコードを作る、など)。

そこで、次の2つの原則を意識します。

第一: ジョブのハンドラは冪等にする。

同じ jobId を(retry や誤動作で)複数回呼んでも、世界が壊れないようにします。具体的には:

  • すべての副作用(DB 書き込み、メール送信、注文作成)は jobId や自然キーに結びつけ、同じステップを既に実施したかをすばやく判定できるようにする。
  • job.status がすでに 'completed' または 'failed' なら、再実行は無視するか、既存の結果を返す。

簡単な防御の例:

export async function processJob(jobId: string) {
  const job = getJob(jobId);
  if (!job) return;

  if (job.status === 'completed' || job.status === 'failed') {
    // ジョブは既に正常完了または失敗で終了しています
    return;
  }

  // ... 残りのコード
}

第二: イベントも冪等であるべき。

すでに event_id と、クライアント側での重複排除の話をしましたが、サーバー側でも注意が必要です。ワーカーの再起動やキューからの復旧時に、同じ job.progress を無闇にスパムしないようにしましょう。

10. あなたのアーキテクチャでキューとワーカーをどこに置くか

図では綺麗ですが、実際にワーカーはどこで動かすのでしょうか。代表的な選択肢がいくつかあります。

統合型ワーカー: MCP サーバーとワーカーが同一プロセス/デプロイ。tool‑call を受け、同じプロセスでワーカーループも動かします。利点はシンプルさ(サービスが少なく、デプロイが容易)。欠点はスケーリング(ワーカーを増やすには MCP サーバー全体をスケールする必要がある)。

分離型ワーカー: MCP サーバーとワーカーを別サービスに分離。その間にキューがあり、イベント用に Pub/Sub を挟むことも。BullMQ/Redis と MCP イベントに関する文脈でよく語られる形です。MCP サーバーは Redis の 'mcp:events' チャンネルを購読し、ワーカーはそこへイベントを publish します。

ハイブリッド: MCP サーバーのあるインスタンスはワーカーも回し、その他のインスタンスは HTTP/SSE のみを担当。Vercel など、常駐バックグラウンドプロセスが苦手な serverless 環境で有効です。

学習用の GiftGenius では、まず統合型(MCP サーバー + 単純なワーカー)で十分です。プロダクションとスケーリングのモジュールに進んだら、ワーカーを別サービスに分離して移行しましょう。

11. 例: GiftGenius の完全な async パイプライン

ユーザーがチャットで次のように書いたとき、何が起きるかを流れで確認しましょう。

「宇宙好きの友人に合う複雑なギフト選定が必要です。過去の購入履歴も考慮してください。」

  1. モデルは start_deep_analysis ツールの呼び出しを選び、受取人プロファイルと予算をパラメータとして渡す。
  2. ツールは GiftJob を DB に 'pending' ステータスで作成し、 それをキューへ投入して、jobId + 確認メッセージを返す。
  3. ChatGPT は分析開始をユーザーに説明し、jobId を GiftGenius のウィジェットに渡す場合もある。
  4. ウィジェットはその jobId のイベントに SSE で購読し、 進捗バーと「データを収集・分析しています」というステータスを表示する。
  5. ワーカーは新しいジョブをキューで検知すると、ステータスを 'in_progress' に更新し、job.started を送る。
  6. 処理中に複数回 job.progress(段階)と job.partial(最初の 23 件のギフト案)を送る。
  7. 外部 API がどこかで落ちた場合、ワーカーは指数バックオフ付きで再試行し、 attempts を更新し、再試行情報を含むイベントを送る。
  8. 最後に job.completed(簡潔なサマリと最終レコメンデーション付き)または job.failed(分かりやすい説明付き)を送る。
  9. ウィジェットはこれらのイベントに基づいて UI を更新し、ChatGPT はテキストのサマリを生成してフォローアップ(「さらに案を表示」「予算を絞る」「ギフトのタイプを変える」)を提案できる。

ユーザー視点では「生きた」長時間プロセスがコントロール下にあります。バックエンド視点では、キュー、ワーカー、retry を伴う健全な async パイプラインです。

12. 小さな演習(自習)

内容を定着させたい場合、GiftGenius について次を試してみてください。

  • 実データベース用の jobs テーブル設計を考える: 必要なインデックス、どのフィールドでフィルタするか(ユーザー別、ステータス別、作成日時別)。
  • ウィジェットが SSE 不可時のフォールバックとしてステータスをポーリングできるよう、 HTTP エンドポイント /api/jobs/:id 用の TypeScript 型を草案する。
  • retry ポリシーを記述する: 試行回数、基本遅延、最終的に失敗したジョブの扱い (単純な dead‑letter テーブル、あるいはロギング + アラート)。

この課題は、プロダクション運用と可観測性のモジュールで「pending ステータスのまま N 分以上滞留しているジョブはいくつか」といったメトリクスを扱う際に役立ちます。

13. 非同期ジョブでよくある落とし穴

よくある誤り1: すべてを tool‑call の同期処理でやる。
最も多い落とし穴は、重い処理をキューなしで1つの MCP ツールに押し込むことです。リクエストが少ないうちは動いて見えますが、負荷が増すか外部 API が遅くなると、タイムアウトやチャットのフリーズ、非常に悪い UX を招きます。数十秒以上かかる可能性のある処理は、最初から jobId 付きの async‑job として設計しましょう。

よくある誤り2: 明確な Job モデルがない。
「キューのメッセージだけ」で済ませ、DB にジョブ状態を持たない開発者もいます。その結果、基本的な問い(「ジョブのステータスは?」「何回試した?」「なぜ落ちた?」)に答えにくくなります。 Job の明確なモデル(statusattemptserrorcreatedAt など)は、デバッグ、監視、UX の土台です。

よくある誤り3: retry がない、または無限に再試行する。
まったく retry せず最初の 500 で落ちるか、while (!success) のように試行回数を制限しないか、両極端になりがちです。前者では短時間の障害で多くのジョブを失い、後者では負荷の「嵐」を引き起こし、外部 API をブロックするリスクもあります。必要なのは中庸です: 有限回の試行 + 指数バックオフ + 一時的エラーと恒久的エラーの切り分け。

よくある誤り4: 非冪等なハンドラ。
各試行で毎回、外部システムに新規レコードを作成したり、同じ支払いを実行したり、同じメールを送信したりすると、retry はすぐ問題化します。ハンドラは、この jobId のジョブがすでに完了しているかを判断し、危険な副作用の重複を避けるべきです。

よくある誤り5: エラー時にイベントを出さない。
ワーカーが想定外の例外で落ち、コンソールにログして終わり……というケースがあります。ユーザーは job.completed を永遠に待ち続け、すでに処理が死んでいることに気付けません。エラーで処理が終わるすべての分岐は、最終的に job.failed を出し、DB の Job ステータスを更新すべきです。これがないと MCP のストリームは一方通行の「ブラックボックス」になります。

よくある誤り6: 進捗イベントが多すぎる。
「正直でありたい」と考えて 1% ごとに job.progress を送ると、ネットワーク、クライアント、MCP サーバーを過負荷にします。進捗はステージの切り替えや大きな差分(例えば 10% ごと)で送るのが良く、それ以外は内部ログにとどめましょう。

よくある誤り7: インメモリキューを本番で使う。
queue: string[]Map を使った学習用の例は、アーキテクチャ理解には良いですが、実運用のシステムではプロセス再起動やサーバーダウンの一撃で崩壊します。本格運用には外部のキューとストレージ(SQS、Pub/Sub、RabbitMQ、Redis Streams など)が必要です。インメモリ方式はローカル開発や簡単なデモに留めましょう。

コメント
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION