
先說說最近看到的一個有意思的項目。有人在 Hacker News 上展示了一個 AI 助手賣點不是聊天對話多順暢而是它自帶一個獨立收件箱inbox。用戶可以往里面投遞任務助手異步消費、處理、回填結果整個流程像一套輕量級的消息隊列。這種設計在 AI Agent 工程實踐里越來越常見當助手不再只是“問答機器人”而需要處理批量寫作、定時巡檢、工單分類、內容審核等異步任務時同步聊天的模式就不夠用了。本文會從零實現一個“帶收件箱的 AI 助手”后端服務只依賴 FastAPI 和 Python 標準庫。我們會逐步拆解收件箱的任務模型、狀態流轉、Worker 消費邏輯以及完整的 API 接口最后給出常見問題排查思路和生產環境落地的建議。無論你是剛開始接觸 AI 工程化還是想給已有助手系統增加異步任務能力這篇教程都值得收藏。1. 背景與核心概念1.1 什么是“帶收件箱的 AI 助手”傳統的 AI 助手通常是一個同步聊天接口用戶發送問題模型返回回答調用結束后整個交互就結束了。這個模式對“聊天”場景沒問題但對任務型場景存在明顯短板。比如用戶一次性提交 20 封郵件讓助手生成摘要如果同步處理客戶端必須長時間等待一旦網絡波動或接口超時前面所有結果都丟失了。“帶收件箱的 AI 助手”借鑒了異步消息系統的設計思路。助手內部維護一個收件箱所有請求先進入收件箱排隊后臺 Worker 不斷從收件箱中取出任務調用模型或工具處理再把結果回寫到對應的任務記錄上。用戶提交請求后拿到一個task_id之后可以通過這個 ID 查詢處理進度和最終結果。從架構角度看收件箱本質上是一個任務隊列只是它需要額外支持任務狀態管理排隊中、處理中、已完成、失敗。按優先級或時間排序取數。任務結果回寫與查詢。失敗重試與錯誤信息記錄。這套設計并不新鮮消息隊列領域已經實踐了多年。但當它被應用到 AI 助手場景時有一個很大的區別AI 模型調用往往是慢操作而且可能失敗必須把“任務狀態”和“處理結果”作為一等公民來管理。1.2 收件箱模式解決的核心問題給 AI 助手引入收件箱模式主要解決四個問題。第一是解耦。調用方只需要把任務投遞到收件箱不需要關心 AI 模型什么時候處理完。后端可以隨時增加消費 Worker也可以平滑升級模型服務調用方是無感的。第二是可靠性。任務被持久化后即使 Worker 進程崩潰任務記錄不會丟失。重啟后可以繼續消費未完成的任務這比同步調用里的“請求丟失”場景可靠得多。第三是可觀測性。所有任務都有明確狀態我們可以方便地統計隊列積壓量、平均處理時長、失敗率甚至對每個任務做審計。第四是并發可控。我們可以限制同時處理的任務數量避免大批量請求瞬間壓垮模型 API也可以結合令牌桶做限流。下面用一個對比表來總結同步聊天與收件箱模式的區別能力維度同步聊天模式收件箱模式請求方式請求-響應提交任務-異步回調/輪詢任務狀態無明確狀態PENDING/PROCESSING/DONE/FAILED持久性依賴客戶端連接任務記錄持久化并發控制較難Worker 數量可控失敗重試需要客戶端重試服務端自動重試適用場景實時對話批量處理、后臺任務、Agent 任務編排1.3 典型應用場景在實際項目中這種模式很適合以下場景。內容生成與摘要批量生成商品文案、新聞摘要、郵件回復草稿。工單分類與回復客服工單進入收件箱AI 自動打標、分配、生成建議回復。數據處理任務從數據庫或文件中抽取數據交給模型結構化再寫回存儲。定時巡檢報告每天定時把運營數據丟進收件箱模型生成日報后推送通知。這些場景的共同點是任務到達時間和處理時間不一定是同步的而且單次處理可能耗時幾十秒甚至幾分鐘。用收件箱模式能最大程度降低系統耦合度。2. 系統架構與消息狀態設計2.1 整體架構我們設計的系統包含四個核心角色API 層、Inbox 存儲層、Worker 消費層、AI 處理服務。下面用一張 ASCII 架構圖表示數據流客戶端 (提交任務/查詢結果) │ ▼ ┌─────────────────┐ │ FastAPI 層 │ /inbox/tasks 提交 │ │ /inbox/tasks 查詢 └─────────────────┘ │ ▼ ┌─────────────────┐ │ Inbox 存儲 │ 任務狀態 內存隊列 └─────────────────┘ │ │ 領取待處理任務 ▼ ┌─────────────────┐ │ AI Worker │ 多線程 / 多進程消費 └─────────────────┘ │ │ 調用模型 / 工具 ▼ ┌─────────────────┐ │ AI 服務 │ LLM API / 本地模型 / 腳本 └─────────────────┘ │ └── 處理完成 → 結果回寫 Inbox → 客戶端查詢API 層負責接收用戶請求把任務寫入 Inbox。Inbox 存儲任務記錄和狀態Worker 定期從中領取任務。領取后Worker 調用 AI 處理服務最后把結果回寫到對應任務記錄上。線程模型上我們的示例采用「API 線程 后臺 Worker 線程」的方式。FastAPI 啟動時拉起一個后臺 Worker 線程Worker 輪詢 Inbox每次領取一個任務。生產環境可以把這個模型替換成多進程 Worker 或獨立部署的任務消費者。2.2 任務生命周期任務在收件箱中會經歷多個狀態。這里把狀態定義清楚是整個系統設計的核心。PENDING任務已進入收件箱等待 Worker 領取。PROCESSING任務被某個 Worker 領取正在調用 AI 處理。DONE處理成功結果字段已回填。FAILED處理多次重試仍然失敗錯誤信息已記錄。狀態流轉可以用下面一段偽代碼表示提交任務 -- PENDING Worker 領取 -- PROCESSING 處理成功 -- DONE 處理失敗且還有重試次數 -- PENDING 處理失敗且達到最大次數 -- FAILED這里把“失敗后重試”和“失敗最終態”區分開非常關鍵。AI 模型接口經常因為網絡抖動、限流、內容審核等原因失敗如果一律進入 FAILED會讓很多本來可以成功的任務白白失敗如果無限重試又會造成成本浪費和隊列堆積。常見方案是設置最大嘗試次數比如 3 次前 2 次失敗回到 PENDING第 3 次失敗進入 FAILED。3. 環境準備與項目結構3.1 運行環境與依賴本文示例代碼使用 Python 3.10主要依賴 FastAPI 和 Uvicorn。數據庫方面先用內存存儲演示后續可以在最佳實踐章節替換為 Redis 或 SQLite。你需要準備的環境如下Python 3.10 或更高版本。一個虛擬環境venv 或 conda 均可。pip 安裝 fastapi、uvicorn、pydantic。版本不需要刻意固定本文示例以常見環境為準重點演示設計思路。下面代碼基于 pydantic v2 編寫如果你使用的是 pydantic v1需要把model_dump(modejson)改回dict()或者直接在模型中寫自定義序列化方法。3.2 項目目錄結構為了方便閱讀我們將代碼拆成幾個模塊結構如下ai-inbox-assistant/ ├── requirements.txt ├── app/ │ ├── __init__.py │ ├── models.py │ ├── inbox.py │ ├── worker.py │ └── main.py └── README.md各個文件職責如下文件職責requirements.txt項目依賴app/models.py任務數據模型、狀態枚舉、優先級枚舉app/inbox.pyInbox 存儲 任務狀態管理 領取策略app/worker.py后臺消費線程領取任務并調用 AI 服務app/main.pyFastAPI 應用注冊路由和生命周期4. 核心模塊設計詳解4.1 消息模型定義先定義收件箱中的任務模型。它比普通隊列消息多了一些業務字段發送方、主題、優先級、內容、狀態、嘗試次數、結果和錯誤信息。字段設計說明id任務唯一標識用 UUID 生成。sender消息來源比如email、user、cron。topic任務主題或分類方便后續篩選。content要交給 AI 處理的原始內容。priority優先級影響 Worker 取數的先后順序。status任務當前狀態。attempts當前已嘗試執行次數用于失敗重試。result處理成功后的結果。error最近一次失敗的錯誤信息。優先級建議使用枚舉這樣在 API 層校驗參數時更安全。我們定義TaskPriority枚舉包含HIGH、NORMAL、LOW三檔。任務狀態用TaskStatus枚舉包含PENDING、PROCESSING、DONE、FAILED四種。4.2 Inbox 存儲與取數策略這里我們用 Python 內置字典作為任務存儲用RLock保證線程安全。收件箱需要提供以下能力add添加任務返回新任務對象。claim_next領取下一個待處理任務。complete完成任務回填結果。fail_or_retry處理失敗判斷是否重試或進入最終失敗態。get按 ID 查詢任務。list列出任務支持按狀態篩選。claim_next是核心方法。Worker 調用它時必須保證“找出任務”和“修改狀態為 PROCESSING”是原子的否則多個 Worker 同時消費時會拿到同一個任務造成重復處理。這里我們在鎖內完成查找和狀態更新保證了單進程內多個線程不會重復領取。取數策略上我們支持按優先級排序。同一優先級的任務按創建時間先后處理這樣既滿足業務緊急度要求又不會讓低優先級任務無限積壓。4.3 Worker 消費端Worker 是一個后臺線程循環執行以下步驟調用claim_next()領取任務。如果當前沒有任務休眠 1 秒再繼續。調用 AI 處理邏輯。成功則調用complete()回填結果。失敗則調用fail_or_retry()記錄錯誤并決定是否重試。Worker 使用daemon線程的原因是不阻塞主進程退出。在真實生產環境中建議用進程管理工具或容器編排來管理多個 Worker而不是單線程。4.4 AI 處理服務抽象為了演示我們把“AI 處理”抽象成handle()方法。真實項目中這個方法內部可以調用 OpenAI 等大模型 API也可以調用本地部署的模型服務還可以執行一段工具腳本。這里有一個設計要點AI 調用一定要設置超時。模型接口的響應時間往往不穩定如果 Worker 因為沒有超時而卡在一個任務上后續所有任務都會被阻塞。我們可以在handle()中顯式設置 HTTP 客戶端超時或者用Signal強制中斷同步調用。5. 完整實現FastAPI 線程 Worker5.1 創建項目與安裝依賴首先創建項目目錄和虛擬環境。mkdir ai-inbox-assistant cd ai-inbox-assistant python -m venv .venv source .venv/bin/activate # Windows 使用 .venv\Scripts\activate創建requirements.txt并寫入以下依賴fastapi0.110 uvicorn[standard]0.29 pydantic2.0安裝依賴pip install -r requirements.txt5.2 定義數據模型文件路徑app/models.pyfrom datetime import datetime, timezone from enum import Enum from typing import Optional from pydantic import BaseModel, Field def utc_now() - datetime: 統一獲取當前 UTC 時間避免重復實現。 return datetime.now(timezone.utc) class TaskStatus(str, Enum): PENDING pending PROCESSING processing DONE done FAILED failed class TaskPriority(str, Enum): HIGH high NORMAL normal LOW low class InboxTask(BaseModel): id: str Field(default_factorylambda: __import__(uuid).uuid4().hex) sender: str unknown topic: str default content: str priority: TaskPriority TaskPriority.NORMAL status: TaskStatus TaskStatus.PENDING created_at: datetime Field(default_factoryutc_now) updated_at: datetime Field(default_factoryutc_now) attempts: int 0 result: Optional[str] None error: Optional[str] None def to_dict(self) - dict: 轉為可直接 JSON 序列化的字典兼容 pydantic v1/v2。 return { id: self.id, sender: self.sender, topic: self.topic, content: self.content, priority: self.priority.value, status: self.status.value, created_at: self.created_at.isoformat(), updated_at: self.updated_at.isoformat(), attempts: self.attempts, result: self.result, error: self.error, }這里重點解釋幾個設計決策。id使用uuid4().hex生成 32 位十六進制字符串足以避免并發提交時的 ID 沖突。也可以直接用str(uuid.uuid4())區別只是是否帶橫線。to_dict()方法統一負責序列化把枚舉值、時間對象轉換為普通字符串這樣接口層在返回響應時不需要關心底層 pydantic 版本差異。attempts字段默認 0表示任務還未被消費。5.3 實現 Inbox 核心邏輯文件路徑app/inbox.pyimport threading from typing import Dict, List, Optional from .models import InboxTask, TaskPriority, TaskStatus class Inbox: 線程安全的內存收件箱。 def __init__(self) - None: self._tasks: Dict[str, InboxTask] {} self._lock threading.RLock() staticmethod def _sort_key(task: InboxTask): 優先級高的任務排在前面相同優先級按創建時間判斷。 priority_order { TaskPriority.HIGH: 0, TaskPriority.NORMAL: 1, TaskPriority.LOW: 2, } return (priority_order.get(task.priority, 1), task.created_at) def add(self, content: str, sender: str unknown, topic: str default, priority: TaskPriority TaskPriority.NORMAL) - InboxTask: 向收件箱添加一個任務。 with self._lock: task InboxTask( sendersender, topictopic, contentcontent, prioritypriority, ) self._tasks[task.id] task return task def claim_next(self) - Optional[InboxTask]: 領取下一個待處理任務并將狀態改為 PROCESSING。 with self._lock: candidates [ task for task in self._tasks.values() if task.status TaskStatus.PENDING ] if not candidates: return None candidates.sort(keyself._sort_key) task candidates[0] task.status TaskStatus.PROCESSING task.attempts 1 task.updated_at __import__(app.models, fromlist[utc_now]).utc_now() return task def complete(self, task_id: str, result: str) - None: 處理成功后回填結果。 with self._lock: task self._tasks.get(task_id) if task is None: raise KeyError(ftask {task_id} not found) task.status TaskStatus.DONE task.result result task.error None task.updated_at __import__(app.models, fromlist[utc_now]).utc_now() def fail_or_retry(self, task_id: str, error: str, max_attempts: int 3) - None: 錯誤處理如果未超過最大執行次數則回到 PENDING否則標記 FAILED。 with self._lock: task self._tasks.get(task_id) if task is None: raise KeyError(ftask {task_id} not found) task.error error task.updated_at __import__(app.models, fromlist[utc_now]).utc_now() if task.attempts max_attempts: task.status TaskStatus.PENDING else: task.status TaskStatus.FAILED def get(self, task_id: str) - Optional[InboxTask]: with self._lock: return self._tasks.get(task_id) def list(self, status: Optional[TaskStatus] None) - List[InboxTask]: with self._lock: tasks list(self._tasks.values()) if status is not None: tasks [t for t in tasks if t.status status] tasks.sort(keylambda t: t.created_at) return tasks這里有一個實現細節需要注意claim_next返回的是任務對象本身而不是副本。這意味著 Worker 在拿到任務對象后即使 Inbox 鎖已經釋放其他線程讀取這個任務時也能看到PROCESSING狀態。這正是我們希望的效果因為它反映了真實的執行狀態。不過要注意由于我們直接修改任務對象的屬性如果 Worker 在執行任務時不小心修改了content等業務字段會產生臟數據。所以在設計約定上Worker 只允許通過complete和fail_or_retry修改任務狀態不要直接操作字段。5.4 實現 Worker 消費端文件路徑app/worker.pyimport threading import time from typing import Optional from .inbox import Inbox from .models import InboxTask class AIWorker(threading.Thread): 后臺消費線程從 Inbox 領取任務、調用模型、回填結果。 def __init__(self, inbox: Inbox, name: str ai-worker, poll_interval: float 1.0): super().__init__(namename, daemonTrue) self.inbox inbox self.poll_interval poll_interval self._stop_event threading.Event() def stop(self) - None: self._stop_event.set() def run(self) - None: while not self._stop_event.is_set(): task: Optional[InboxTask] self.inbox.claim_next() if task is None: self._stop_event.wait(self.poll_interval) continue try: result self.handle(task) self.inbox.complete(task.id, result) except Exception as exc: self.inbox.fail_or_retry(task.id, str(exc)) def handle(self, task: InboxTask) - str: 核心 AI 處理函數可替換為真實模型 API 調用。 # 模擬耗時操作生產環境替換成 LLM API / 本地模型推理 time.sleep(0.5) return f[{task.topic}] {task.content[:20]} 的 AI 摘要已生成Worker 中最容易踩坑的是異常處理邊界。handle()中任何異常都會觸發fail_or_retry()但這個邏輯需要與重試策略配合。如果任務是“永久性錯誤”例如內容包含非法字符導致模型拒絕處理重試多少次都不成功反而會浪費資源。所以在生產系統中handle()內部應該區分臨時錯誤和永久錯誤永久錯誤直接拋出特定異常由調用方判斷是一次性失敗還是繼續重試。這個示例中的poll_interval是 1 秒在演示環境可以接受。生產環境通常用消息隊列的阻塞讀取或者長輪詢避免無意義的輪詢開銷。5.5 編寫 FastAPI 接口文件路徑app/main.pyfrom contextlib import asynccontextmanager from typing import Optional from fastapi import FastAPI, HTTPException, Query from .inbox import Inbox from .models import InboxTask, TaskStatus from .worker import AIWorker inbox Inbox() worker: Optional[AIWorker] None asynccontextmanager async def lifespan(app: FastAPI): global worker worker AIWorker(inbox, nameai-worker) worker.start() yield if worker is not None: worker.stop() app FastAPI( titleAI Assistant with Inbox, description一個自帶收件箱的 AI 助手服務, version0.1.0, lifespanlifespan, ) class TaskCreateRequest: def __init__(self, content: str, sender: str unknown, topic: str default, priority: str normal): self.content content self.sender sender self.topic topic self.priority priority from pydantic import BaseModel class TaskCreateBody(BaseModel): content: str sender: str unknown topic: str default priority: str normal class TaskListResponse(BaseModel): items: list[dict] app.post(/inbox/tasks, status_code201) def create_task(body: TaskCreateBody) - dict: 提交一個新任務到收件箱。 from .models import TaskPriority try: priority TaskPriority(body.priority) except ValueError: raise HTTPException(status_code422, detailf無效優先級: {body.priority}) task inbox.add( contentbody.content, senderbody.sender, topicbody.topic, prioritypriority, ) return {task_id: task.id, status: task.status.value} app.get(/inbox/tasks) def list_tasks( status: Optional[TaskStatus] Query(defaultNone), sender: Optional[str] Query(defaultNone), ) - TaskListResponse: 列出收件箱任務支持按狀態和發送方篩選。 tasks inbox.list(statusstatus) if sender: tasks [t for t in tasks if t.sender sender] return TaskListResponse(items[t.to_dict() for t in tasks]) app.get(/inbox/tasks/{task_id}) def get_task(task_id: str) - dict: 查詢單個任務狀態和結果。 task: Optional[InboxTask] inbox.get(task_id) if task is None: raise HTTPException(status_code404, detail任務不存在) return task.to_dict()代碼里保留了TaskCreateRequest這個舊類其實是不需要的可以去掉。我在這里故意保留是因為實際開發中經常會有“寫多了再清理”的情況。正式代碼建議直接刪掉只保留 Pydantic 模型。接口設計有三個核心點。第一POST /inbox/tasks返回task_id而不是完整處理結果。客戶端拿到任務 ID 后可以通過GET /inbox/tasks/{task_id}輪詢結果。這是異步任務接口的標準做法。第二查詢接口支持按status和sender過濾方便業務側按狀態或來源查看收件箱內容。第三狀態枚舉通過 Query 參數接收時FastAPI 會自動做參數校驗。如果傳入非法狀態返回 422不需要我們手寫校驗邏輯。5.6 啟動服務并驗證現在啟動服務。uvicorn app.main:app --reload --port 8000看到如下輸出說明啟動成功INFO: Uvicorn running on http://127.0.0.1:8000 INFO: Application startup complete.FastAPI 會自動生成交互式文檔訪問http://127.0.0.1:8000/docs可以查看所有接口。6. 運行演示與結果說明6.1 提交任務打開另一個終端使用 curl 提交兩個測試任務一個高優先級一個普通優先級。curl -X POST http://127.0.0.1:8000/inbox/tasks \ -H Content-Type: application/json \ -d {content: 請總結本周運營數據, sender: cron, topic: report, priority: high}預期輸出{task_id:9f7b2f6d0c9a4e6f9c48e0c9ae62da21,status:pending}再提交一個普通任務curl -X POST http://127.0.0.1:8000/inbox/tasks \ -H Content-Type: application/json \ -d {content: 生成一封客戶回復郵件, sender: user, topic: email, priority: normal}6.2 查詢任務列表任務提交后立即查詢列表可能看到部分任務處于pending部分處于processing取決于 Worker 的處理速度。curl http://127.0.0.1:8000/inbox/tasks輸出示例{ items: [ { id: 9f7b2f6d0c9a4e6f9c48e0c9ae62da21, sender: cron, topic: report, content: 請總結本周運營數據, priority: high, status: done, created_at: 2025-01-01T10:00:0000:00, updated_at: 2025-01-01T10:00:0100:00, attempts: 1, result: [report] 請總結本周運營數據 的 AI 摘要已生成, error: null } ] }注意attempts字段已經變成 1說明 Worker 領取并處理過一次。result字段已經回填生成結果。6.3 查詢單個任務結果根據之前拿到的task_id查詢單個任務curl http://127.0.0.1:8000/inbox/tasks/9f7b2f6d0c9a4e6f9c48e0c9ae62da21輸出與列表中的單個項目一致。到這里一個最小的“帶收件箱的 AI 助手”已經可以跑通了。7. 常見問題與排查思路實際開發中你會遇到各種預期外的情況。下面整理了一些高頻問題。問題現象常見原因解決思路任務一直 pending狀態不變Worker 線程沒有啟動或提前退出檢查 lifespan 是否生效打印 Worker 啟動日志多個 Worker 重復處理同一任務領取任務和修改狀態不是原子操作在鎖/事務中完成狀態更新使用分布式鎖任務失敗后頻繁重試沒有區分臨時錯誤和永久錯誤定義可重試異常永久錯誤直接標記 FAILED服務重啟后任務丟失任務存儲在內存中引入 Redis Streams、SQLite、PostgreSQL 持久化API 返回 422狀態或優先級參數傳錯核對枚舉值大小寫參考 /docs 接口文檔模型調用超時導致 Worker 卡死外部接口沒有設置超時為 AI 調用設置超時時間并配合重試策略Uvicorn 啟動報 lifespan 錯誤代碼縮進或局部變量問題檢查 lifespan 上下文管理器結構啟動日志會顯示堆棧7.1 任務一直處于 pending 狀態出現這個現象首先檢查 Worker 是否在運行。在啟動日志中看不到 Worker 相關信息時往往是 lifespan 生命周期沒有掛載正確。FastAPI 舊版本常見做法是app.on_event(startup)新版開始推薦lifespan上下文管理器。如果你使用的是較老版本 FastAPI可以改回 startup 事件寫法但要注意不同版本的兼容性。還可以在 Worker 的run()方法最開始加一行打印日志比如print([worker] started)這樣能很快確認線程是否啟動。7.2 任務重復消費在單進程多線程模型中claim_next因為有RLock保護不會出現重復領取。但在多進程部署時每個進程都有自己的 Inbox 實例任務存儲不共享這時候問題會變成“各進程各處理各的”而不是重復消費同一個任務。真正的重復消費風險發生在任務存儲是共享的比如 Redis但領取時沒有用原子操作。解決方法有兩種在 Inbox 存儲層使用帶條件的原子更新例如 Redis Lua 腳本或 SQLUPDATE ... WHERE statuspending。在 Worker 處理結果回寫時使用冪等 ID 校驗防止重復寫入結果。對于 AI 任務重復消費不只是資源浪費還可能導致重復扣費和重復生成內容所以冪等設計要提前做。7.3 模型調用超時模型 API 是外部依賴它的延遲不可控。如果不設置超時一個慢請求可能讓 Worker 長期阻塞。常見做法有在網絡請求庫層面設置timeout比如requests.post(url, timeout(3, 30))。在多線程 Worker 中用Future.get(timeout...)控制單個任務執行時長。為任務設置最大執行時間超過閾值的任務重新進入隊列或直接標記失敗。8. 最佳實踐與工程建議演示代碼跑通后如果要在生產環境落地下面這些點非常關鍵。8.1 存儲層選型內存字典最明顯的缺點是重啟丟數據。生產環境推薦替換為以下方案之一。存儲方案適合場景優點注意點Redis Streams中高吞吐任務隊列天然支持消息持久化、消費者組需要處理 Stream 的消息過期和積壓Redis List BRPOP簡單任務隊列實現簡單阻塞讀取缺少消費者 ACK需要額外設計SQLite 狀態列低并發單機任務零額外依賴方便審計寫并發有限需要適當加鎖PostgreSQL SKIP LOCKED中大型系統支持事務和 SKIP LOCKED 避免重復消費需要數據庫連接池如果你已經有 RabbitMQ 或 Kafka 基礎設施也可以直接把它們作為任務隊列但要在消息體里保留task_id和完整錯誤信息。8.2 冪等與重試策略AI 調用通常涉及成本重試策略必須謹慎。建議按以下原則設計為每個任務生成全局唯一request_id發往模型服務時攜帶該 ID。網絡超時、限流、5xx 等臨時錯誤允許重試。內容不合法、參數錯誤等永久錯誤不要重試。設置最大嘗試次數默認為 3避免無限重試。使用指數退避策略比如第 1 次等 2 秒第 2 次等 4 秒第 3 次等 8 秒。在當前的fail_or_retry方法中最簡單的指數退避可以放在 Worker 內部實現重試前time.sleep(backoff)。8.3 超時與死信任務長時間處于PROCESSING狀態可能是 Worker 崩潰導致的任務“死亡”。生產環境需要引入“死信”機制。可以每隔一段時間掃描狀態為PROCESSING但updated_at超過 10 分鐘的任務將它們重新置為PENDING或標記為FAILED并記錄告警。這個掃描任務通常由定時調度器執行。8.4 安全與鑒權收件箱中可能包含敏感數據比如客戶郵件、業務報告文本。接口不能裸奔在公網上。建議在 FastAPI 中配置 API Key 或 OAuth2 鑒權。對任務內容加密存儲。查詢接口做權限校驗普通用戶只能查詢自己提交的任務不能查看他人的任務內容。記錄每個請求的操作人、時間和任務 ID以便審計。8.5 AI 調用成本控制當收件箱堆積大量任務時如果不做控制模型 API 賬單會很快飆升。控制成本可以從幾個方向入手任務入庫前進行內容長度限制和去重。對相同或近似內容做緩存命中后直接返回歷史結果。給 Worker 加速率限制防止瞬間請求過多導致模型 API 限流。流式讀取大文本時先做預處理