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