戰(zhàn)】FastAPI + LlamaIndex Agent 流式對(duì)話踩坑實(shí)錄)
FastAPI LlamaIndex Agent 流式對(duì)話踩坑實(shí)錄取消標(biāo)志殘留、中間件緩沖、事件丟失全記錄文章目錄FastAPI LlamaIndex Agent 流式對(duì)話踩坑實(shí)錄取消標(biāo)志殘留、中間件緩沖、事件丟失全記錄前言一、項(xiàng)目背景事件類型二、問(wèn)題一BaseHTTPMiddleware 導(dǎo)致流式響應(yīng)被緩沖2.1 問(wèn)題現(xiàn)象2.2 原因分析2.3 解決方案2.4 關(guān)鍵結(jié)論三、問(wèn)題二asyncio.Queue 間接層導(dǎo)致事件丟失3.1 問(wèn)題現(xiàn)象3.2 原因分析3.3 解決方案3.4 關(guān)鍵結(jié)論四、問(wèn)題三Redis 取消標(biāo)志殘留導(dǎo)致新對(duì)話被誤判核心問(wèn)題4.1 問(wèn)題現(xiàn)象4.2 原因分析4.3 解決方案4.4 為什么這樣修復(fù)有效4.5 關(guān)鍵結(jié)論五、完整修復(fù)代碼5.1 agent_service.py核心修改5.2 agent_controller.py傳入 Redis 實(shí)例5.3 中間件修復(fù)純 ASGI 實(shí)現(xiàn)六、排查經(jīng)驗(yàn)總結(jié)6.1 流式響應(yīng)排查清單6.2 關(guān)鍵日志模板6.3 最佳實(shí)踐七、總結(jié)前言最近在基于 FastAPI LlamaIndex 構(gòu)建 Agent 智能對(duì)話系統(tǒng)時(shí)遇到了一個(gè)典型的流式響應(yīng)問(wèn)題Agent 對(duì)話只返回start和end(cancelled)事件中間的內(nèi)容全部丟失。經(jīng)過(guò)多輪排查和修復(fù)最終定位到三個(gè)獨(dú)立但相互關(guān)聯(lián)的問(wèn)題。本文完整記錄問(wèn)題現(xiàn)象、排查過(guò)程和解決方案希望能幫助遇到類似問(wèn)題的同學(xué)少走彎路。一、項(xiàng)目背景Web 框架FastAPI UvicornAI 框架LlamaIndexFunctionAgent流式協(xié)議NDJSONapplication/x-ndjson每行一個(gè) JSON 對(duì)象{type: xxx, payload: {...}}取消機(jī)制通過(guò) Redis 標(biāo)志位實(shí)現(xiàn)用戶點(diǎn)擊取消時(shí)設(shè)置agent_cancel:{session_id} 1Agent 事件循環(huán)中輪詢檢測(cè)事件類型事件類型說(shuō)明start對(duì)話開(kāi)始攜帶sessionIdmessageAI 回答內(nèi)容增量片段sources檢索來(lái)源切片列表end對(duì)話結(jié)束error異常信息二、問(wèn)題一BaseHTTPMiddleware 導(dǎo)致流式響應(yīng)被緩沖2.1 問(wèn)題現(xiàn)象Agent 流式接口返回的StreamingResponse在客戶端只能收到start事件后續(xù)的message、sources、end等事件全部丟失。2.2 原因分析項(xiàng)目中有兩個(gè)中間件使用了 Starlette 的BaseHTTPMiddleware# 中間件 1上下文清理classContextCleanupMiddleware(BaseHTTPMiddleware):asyncdefdispatch(self,request,call_next):responseawaitcall_next(request)RequestContext.clear_all()returnresponse# 中間件 2響應(yīng)頭追加classApiResponseHeaderMiddleware(BaseHTTPMiddleware):asyncdefdispatch(self,request,call_next):responseawaitcall_next(request)# 追加自定義響應(yīng)頭...returnresponseBaseHTTPMiddleware的call_next()內(nèi)部會(huì)攔截響應(yīng)流導(dǎo)致StreamingResponse的 NDJSON chunk 被緩沖或提前終止。這是 Starlette 的已知問(wèn)題BaseHTTPMiddleware不適合處理流式響應(yīng)。2.3 解決方案將兩個(gè)中間件從BaseHTTPMiddleware轉(zhuǎn)為純 ASGI 實(shí)現(xiàn)直接傳遞send函數(shù)不攔截響應(yīng)流# 純 ASGI 中間件上下文清理classContextCleanupMiddleware:def__init__(self,app)-None:self.appappasyncdef__call__(self,scope,receive,send)-None:ifscope[type]!http:awaitself.app(scope,receive,send)returntry:awaitself.app(scope,receive,send)finally:RequestContext.clear_all()# 純 ASGI 中間件響應(yīng)頭追加classApiResponseHeaderMiddleware:def__init__(self,app)-None:self.appappasyncdef__call__(self,scope,receive,send)-None:ifscope[type]!http:awaitself.app(scope,receive,send)returnrequestRequest(scope,receive)api_response_headersgetattr(request.state,api_response_headers,None)ifnotapi_response_headers:awaitself.app(scope,receive,send)returnasyncdefsend_with_headers(message)-None:ifmessage[type]http.response.start:headersdict(message.get(headers,[]))forkey,valueinapi_response_headers.items():headers[key.encode(utf-8)]value.encode(utf-8)message{**message,headers:list(headers.items())}awaitsend(message)awaitself.app(scope,receive,send_with_headers)2.4 關(guān)鍵結(jié)論凡是涉及流式響應(yīng)SSE、NDJSON、WebSocket 等的 FastAPI 項(xiàng)目中間件必須使用純 ASGI 實(shí)現(xiàn)不能使用BaseHTTPMiddleware。三、問(wèn)題二asyncio.Queue 間接層導(dǎo)致事件丟失3.1 問(wèn)題現(xiàn)象中間件修復(fù)后流式接口仍然只返回start事件。后端日志顯示 Agent 正常執(zhí)行RAG 檢索命中、LLM 調(diào)用成功但事件無(wú)法送達(dá)客戶端。3.2 原因分析之前的架構(gòu)使用了asyncio.Queue 后臺(tái)任務(wù)來(lái)解耦 Agent 執(zhí)行和 HTTP 流Agent 事件 → queue.put() → queue.get() → yield → 中間件 → 客戶端這個(gè)架構(gòu)在 Agent 事件和 HTTP 流之間增加了間接層導(dǎo)致事件在 queue 傳遞過(guò)程中丟失或阻塞。3.3 解決方案徹底去掉 Queue改為直接流式傳輸# 修復(fù)后的架構(gòu)Agent 事件 →yield→ 中間件 → 客戶端classmethodasyncdefchat_stream(cls,db,request,user_id,app_redisNone):# ... setup 代碼 ...# 創(chuàng)建 Agent 并運(yùn)行agentAgentFactory.create_agent(...)handleragent.run(user_msgrequest.query,chat_historyllama_messages)yieldcls._ndjson(start,{sessionId:actual_session_id})# 直接從 Agent 事件流轉(zhuǎn)發(fā)給客戶端asyncforeventinhandler.stream_events():event_typetype(event).__name__ifevent_typeAgentStream:deltagetattr(event,delta,)ifdelta:full_answerdeltayieldcls._ndjson(message,{content:delta})elifevent_typeToolCallResult:# 提取來(lái)源...passyieldcls._ndjson(end,{})3.4 關(guān)鍵結(jié)論對(duì)于 LlamaIndex Agent 的流式場(chǎng)景直接從handler.stream_events()yield 事件給客戶端即可不需要 Queue 間接層。Queue 架構(gòu)適合需要復(fù)雜的生產(chǎn)者-消費(fèi)者模式但在簡(jiǎn)單的流式轉(zhuǎn)發(fā)場(chǎng)景中反而增加了不必要的復(fù)雜度和出錯(cuò)概率。四、問(wèn)題三Redis 取消標(biāo)志殘留導(dǎo)致新對(duì)話被誤判核心問(wèn)題4.1 問(wèn)題現(xiàn)象用戶先取消一次對(duì)話然后在同一個(gè)會(huì)話中再次提問(wèn)新對(duì)話立即返回{type: end, payload: {cancelled: true}}完全沒(méi)有內(nèi)容。4.2 原因分析取消機(jī)制使用 Redis 標(biāo)志位# 取消接口awaitredis.set(fagent_cancel:{session_id},1,ex3600)# 事件循環(huán)中檢測(cè)flagawaitredis.get(fagent_cancel:{session_id})ifflag1:yieldcls._ndjson(end,{cancelled:True})return問(wèn)題 1標(biāo)志殘留取消標(biāo)志的 TTL 是 3600 秒1 小時(shí)。用戶取消后標(biāo)志留在 Redis 中。同一會(huì)話的后續(xù)請(qǐng)求會(huì)檢測(cè)到這個(gè)殘留標(biāo)志被誤判為已取消。問(wèn)題 2時(shí)序競(jìng)爭(zhēng)即使在新對(duì)話開(kāi)始時(shí)清除標(biāo)志仍然存在時(shí)序競(jìng)爭(zhēng)時(shí)間線 13.638 → 請(qǐng)求 2 開(kāi)始 setup保存問(wèn)題、構(gòu)建記憶、創(chuàng)建 Agent... 14.228 → 用戶點(diǎn)擊取消仍在請(qǐng)求 2 的 setup 階段 14.467 → 取消標(biāo)志被寫(xiě)入 Redis 14.500 → 請(qǐng)求 2 的 setup 結(jié)束 14.504 → 請(qǐng)求 2 的事件循環(huán)檢測(cè)到標(biāo)志 → 誤判如果把清除標(biāo)志的代碼放在 setup之前clear 在 13.638 執(zhí)行 → 標(biāo)志還不存在查了個(gè)空cancel 在 14.467 設(shè)置標(biāo)志事件循環(huán)在 14.504 檢測(cè)到標(biāo)志 → 誤判4.3 解決方案雙保險(xiǎn)策略將清除標(biāo)志的代碼移到 setup 之后、事件循環(huán)之前檢測(cè)到取消后立即清除標(biāo)志classmethodasyncdefchat_stream(cls,db,request,user_id,app_redisNone):actual_session_idrequest.session_idorstr(uuid.uuid4())# 獲取 Redis 實(shí)例redisapp_redisorawaitcls._get_redis(db)# Setup 階段 # 1. 保存用戶問(wèn)題# 2. 構(gòu)建會(huì)話記憶# 3. 創(chuàng)建 Agent# 4. 啟動(dòng) Agent 運(yùn)行# yieldcls._ndjson(start,{sessionId:actual_session_id})# ★ 關(guān)鍵修復(fù) 1setup 完成后、事件循環(huán)開(kāi)始前清除殘留的取消標(biāo)志ifredis:cancel_keyfagent_cancel:{actual_session_id}old_flagawaitredis.get(cancel_key)ifold_flag:awaitredis.delete(cancel_key)logger.info(f已清除殘留取消標(biāo)志:{cancel_key}(舊值{old_flag}))try:asyncforeventinhandler.stream_events():# 檢查用戶是否已取消ifredisandawaitcls._is_cancelled(redis,actual_session_id):# ★ 關(guān)鍵修復(fù) 2檢測(cè)到取消后立即清除標(biāo)志try:awaitredis.delete(fagent_cancel:{actual_session_id})exceptException:passlogger.info(f用戶已取消對(duì)話: session_id{actual_session_id})awaitcls._save_cancelled_answer(...)yieldcls._ndjson(end,{cancelled:True})return# 處理事件...4.4 為什么這樣修復(fù)有效修復(fù)后的時(shí)序時(shí)間線修復(fù)后 13.638 → 請(qǐng)求 2 開(kāi)始 setup 14.228 → 用戶點(diǎn)擊取消 14.467 → 取消標(biāo)志被寫(xiě)入 Redis 14.500 → setup 結(jié)束清除取消標(biāo)志 ← 標(biāo)志被清除 14.504 → 事件循環(huán)開(kāi)始 → 沒(méi)有標(biāo)志 → 正常運(yùn)行 ?清除標(biāo)志的代碼從 setup之前移到 setup之后確保了即使取消請(qǐng)求在 setup 期間到達(dá)并設(shè)置了標(biāo)志clear 也會(huì)在事件循環(huán)開(kāi)始前把它清掉事件循環(huán)開(kāi)始時(shí)看到的永遠(yuǎn)是干凈的狀態(tài)4.5 關(guān)鍵結(jié)論取消標(biāo)志的生命周期應(yīng)該與請(qǐng)求綁定而不是與會(huì)話綁定。每次新請(qǐng)求開(kāi)始時(shí)清除舊標(biāo)志每次取消被處理后也立即清除標(biāo)志確保標(biāo)志不會(huì)殘留影響后續(xù)請(qǐng)求。五、完整修復(fù)代碼5.1 agent_service.py核心修改classmethodasyncdefchat_stream(cls,db:AsyncSession,request,user_id:int,app_redisNone)-AsyncGenerator[str,None]:Agent 流式對(duì)話入口frommodule_rag.service.rag_chat_history_serviceimportRagChatHistoryService actual_session_idrequest.session_idorstr(uuid.uuid4())# 0. 獲取 Redis 實(shí)例優(yōu)先使用應(yīng)用級(jí) Redisredisapp_redisifnotredis:try:redisawaitcls._get_redis(db)exceptException:redisNone# Setup 階段保存問(wèn)題、構(gòu)建記憶、創(chuàng)建 Agent 等# ... 省略 setup 代碼 ...# 啟動(dòng) AgentagentAgentFactory.create_agent(...)handleragent.run(user_msgrequest.query,chat_historyllama_messages)yieldcls._ndjson(start,{sessionId:actual_session_id})# ★ 修復(fù)setup 完成后、事件循環(huán)前清除殘留取消標(biāo)志ifredis:try:cancel_keyfagent_cancel:{actual_session_id}old_flagawaitredis.get(cancel_key)ifold_flag:awaitredis.delete(cancel_key)logger.info(f已清除殘留取消標(biāo)志:{cancel_key}(舊值{old_flag}))exceptExceptionasclear_err:logger.warning(f清除取消標(biāo)志失敗:{clear_err})try:asyncforeventinhandler.stream_events():ifredisandawaitcls._is_cancelled(redis,actual_session_id):# ★ 修復(fù)檢測(cè)到取消后立即清除標(biāo)志try:awaitredis.delete(fagent_cancel:{actual_session_id})exceptException:passawaitcls._save_cancelled_answer(...)yieldcls._ndjson(end,{cancelled:True})returnevent_typetype(event).__name__ifevent_typeAgentStream:deltagetattr(event,delta,)ifdelta:yieldcls._ndjson(message,{content:delta})yieldcls._ndjson(end,{})# 后臺(tái)保存回答fire-and-forgetcls._post_chat_tasks(...)exceptExceptionasstream_err:yieldcls._ndjson(error,{message:str(stream_err)})5.2 agent_controller.py傳入 Redis 實(shí)例agent_controller.post(/chat/stream)asyncdefagent_chat_stream(request:Request,chat_req:AgentChatRequestModel,query_db:Annotated[AsyncSession,DBSessionDependency()],current_user:Annotated[CurrentUserModel,CurrentUserDependency()],)-StreamingResponse:user_idcurrent_user.user.user_id# 獲取應(yīng)用級(jí) Redis 實(shí)例與取消接口使用同一個(gè)連接app_redisrequest.app.state.redisifhasattr(request.app.state,redis)elseNoneevent_streamAgentService.chat_stream(query_db,chat_req,user_id,app_redisapp_redis)returnStreamingResponse(contentevent_stream,media_typeapplication/x-ndjson)5.3 中間件修復(fù)純 ASGI 實(shí)現(xiàn)# context_middleware.pyclassContextCleanupMiddleware:上下文清理中間件純 ASGI 實(shí)現(xiàn)def__init__(self,app)-None:self.appappasyncdef__call__(self,scope,receive,send)-None:ifscope[type]!http:awaitself.app(scope,receive,send)returntry:awaitself.app(scope,receive,send)finally:RequestContext.clear_all()六、排查經(jīng)驗(yàn)總結(jié)6.1 流式響應(yīng)排查清單檢查中間件鏈所有BaseHTTPMiddleware都可能緩沖流式響應(yīng)檢查事件傳遞路徑Agent → Queue → yield → 中間件 → 客戶端每一環(huán)都可能丟失事件檢查 Redis 狀態(tài)取消標(biāo)志、會(huì)話標(biāo)志等是否殘留檢查時(shí)序異步場(chǎng)景下的競(jìng)爭(zhēng)條件特別是 setup 階段和事件循環(huán)之間的時(shí)間窗口6.2 關(guān)鍵日志模板# 請(qǐng)求開(kāi)始logger.info(f[AGENT][STREAM] 開(kāi)始: session_id{actual_session_id})# 清除標(biāo)志logger.info(f[AGENT][STREAM] 已清除殘留取消標(biāo)志:{cancel_key}(舊值{old_flag}))# 事件循環(huán)logger.info(f[AGENT][STREAM] 事件 #{event_count}:{event_type})# 取消檢測(cè)logger.info(f[AGENT] _is_cancelled: key{cancel_key}, flag{flag})6.3 最佳實(shí)踐場(chǎng)景推薦做法避免做法流式中間件純 ASGI 實(shí)現(xiàn)BaseHTTPMiddlewareAgent 流式傳輸直接yield事件asyncio.Queue間接層取消標(biāo)志管理請(qǐng)求開(kāi)始時(shí)清除取消后立即清除依賴 TTL 自動(dòng)過(guò)期Redis 實(shí)例共享使用request.app.state.redis每次創(chuàng)建新連接七、總結(jié)本文記錄了 FastAPI LlamaIndex Agent 流式對(duì)話的三個(gè)典型問(wèn)題BaseHTTPMiddleware 緩沖流式響應(yīng)→ 轉(zhuǎn)為純 ASGI 中間件Queue 間接層導(dǎo)致事件丟失→ 直接流式傳輸Redis 取消標(biāo)志殘留 時(shí)序競(jìng)爭(zhēng)→ 雙保險(xiǎn)清除策略這三個(gè)問(wèn)題獨(dú)立存在但相互關(guān)聯(lián)任何一個(gè)都可能導(dǎo)致流式響應(yīng)異常。排查時(shí)需要從中間件鏈、事件傳遞路徑、Redis 狀態(tài)、時(shí)序競(jìng)爭(zhēng)等多個(gè)維度綜合分析。希望本文能幫助遇到類似問(wèn)題的同學(xué)快速定位和解決。如果覺(jué)得有用歡迎點(diǎn)贊、收藏、轉(zhuǎn)發(fā)