LangChain 与 FastAPI 集成:用流式 SSE 将 Agent 封装为 REST API

发布时间:2026/7/26 18:45:59
LangChain 与 FastAPI 集成:用流式 SSE 将 Agent 封装为 REST API LangChain 与 FastAPI 集成用流式 SSE 将 Agent 封装为 REST API一、深度引言与场景痛点大家好我是赵咕咕。我发现一个很有意思的现象很多工程师花了两周打磨 Agent 的逻辑——工具调用、Prompt 优化、记忆管理都做得很精细了——然后到了怎么让前端调用这一步随意起个 Flask 单线程就跑或者直接扔一个同步的POST /chat接口完事。这有两个问题。第一个是体验问题用户发一条消息等 15 秒页面一动不动然后咣当一下返回整个结果。第二个是运维问题Flask 单线程根本扛不住并发Agent 一次工具调用链跑 20 秒下一个请求就等着排队。Agent 的服务化不是在模型外面包一层 HTTP。它需要处理流式输出SSE、会话管理、并发控制、超时优雅关闭。这篇文章我用 FastAPI LangChain SSE 把这些问题逐个解决给出一个可以直上生产的环境。二、底层机制与原理深度剖析2.1 为什么 Agent 的 API 封装比普通 LLM 调用复杂普通 LLM 调用的流程是请求进来 → 调 OpenAI API → 流式返回。一条线没有分支。但 Agent 的推理过程是一棵树用户发帮我查一下今天的天气预报然后发个 Slack 消息给团队Agent 思考 → 决定调天气 API → 拿到结果 → 再思考 → 决定调 Slack API → 发送 → 总结回复中间可能因为工具调用失败而重试可能因为信息不足而反问用户如果把这一整棵树都跑完再返回结果用户的等待时间 所有工具调用的总延迟。SSE 的价值在于每一步的输出都实时推送。用户在等待工具调用时能看到 Agent 在思考什么、调了什么工具、拿到了什么结果——这个透明度极大提升了体验。2.2 SSE 流式推送的完整架构这个时序图揭示了几个关键设计事件类型分离前端需要知道每一步是什么——是 token 流、工具调用开始、工具调用结束、还是整个会话结束。不同事件类型前端可以做不同的 UI 渲染。会话透传每次请求都携带session_id服务端根据它加载历史对话和工具调用记录。这样 Agent 的记忆才能跨请求保持。异步非阻塞整个链路从 HTTP 接收到 Agent 推理再到 SSE 推送全程async/await不占用线程。2.3 FastAPI SSE 的技术要点SSEServer-Sent Events是用Content-Type: text/event-stream响应头声明的长连接。FastAPI 的StreamingResponse天然支持。几个容易踩的坑连接保持Nginx 默认 60s 超时Agent 推理可能超过这个时间。需要调大proxy_read_timeout。前端断连SSE 是基于 HTTP 的长连接。前端关掉页面或者刷新连接断开。服务端需要通过asyncio.CancelledError感知并优雅终止 Agent 推理。并发模型每个 SSE 连接是一个独立的 asyncio Task。FastAPI 的 event loop 可以管理上千个并发连接但要确保 Agent 操作的都是 async 的否则阻塞 event loop。三、生产级代码实现import asyncio import json import logging import uuid from contextlib import asynccontextmanager from typing import Any, AsyncIterator from fastapi import FastAPI, HTTPException from fastapi.responses import StreamingResponse from pydantic import BaseModel, Field from langchain_openai import ChatOpenAI from langchain.agents import AgentExecutor, create_openai_tools_agent from langchain_core.prompts import ChatPromptTemplate, MessagesPlaceholder from langchain_core.messages import HumanMessage, AIMessage from langchain_core.tools import tool logger logging.getLogger(__name__) # ── 数据模型 ─────────────────────────────────────────── class ChatRequest(BaseModel): message: str Field(..., min_length1, description用户输入) session_id: str Field(default_factorylambda: uuid.uuid4().hex[:12]) # ── 会话管理 ─────────────────────────────────────────── class SessionManager: 管理 Agent 会话的创建、查找和历史维护。 def __init__(self, max_history: int 20): self._sessions: dict[str, list[Any]] {} self._max_history max_history def get_or_create(self, session_id: str) - list[Any]: if session_id not in self._sessions: self._sessions[session_id] [] return self._sessions[session_id] def append(self, session_id: str, message: Any) - None: history self._sessions.setdefault(session_id, []) history.append(message) # 防止历史过长超出模型上下文窗口 if len(history) self._max_history * 2: self._sessions[session_id] history[-self._max_history * 2:] def cleanup(self, session_id: str) - None: self._sessions.pop(session_id, None) # ── 工具定义示例 ─────────────────────────────────── tool async def get_weather(city: str) - str: 查询指定城市的天气信息。 # 实际项目里调真实 API weather_data {北京: 晴, 25°C, 上海: 多云, 28°C, 深圳: 阵雨, 30°C} await asyncio.sleep(0.5) # 模拟网络延迟 return weather_data.get(city, f未找到{city}的天气数据) tool async def send_slack(channel: str, message: str) - str: 向 Slack 频道发送消息。 await asyncio.sleep(0.3) return f消息已发送到 #{channel}: {message[:50]}... # ── Agent 工厂 ───────────────────────────────────────── class AgentFactory: 创建带工具集的 Agent。 def __init__(self, model: str gpt-4o): self._llm ChatOpenAI( modelmodel, temperature0, streamingTrue, # 关键启用流式 ) self._tools [get_weather, send_slack] self._prompt ChatPromptTemplate.from_messages([ (system, 你是一个智能助手。使用工具来回答问题逐步推理。), MessagesPlaceholder(variable_namechat_history, optionalTrue), (human, {input}), MessagesPlaceholder(variable_nameagent_scratchpad), ]) def create(self) - AgentExecutor: agent create_openai_tools_agent(self._llm, self._tools, self._prompt) return AgentExecutor( agentagent, toolsself._tools, verboseFalse, max_iterations10, handle_parsing_errorsTrue, ) # ── SSE 事件序列化 ───────────────────────────────────── class SSEEvent: SSE 事件格式化。 staticmethod def format(event_type: str, data: dict[str, Any]) - str: payload json.dumps({type: event_type, **data}, ensure_asciiFalse) return fdata: {payload}\n\n staticmethod def done() - str: return data: {\type\: \done\}\n\n staticmethod def error(message: str) - str: payload json.dumps({type: error, message: message}, ensure_asciiFalse) return fdata: {payload}\n\n # ── FastAPI 应用 ─────────────────────────────────────── asynccontextmanager async def lifespan(app: FastAPI): 应用生命周期管理。 app.state.sessions SessionManager(max_history20) app.state.agent_factory AgentFactory() logger.info(Agent API 服务已启动) yield logger.info(Agent API 服务正在关闭) app FastAPI(titleAgent API, lifespanlifespan) app.post(/chat) async def chat(req: ChatRequest) - StreamingResponse: 流式 Agent 对话接口。 async def event_stream() - AsyncIterator[str]: session_id req.session_id history app.state.sessions.get_or_create(session_id) try: # 1. 发送会话就绪事件 yield SSEEvent.format(session_ready, {session_id: session_id}) # 2. 创建 Agent agent app.state.agent_factory.create() # 3. 用 astream_events 获取精细事件流 # 注意: astream_events 在 astream_log 之后版本可能变化 async for event in agent.astream_events( { input: req.message, chat_history: history, }, versionv2, ): kind event[event] if kind on_chat_model_stream: # LLM 逐 token 推送 chunk event[data][chunk] if hasattr(chunk, content) and chunk.content: yield SSEEvent.format(token, {content: chunk.content}) elif kind on_tool_start: # 工具开始调用 yield SSEEvent.format(tool_start, { tool: event.get(name, unknown), input: event[data].get(input, {}), }) elif kind on_tool_end: # 工具调用结束 yield SSEEvent.format(tool_end, { tool: event.get(name, unknown), output: str(event[data].get(output, ))[:500], }) elif kind on_chain_end and event.get(name) AgentExecutor: # Agent 推理完成保存历史 output event[data].get(output, ) app.state.sessions.append(session_id, HumanMessage(contentreq.message)) app.state.sessions.append(session_id, AIMessage(contentoutput)) yield SSEEvent.done() except asyncio.CancelledError: # 前端断开连接优雅退出 logger.info(SSE 连接被客户端取消: session%s, session_id) yield SSEEvent.error(连接已取消) except Exception as e: logger.exception(Agent 推理失败: session%s, session_id) yield SSEEvent.error(f内部错误: {str(e)[:200]}) return StreamingResponse( event_stream(), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, # 禁用 Nginx 缓冲 }, ) app.get(/sessions/{session_id}/history) async def get_history(session_id: str): 获取会话历史。 history app.state.sessions.get_or_create(session_id) return { session_id: session_id, messages: [ {role: user if isinstance(m, HumanMessage) else assistant, content: m.content} for m in history ], } app.delete(/sessions/{session_id}) async def clear_session(session_id: str): 清除会话。 app.state.sessions.cleanup(session_id) return {status: cleared, session_id: session_id} if __name__ __main__: import uvicorn uvicorn.run(app, host0.0.0.0, port8000)关键设计点astream_events是精髓LangChain 的astream返回的是每一步最终输出但astream_events返回的是每一步的内部事件。token 流、工具调用开始/结束、链结束——这些事件让前端能做精细的 UI 渲染loading spinner、工具调用卡片、流式文本。X-Accel-Buffering: no如果前面有 Nginx 反代不加这个头 Nginx 会把 SSE 的输出缓冲起来等 Agent 全部跑完才一次性发给前端——流式就白做了。asyncio.CancelledError处理SSE 是长连接前端断连时 asyncio Task 会收到取消信号。捕获它做清理而不是报错。会话历史用 AIMessage/HumanMessage 原生类型LangChain Agent 的chat_history参数需要 LangChain 的 Message 类型不要自己搞一套数据结构转换。四、边界分析与架构权衡4.1 SSE vs WebSocket vs Polling方案优势劣势适用场景SSEHTTP 协议原生走 CDN 无压力自动重连单向推送需额外 POST 发消息Agent 对话一次请求多次推送WebSocket双向通信低延迟需要自己管理重连、心跳代理配置复杂实时协作、多轮交互频繁Polling最简单兼容性好浪费带宽延迟高不需要实时反馈的场景Agent 对话这个场景SSE 是最合适的——前端 POST 一次消息服务端一路流式推回结果。没有双向通信的需求WebSocket 的复杂度是多余的。4.2 生产级的并发与压测考量一个 Agent 推理会占用 LLM API 的连接和本地 asyncio Task。压测时要关注的指标最大并发数取决于 LLM API 的 rate limit 和本地 CPU/内存。一般单机 50-100 并发 Agent 对话是比较安全的范围。背压处理当并发满时新请求应该返回 429Too Many Requests而不是排队等。前端看到 429 可以提示用户稍后重试。Token 级别的速率限制除了并发数还要限制每用户每分钟的 Token 消耗防止单个用户打爆预算。4.3 会话持久化上面的实现用了内存字典_sessions: dict。单机部署完全够用但要做持久化的话有几种选择Redis存会话历史 JSON设置 TTL。适合多实例部署。SQLite/Postgres存结构化消息记录方便后续做分析和评估。LangChain 的 BaseChatMessageHistoryLangChain 有内置的 Redis/Postgres ChatMessageHistory 实现无缝对接。4.4 超时与资源释放Agent 推理可能因为工具调用卡住而无限等待。加超时是必须的try: async for event in agent.astream_events(...): ... except asyncio.TimeoutError: yield SSEEvent.error(推理超时请简化问题重试)建议对单次 Agent 推理设置 120 秒超时同时对单个工具调用设置 15 秒超时。哪个环节超时就在哪个环节终止不要一刀切。五、总结把 Agent 封装成 REST API看起来简单做好细节不简单。三个核心经验SSE 不是可选是必须。Agent 推理时间长流式推送让用户能看到进度容忍度从 5 秒提升到 30 秒以上。用astream_events而不是astream拿到每个事件级别的粒度。会话管理要提前设计。Agent 的记忆不是请求结束时消失的——历史对话、上一步工具调用结果都要在下次请求时加载回来。内存字典起步Redis 兜底。异常路径优先考虑。Agent 推理可能失败、工具调用可能超时、前端可能断连。正常路径跑通 30 分钟异常路径想清楚要花 3 小时。早想早安心。Agent 的服务化是 Agent 从玩具到产品的关键一跳。花点时间把流式推送、会话管理和异常处理做好这个 API 就能在线上稳稳地跑起来。下一篇预告RAG 服务的 API 密钥怎么管聊聊 Infisical 和 Vault 的工程实践。