CodeGym /課程 /ChatGPT Apps /非同步任務:佇列、工作者、重試 (retry)

非同步任務:佇列、工作者、重試 (retry)

ChatGPT Apps
等級 13 , 課堂 3
開放

1. 為什麼在 ChatGPT App 需要非同步任務

如果世界很理想,你的每個 MCP 工具都能在幾百毫秒內完成。但在現實中,有趣的東西通常既漫長又吃力:

  • 解析包含使用者購買歷史的大型 CSV;
  • 彙整多個外部 API 的資料,而每個 API 不是睡著就是回應 503
  • 建構包含許多中間步驟的複雜推薦;
  • 產生大型報告與簡報。

若嘗試把這些全部塞進單一同步的 tool‑call,你會遇到三個問題。

第一,逾時。ChatGPT 會話、HTTP 基礎設施、MCP 客戶端——這些都不是為「五分鐘後才回應」而設計。伺服器把連線掛太久,對 ChatGPT 和使用者來說都像是「卡住」。

第二,負載管理。如果同時有一百位使用者啟動「新年禮物超級分析」,你不會希望 MCP 伺服器在 HTTP 執行緒裡同步地扛著一百個長任務。你需要一層能吸收尖峰、將任務排入佇列並用多個工作者處理的中介。

第三,UX。使用者在 GiftGenius 小工具按下按鈕後,盯著一個轉圈圈 40 秒——感覺就像早年的網銀。我們更喜歡「快速回覆 + 進度 + 可取消」的模式。

這些問題可以用一個通用流程解決:「啟動 → 佇列 → 背景 → 事件」。

2. MCP 情境下的 async‑job 基本架構

以 GiftGenius 為例。假設有個新的重型情境:「依據購買歷史與朋友的社群資料做深度偏好分析」。這類操作可能會跑好幾分鐘,因此:

  1. MCP 工具(tool)接收模型傳來的請求參數。
  2. 它不會立刻把所有事算完,而是先在資料庫中建立一筆 Job 紀錄。
  3. 把任務放入佇列。
  4. 立刻回覆 ChatGPT:「分析已啟動,這是 jobId」。
  5. 背景工作者從佇列取出任務,執行繁重工作,過程中發送 MCP 事件 job.progressjob.partial, 最後發送 job.completedjob.failed

從架構角度看,大致如下:

flowchart LR
    subgraph ChatGPT
      U[使用者] --> GPT[模型 + ChatGPT UI]
    end

    GPT -->|call_tool analyze_preferences| MCP[MCP 伺服器]

    subgraph Backend
      MCP -->|建立 Job| DB[(任務資料庫)]
      MCP -->|enqueue| Q[佇列]
      W[工作者] -->|取出任務| Q
      W -->|更新狀態/進度| DB
      W -->|MCP 事件: job.progress/job.completed| MCP
    end

    MCP -->|SSE 事件| GPT

重點:MCP 伺服器不一定是單體。它常常只是你內部非同步基礎設施的門面:接收 tool‑call、建立 job 並發送事件,而繁重的工作由獨立的工作者程序完成。

3. 非同步任務的資料模型

從簡單的 Job 模型開始。我們會用 TypeScript 與假想的 Node/MCP 伺服器,讓你直接看到它如何融入你的技術棧。

最簡單的記憶體/資料庫模型可以長這樣:

// 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 工具需要很快:最多對資料庫/佇列做一兩個請求;
  • 它回傳包含 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>();   // 記憶體中的「資料庫」
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);
}

接著是一個會定期查看佇列、取出並處理 job 的原始工作者:

// 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); // 每秒檢查一次佇列
}

這是教學範例。在實務上,會用真正的佇列來在有新訊息時「喚醒」工作者。 但核心概念很清楚:工作者與 MCP 工具解耦,於背景運作,並透過事件與 MCP 伺服器通訊。

6. 由工作者產生 MCP 事件

先前的講座中你已看過 MCP 事件的格式:型別(type)、唯一的 event_idtimestampjob_idpayload。 現在示範工作者如何呼叫 helper emitJobEvent, 而它會把事件透過 MCP 伺服器的 SSE 通道送到 ChatGPT。

簡單的 helper 範例:

// 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);
}

而 MCP 伺服器內的 sendMcpEvent 會知道要如何把事件推送進 MCP SDK 的 SSEServerTransport:例如透過本機事件匯流排或 Redis Pub/Sub,就像我們在模組 12 中談過的。

關鍵想法:工作者不會直接與 ChatGPT 通訊。它只與 MCP 伺服器通訊,而後者維持 SSE 連線並把事件轉發給客戶端。

7. 工作者的進度與部分結果

接著來看最有意思的:進度與部分結果。在 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 後,更新清單並移除載入指示器。 這與我們在第 3 講討論的 UX 模式完美契合。

8. 工作者的重試邏輯 (retry)

現在進入最讓人緊張的部分:錯誤與重試。

想像工作者在處理任務時要呼叫外部商品清單 API,而該 API 偶爾回應 500429。 第一次失敗就放棄很奇怪;但無限次重試也不行:你會把自己或外部服務打爆。

我們需要一個帶有指數延遲且有次數上限的重試策略。

先做個錯誤分類,之後在本課程也會用到:

  • 暫時性(transient)—— 逾時、500503429
  • 永久性(permanent)—— 錯誤的輸入、資源不存在;
  • 致命(bug)—— 程式錯誤、TypeError、非預期例外。

只有暫時性錯誤值得重試。其他應如實標記為 'failed'

簡化成一個 helper:

// 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; // 指數退避
}

現在更新工作者,讓它考量 attempts(存於 GiftJob):

// 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(() => {
        // 在真正的佇列中,你會以延遲「重新 enqueue」這個任務
        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 存在任務本身——這對記錄與可觀測性很方便 (在圖表上也能清楚看出有多少任務經歷了重試)。

第二,每次重試時我們都會送出帶有「第 N 次嘗試」資訊的 job.progress。 模型可以用這些資訊向使用者解釋:「禮物服務回應不穩定,我會再試一次」。

第三,我們保證最終一定會送出 job.completedjob.failed。 不會有任務「不死不活」掛著不動。

取消('canceled')也是重要狀態。教學範例中我們不實作, 但在生產環境通常由使用者主動觸發(小工具上的「取消」按鈕)或因逾時而觸發。 這種情況下,工作者在下一次從佇列取到任務時看到 status'canceled', 就不會啟動處理,MCP 伺服器會發送最終事件 job.canceled

9. 冪等性與重試:避免重複副作用

一旦引入重試,就會有「同一件事做兩次」的風險。在商務模組中這可能很嚴重(例如重複扣款), 在 GiftGenius 中也可能出問題:對朋友寄出兩封相同的信、在內部分析系統中重複寫入等。

因此請把握兩個原則。

第一:任務處理器必須是冪等的。

如果你以相同的 jobId 呼叫它多次(因重試或錯誤),世界不該崩壞。為了達成此目標:

  • 所有副作用(寫入資料庫、寄送郵件、建立訂單)都應綁定到 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,也啟動 worker loop。優點是簡單:服務更少、部署容易。 缺點是擴充:要增加工作者,就得把整個 MCP 伺服器一起擴大。

獨立工作者:MCP 伺服器是一個服務,工作者是另一個。 兩者之間是佇列,可能再加上事件的 Pub/Sub。這在 BullMQ/Redis 與 MCP 事件的文章裡常見: MCP 伺服器訂閱 Redis 頻道 'mcp:events',工作者把事件發佈到那裡。

混合式:一個 MCP 伺服器實例同時跑工作者,其餘實例只負責 HTTP/SSE。 若你部署在 Vercel 或其他 serverless 平台,而長駐的背景程序不太友善,這會很有用。

我們在教學用的 GiftGenius 先採第一種:MCP 伺服器 + 單一簡單的程序內工作者。 當你進入生產與擴充的模組後,可以把工作者遷移成獨立服務。

11. 範例:GiftGenius 的完整 async 流水線

讓我們串起來看看,當使用者在聊天輸入:

「我需要為一位太空迷挑選複雜的禮物,並考量他的過往購買紀錄」。

  1. 模型決定呼叫工具 start_deep_analysis,帶入收禮者的個人檔案與預算參數。
  2. 工具在資料庫建立一筆狀態為 'pending'GiftJob, 把它放入佇列,並回傳 jobId + 確認訊息。
  3. ChatGPT 向使用者說明分析已啟動,並可把 jobId 傳給 GiftGenius 小工具。
  4. 小工具透過 SSE 訂閱該 jobId,顯示進度列與「正在收集與分析資料」的狀態。
  5. 工作者在佇列看到新任務,將狀態更新為 'in_progress',並發送 job.started
  6. 過程中會多次發送 job.progress(階段)與 job.partial(前 23 個禮物)。
  7. 若途中外部 API 掛掉,工作者會以指數退避重試,更新 attempts,並發送含重試資訊的事件。
  8. 最後要嘛發送帶摘要與最終建議的 job.completed, 要嘛發送含清楚說明的 job.failed
  9. 小工具依據這些事件更新 UI,而 ChatGPT 可以生成文字摘要並提供後續動作: 「顯示更多點子」、「縮小預算」、「更換禮物類型」。

對使用者而言,這是一個「有生命」且可掌控的長流程。 對後端而言,這是一條正常的 async 流水線:佇列、工作者與重試。

12. 小練習(自我實作)

若想加強練習,請在 GiftGenius 動手試試:

  • 設計真實資料庫的 jobs 表結構: 需要哪些索引、哪些欄位會用在篩選(依使用者、依狀態、依建立時間)。
  • 草擬 HTTP 端點 /api/jobs/:id 的 TypeScript 型別, 以便小工具在 SSE 不可用時最少能輪詢狀態。
  • 描述你的重試策略:嘗試次數、基礎延遲、對仍然失敗的任務要做什麼 (簡單的 dead‑letter 表或是記錄 + 警示)。

這份作業在之後的生產與可觀測性模組會用到,屆時我們會討論像 「有多少任務在 pending 狀態超過 N 分鐘」這類指標。

13. 處理非同步任務的常見錯誤

錯誤 #1:在 tool‑call 中把一切都做成同步。
最常見的陷阱就是試圖把所有繁重工作塞進單一 MCP 工具而不使用佇列。 在請求量少時似乎沒問題;一旦負載上升或外部 API 變慢,就會遇到逾時、 聊天卡頓與非常糟的 UX。凡是可能持續數十秒以上的操作, 最好一開始就設計成帶有 jobId 的 async‑job。

錯誤 #2:沒有明確的 Job 模型。
有時開發者只想靠「佇列中的訊息」撐過去,而不把任務狀態存入資料庫。結果連基本問題都很難回答: 「任務狀態是什麼?」「我們嘗試了幾次?」「為什麼會失敗?」。 明確的 Job 模型,加上 statusattemptserrorcreatedAt 等欄位,是除錯、監控與 UX 的基礎。

錯誤 #3:沒有重試,或相反地,無限重試。
有人完全不做重試,第一次遇到 500 就失敗;有人寫了 while (!success) 而不限制嘗試次數。前者會因短暫故障而丟失大量任務,後者會掀起負載「風暴」, 甚至把外部 API 掛死。需要合理的中庸:有限次數 + 指數延遲 + 區分暫時性與永久性錯誤。

錯誤 #4:非冪等的處理器。
如果每次重試都在外部系統新建一筆紀錄而不檢查、重複執行同一筆付款或寄出同一封信—— 重試很快就會變成問題。處理器必須能判斷某個 jobId 是否已成功完成, 並避免重複產生危險的副作用。

錯誤 #5:錯誤時沒有事件。
常見情況是工作者因非預期例外而崩潰,只把錯誤印到 console 就結束。使用者持續等待 job.completed,卻不知道其實早就失敗了。任何失敗結束的分支, 最後都應導向 job.failed,並更新資料庫中的 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