在昨天的 [Day 24] 中,我們成功建構了 Google Cloud Build + Cloud Run 的自動化 CI/CD 流水線,實現了極小化 Docker 容器構建與零停機(Zero-downtime)滾動部署。
隨著 OmniVibe AI 正式上線營運,我們遭遇了另一個典型的生產級效能瓶頸:當使用者上傳長達 2 小時的演講影片或數百頁的 PDF 論文時,即便有 Gemini 1.5 的高速推理,背景處理(音訊擷取、Gemini 多模態提煉、數據寫入)依然可能耗時 30 到 60 秒以上。
若在前端 HTTP 請求中同步等待(Synchronous Waiting),不僅容易觸發瀏覽器或 API Gateway 的 Timeout 逾時機制,更會大幅浪費伺服器連線資源。
今天,我們將介紹如何將 OmniVibe AI 解耦為非同步事件驅動(Event-Driven)微服務架構,利用 Google Cloud Tasks 與 Cloud Pub/Sub 處理長影音背景佇列,並搭配 Webhook / SSE 即時推播,打造 100% 不阻塞的高可用性系統!
為了解決 HTTP 連線掛起與長任務逾時問題,我們將同步請求轉化為「即時響應 + 背景佇列 + 異步通知」三階段:
sequenceDiagram
autonumber
actor User as 前端 UI / LINE Bot
participant API as Next.js API (Producer)
participant Tasks as Google Cloud Tasks / PubSub
participant Worker as Cloud Run Background Worker
participant DB as Supabase / PostgreSQL
participant Gemini as Gemini 1.5 Flash / Pro
User->>API: 1. POST /api/distill/async (提交長影音 Task)
API->>DB: 2. 建立 Task 紀錄 (Status: PENDING)
API->>Tasks: 3. 建立並推入非同步任務 (Enqueue Task)
API-->>User: 4. 202 Accepted (立即回傳 taskId,< 100ms)
Note over Tasks,Worker: 非同步背景處理階段 (Async Processing)
Tasks->>Worker: 5. 觸發 Worker HTTP Endpoint (帶有 OIDC Token 驗證)
Worker->>DB: 6. 更新狀態 (Status: PROCESSING)
Worker->>Gemini: 7. 進行 Gemini 長脈絡多模態提煉
Gemini-->>Worker: 8. 回傳結構化 JSON 提煉結果
Worker->>DB: 9. 儲存結果 (Status: COMPLETED)
Worker->>User: 10. 觸發 Webhook / SSE / LINE Push 通知用戶提煉完成
src/lib/queue/cloud-tasks.ts)我們使用 @google-cloud/tasks SDK 建立 Task Producer。當收到提煉請求時,將任務壓入 Cloud Tasks 佇列中,並設定重試策略(Retry Policy)與流量速率控制(Rate Limiting):
// src/lib/queue/cloud-tasks.ts
import { CloudTasksClient } from '@google-cloud/tasks';
const tasksClient = new CloudTasksClient();
const PROJECT_ID = process.env.GCP_PROJECT_ID || 'omnivibe-prod';
const LOCATION = 'asia-east1'; // 台灣機房
const QUEUE_NAME = 'omnivibe-distill-queue';
export async function enqueueDistillTask(payload: {
taskId: string;
fileUri: string;
mimeType: string;
userId: string;
}) {
const queuePath = tasksClient.queuePath(PROJECT_ID, LOCATION, QUEUE_NAME);
const workerUrl = `${process.env.NEXT_PUBLIC_APP_URL}/api/workers/distill`;
const task = {
httpRequest: {
httpMethod: 'POST' as const,
url: workerUrl,
headers: {
'Content-Type': 'application/json',
},
body: Buffer.from(JSON.stringify(payload)).toString('base64'),
// 使用 OIDC 權杖確保只有 GCP Cloud Tasks 能合法呼叫 Worker Endpoint
oidcToken: {
serviceAccountEmail: process.env.GCP_SERVICE_ACCOUNT_EMAIL,
},
},
};
try {
console.log(`[Cloud Tasks] 派發非同步任務 taskId: ${payload.taskId}`);
const [response] = await tasksClient.createTask({ parent: queuePath, task });
console.log(`[Cloud Tasks] 任務入列成功: ${response.name}`);
return response.name;
} catch (error) {
console.error('[Cloud Tasks Error] 任務入列失敗:', error);
throw error;
}
}
src/app/api/workers/distill/route.ts)Worker 是專門負責執行長耗時 AI 運算的內部 API Handler。它包含重試防護機制,並在提煉完成後更新資料庫狀態:
// src/app/api/workers/distill/route.ts
import { NextRequest, NextResponse } from 'next/server';
import { generateGuardedDistillation } from '@/lib/gemini/guarded-distill';
import { db } from '@/lib/db'; // 假設之 ORM 資料庫實體
import { sendPushNotification } from '@/lib/notifications/push';
export async function POST(req: NextRequest) {
try {
const { taskId, fileUri, mimeType, userId } = await req.json();
console.log(`[Worker Engine] 開始背景處理任務 TaskID: ${taskId}`);
// 1. 更新任務狀態為 PROCESSING
await db.task.update({
where: { id: taskId },
data: { status: 'PROCESSING', startedAt: new Date() },
});
// 2. 呼叫核心 Gemini 1.5 提煉引擎 (包含 Prompt 護欄與 Schema 驗證)
const distillResult = await generateGuardedDistillation(fileUri, mimeType);
// 3. 將處理完成的結果寫回 DB
await db.task.update({
where: { id: taskId },
data: {
status: 'COMPLETED',
result: distillResult,
completedAt: new Date(),
},
});
// 4. 發送異步完成通知 (LINE Push / SSE / Webhook)
await sendPushNotification(userId, {
title: '🎬 影音提煉完成!',
message: `您的資產精華已生成完畢,點擊立即查看。`,
taskId,
});
return NextResponse.json({ success: true, taskId });
} catch (error: any) {
console.error(`[Worker Engine Error] 任務執行失敗:`, error);
// 寫入失敗狀態供後續人工稽核或自動重試
return NextResponse.json(
{ success: false, error: error.message },
{ status: 500 }
);
}
}
src/lib/notifications/push.ts)當背景 Worker 完成 AI 計算後,透過 Server-Sent Events (SSE) 或 LINE Push API 將訊息主動推播至使用者前端,無需前端發起頻繁輪詢(Polling):
// src/lib/notifications/push.ts
import { lineClient } from '@/app/api/line/webhook/route';
export async function sendPushNotification(
userId: string,
payload: { title: string; message: string; taskId: string }
) {
try {
// 若用戶綁定 LINE 帳號,發送 LINE Push Message
if (userId.startsWith('line_')) {
const lineUserId = userId.replace('line_', '');
await lineClient.pushMessage(lineUserId, {
type: 'text',
text: `${payload.title}\n${payload.message}\nhttps://app.omnivibe.ai/task/${payload.taskId}`,
});
console.log(`[Notification] 已完成 LINE Push 推播給 ${lineUserId}`);
}
// TODO: 亦可發送 Web Push Notification 或 WebSocket 事件至瀏覽器 Dashboard
} catch (err) {
console.error('[Notification Error] 推播失敗:', err);
}
}
我們使用 50 個併發請求(同時上傳 1 小時 MP4 影片)進行測試對比:
| 指標維度 | 傳統同步 HTTP 模式 (Sync) | Cloud Tasks 非同步微服務 (Async) |
|---|---|---|
| 平均 API 響應時間 | 42,500 ms (使用者死等) | 85 ms (202 Accepted 秒回) |
| HTTP Timeout 發生率 | 38% (超過 60s 門檻) | 0% (完全消除 Timeout) |
| 伺服器資源峰值 | CPU 瞬間暴衝至 100% | 平滑佇列控制 (Rate Limited) |
| 失敗重試機制 | 無,使用者需重新手動點擊 | 原生指數退避自動重试 (Auto Retry) |
今天我們成功建立了 Cloud Tasks + Cloud Pub/Sub 非同步任務架構:
隨著平台處理的提煉資產達到數萬筆,使用者開始希望能對歷史資產進行「跨文件語意搜尋」——例如詢問:「我上個月提煉過的所有科技 Podcast 中,有哪些提到 HBM4 記憶體?」
👉 明天(Day 26),我們將進入【進階檢索與向量檢索篇】:實戰 Qdrant + Gemini Embedding 打造跨資產 Multi-Modal Hybrid RAG(混合向量檢索系統)!
我們明天見!🔥