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