尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

【RAG实战】FastAPI + LlamaIndex Agent 流式对话踩坑实录

【RAG实战】FastAPI + LlamaIndex Agent 流式对话踩坑实录 FastAPI LlamaIndex Agent 流式对话踩坑实录取消标志残留、中间件缓冲、事件丢失全记录文章目录FastAPI LlamaIndex Agent 流式对话踩坑实录取消标志残留、中间件缓冲、事件丢失全记录前言一、项目背景事件类型二、问题一BaseHTTPMiddleware 导致流式响应被缓冲2.1 问题现象2.2 原因分析2.3 解决方案2.4 关键结论三、问题二asyncio.Queue 间接层导致事件丢失3.1 问题现象3.2 原因分析3.3 解决方案3.4 关键结论四、问题三Redis 取消标志残留导致新对话被误判核心问题4.1 问题现象4.2 原因分析4.3 解决方案4.4 为什么这样修复有效4.5 关键结论五、完整修复代码5.1 agent_service.py核心修改5.2 agent_controller.py传入 Redis 实例5.3 中间件修复纯 ASGI 实现六、排查经验总结6.1 流式响应排查清单6.2 关键日志模板6.3 最佳实践七、总结前言最近在基于 FastAPI LlamaIndex 构建 Agent 智能对话系统时遇到了一个典型的流式响应问题Agent 对话只返回start和end(cancelled)事件中间的内容全部丢失。经过多轮排查和修复最终定位到三个独立但相互关联的问题。本文完整记录问题现象、排查过程和解决方案希望能帮助遇到类似问题的同学少走弯路。一、项目背景Web 框架FastAPI UvicornAI 框架LlamaIndexFunctionAgent流式协议NDJSONapplication/x-ndjson每行一个 JSON 对象{type: xxx, payload: {...}}取消机制通过 Redis 标志位实现用户点击取消时设置agent_cancel:{session_id} 1Agent 事件循环中轮询检测事件类型事件类型说明start对话开始携带sessionIdmessageAI 回答内容增量片段sources检索来源切片列表end对话结束error异常信息二、问题一BaseHTTPMiddleware 导致流式响应被缓冲2.1 问题现象Agent 流式接口返回的StreamingResponse在客户端只能收到start事件后续的message、sources、end等事件全部丢失。2.2 原因分析项目中有两个中间件使用了 Starlette 的BaseHTTPMiddleware# 中间件 1上下文清理classContextCleanupMiddleware(BaseHTTPMiddleware):asyncdefdispatch(self,request,call_next):responseawaitcall_next(request)RequestContext.clear_all()returnresponse# 中间件 2响应头追加classApiResponseHeaderMiddleware(BaseHTTPMiddleware):asyncdefdispatch(self,request,call_next):responseawaitcall_next(request)# 追加自定义响应头...returnresponseBaseHTTPMiddleware的call_next()内部会拦截响应流导致StreamingResponse的 NDJSON chunk 被缓冲或提前终止。这是 Starlette 的已知问题BaseHTTPMiddleware不适合处理流式响应。2.3 解决方案将两个中间件从BaseHTTPMiddleware转为纯 ASGI 实现直接传递send函数不拦截响应流# 纯 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 中间件响应头追加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 关键结论凡是涉及流式响应SSE、NDJSON、WebSocket 等的 FastAPI 项目中间件必须使用纯 ASGI 实现不能使用BaseHTTPMiddleware。三、问题二asyncio.Queue 间接层导致事件丢失3.1 问题现象中间件修复后流式接口仍然只返回start事件。后端日志显示 Agent 正常执行RAG 检索命中、LLM 调用成功但事件无法送达客户端。3.2 原因分析之前的架构使用了asyncio.Queue 后台任务来解耦 Agent 执行和 HTTP 流Agent 事件 → queue.put() → queue.get() → yield → 中间件 → 客户端这个架构在 Agent 事件和 HTTP 流之间增加了间接层导致事件在 queue 传递过程中丢失或阻塞。3.3 解决方案彻底去掉 Queue改为直接流式传输# 修复后的架构Agent 事件 →yield→ 中间件 → 客户端classmethodasyncdefchat_stream(cls,db,request,user_id,app_redisNone):# ... setup 代码 ...# 创建 Agent 并运行agentAgentFactory.create_agent(...)handleragent.run(user_msgrequest.query,chat_historyllama_messages)yieldcls._ndjson(start,{sessionId:actual_session_id})# 直接从 Agent 事件流转发给客户端asyncforeventinhandler.stream_events():event_typetype(event).__name__ifevent_typeAgentStream:deltagetattr(event,delta,)ifdelta:full_answerdeltayieldcls._ndjson(message,{content:delta})elifevent_typeToolCallResult:# 提取来源...passyieldcls._ndjson(end,{})3.4 关键结论对于 LlamaIndex Agent 的流式场景直接从handler.stream_events()yield 事件给客户端即可不需要 Queue 间接层。Queue 架构适合需要复杂的生产者-消费者模式但在简单的流式转发场景中反而增加了不必要的复杂度和出错概率。四、问题三Redis 取消标志残留导致新对话被误判核心问题4.1 问题现象用户先取消一次对话然后在同一个会话中再次提问新对话立即返回{type: end, payload: {cancelled: true}}完全没有内容。4.2 原因分析取消机制使用 Redis 标志位# 取消接口awaitredis.set(fagent_cancel:{session_id},1,ex3600)# 事件循环中检测flagawaitredis.get(fagent_cancel:{session_id})ifflag1:yieldcls._ndjson(end,{cancelled:True})return问题 1标志残留取消标志的 TTL 是 3600 秒1 小时。用户取消后标志留在 Redis 中。同一会话的后续请求会检测到这个残留标志被误判为已取消。问题 2时序竞争即使在新对话开始时清除标志仍然存在时序竞争时间线 13.638 → 请求 2 开始 setup保存问题、构建记忆、创建 Agent... 14.228 → 用户点击取消仍在请求 2 的 setup 阶段 14.467 → 取消标志被写入 Redis 14.500 → 请求 2 的 setup 结束 14.504 → 请求 2 的事件循环检测到标志 → 误判如果把清除标志的代码放在 setup之前clear 在 13.638 执行 → 标志还不存在查了个空cancel 在 14.467 设置标志事件循环在 14.504 检测到标志 → 误判4.3 解决方案双保险策略将清除标志的代码移到 setup 之后、事件循环之前检测到取消后立即清除标志classmethodasyncdefchat_stream(cls,db,request,user_id,app_redisNone):actual_session_idrequest.session_idorstr(uuid.uuid4())# 获取 Redis 实例redisapp_redisorawaitcls._get_redis(db)# Setup 阶段 # 1. 保存用户问题# 2. 构建会话记忆# 3. 创建 Agent# 4. 启动 Agent 运行# yieldcls._ndjson(start,{sessionId:actual_session_id})# ★ 关键修复 1setup 完成后、事件循环开始前清除残留的取消标志ifredis:cancel_keyfagent_cancel:{actual_session_id}old_flagawaitredis.get(cancel_key)ifold_flag:awaitredis.delete(cancel_key)logger.info(f已清除残留取消标志:{cancel_key}(旧值{old_flag}))try:asyncforeventinhandler.stream_events():# 检查用户是否已取消ifredisandawaitcls._is_cancelled(redis,actual_session_id):# ★ 关键修复 2检测到取消后立即清除标志try:awaitredis.delete(fagent_cancel:{actual_session_id})exceptException:passlogger.info(f用户已取消对话: session_id{actual_session_id})awaitcls._save_cancelled_answer(...)yieldcls._ndjson(end,{cancelled:True})return# 处理事件...4.4 为什么这样修复有效修复后的时序时间线修复后 13.638 → 请求 2 开始 setup 14.228 → 用户点击取消 14.467 → 取消标志被写入 Redis 14.500 → setup 结束清除取消标志 ← 标志被清除 14.504 → 事件循环开始 → 没有标志 → 正常运行 ✓清除标志的代码从 setup之前移到 setup之后确保了即使取消请求在 setup 期间到达并设置了标志clear 也会在事件循环开始前把它清掉事件循环开始时看到的永远是干净的状态4.5 关键结论取消标志的生命周期应该与请求绑定而不是与会话绑定。每次新请求开始时清除旧标志每次取消被处理后也立即清除标志确保标志不会残留影响后续请求。五、完整修复代码5.1 agent_service.py核心修改classmethodasyncdefchat_stream(cls,db:AsyncSession,request,user_id:int,app_redisNone)-AsyncGenerator[str,None]:Agent 流式对话入口frommodule_rag.service.rag_chat_history_serviceimportRagChatHistoryService actual_session_idrequest.session_idorstr(uuid.uuid4())# 0. 获取 Redis 实例优先使用应用级 Redisredisapp_redisifnotredis:try:redisawaitcls._get_redis(db)exceptException:redisNone# Setup 阶段保存问题、构建记忆、创建 Agent 等# ... 省略 setup 代码 ...# 启动 AgentagentAgentFactory.create_agent(...)handleragent.run(user_msgrequest.query,chat_historyllama_messages)yieldcls._ndjson(start,{sessionId:actual_session_id})# ★ 修复setup 完成后、事件循环前清除残留取消标志ifredis:try:cancel_keyfagent_cancel:{actual_session_id}old_flagawaitredis.get(cancel_key)ifold_flag:awaitredis.delete(cancel_key)logger.info(f已清除残留取消标志:{cancel_key}(旧值{old_flag}))exceptExceptionasclear_err:logger.warning(f清除取消标志失败:{clear_err})try:asyncforeventinhandler.stream_events():ifredisandawaitcls._is_cancelled(redis,actual_session_id):# ★ 修复检测到取消后立即清除标志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,{})# 后台保存回答fire-and-forgetcls._post_chat_tasks(...)exceptExceptionasstream_err:yieldcls._ndjson(error,{message:str(stream_err)})5.2 agent_controller.py传入 Redis 实例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# 获取应用级 Redis 实例与取消接口使用同一个连接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 中间件修复纯 ASGI 实现# context_middleware.pyclassContextCleanupMiddleware:上下文清理中间件纯 ASGI 实现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()六、排查经验总结6.1 流式响应排查清单检查中间件链所有BaseHTTPMiddleware都可能缓冲流式响应检查事件传递路径Agent → Queue → yield → 中间件 → 客户端每一环都可能丢失事件检查 Redis 状态取消标志、会话标志等是否残留检查时序异步场景下的竞争条件特别是 setup 阶段和事件循环之间的时间窗口6.2 关键日志模板# 请求开始logger.info(f[AGENT][STREAM] 开始: session_id{actual_session_id})# 清除标志logger.info(f[AGENT][STREAM] 已清除残留取消标志:{cancel_key}(旧值{old_flag}))# 事件循环logger.info(f[AGENT][STREAM] 事件 #{event_count}:{event_type})# 取消检测logger.info(f[AGENT] _is_cancelled: key{cancel_key}, flag{flag})6.3 最佳实践场景推荐做法避免做法流式中间件纯 ASGI 实现BaseHTTPMiddlewareAgent 流式传输直接yield事件asyncio.Queue间接层取消标志管理请求开始时清除取消后立即清除依赖 TTL 自动过期Redis 实例共享使用request.app.state.redis每次创建新连接七、总结本文记录了 FastAPI LlamaIndex Agent 流式对话的三个典型问题BaseHTTPMiddleware 缓冲流式响应→ 转为纯 ASGI 中间件Queue 间接层导致事件丢失→ 直接流式传输Redis 取消标志残留 时序竞争→ 双保险清除策略这三个问题独立存在但相互关联任何一个都可能导致流式响应异常。排查时需要从中间件链、事件传递路径、Redis 状态、时序竞争等多个维度综合分析。希望本文能帮助遇到类似问题的同学快速定位和解决。如果觉得有用欢迎点赞、收藏、转发
返回列表