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 的技术要点
SSE(Server-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_length=1, description="用户输入") session_id: str = Field(default_factory=lambda: 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( model=model, temperature=0, streaming=True, # 关键:启用流式 ) self._tools = [get_weather, send_slack] self._prompt = ChatPromptTemplate.from_messages([ ("system", "你是一个智能助手。使用工具来回答问题,逐步推理。"), MessagesPlaceholder(variable_name="chat_history", optional=True), ("human", "{input}"), MessagesPlaceholder(variable_name="agent_scratchpad"), ]) def create(self) -> AgentExecutor: agent = create_openai_tools_agent(self._llm, self._tools, self._prompt) return AgentExecutor( agent=agent, tools=self._tools, verbose=False, max_iterations=10, handle_parsing_errors=True, ) # ── SSE 事件序列化 ───────────────────────────────────── class SSEEvent: """SSE 事件格式化。""" @staticmethod def format(event_type: str, data: dict[str, Any]) -> str: payload = json.dumps({"type": event_type, **data}, ensure_ascii=False) return f"data: {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_ascii=False) return f"data: {payload}\n\n" # ── FastAPI 应用 ─────────────────────────────────────── @asynccontextmanager async def lifespan(app: FastAPI): """应用生命周期管理。""" app.state.sessions = SessionManager(max_history=20) app.state.agent_factory = AgentFactory() logger.info("Agent API 服务已启动") yield logger.info("Agent API 服务正在关闭") app = FastAPI(title="Agent API", lifespan=lifespan) @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, }, version="v2", ): 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(content=req.message)) app.state.sessions.append(session_id, AIMessage(content=output)) 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_type="text/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, host="0.0.0.0", port=8000)关键设计点:
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
| 方案 | 优势 | 劣势 | 适用场景 |
|---|---|---|---|
| SSE | HTTP 协议原生,走 CDN 无压力,自动重连 | 单向推送,需额外 POST 发消息 | Agent 对话(一次请求多次推送) |
| WebSocket | 双向通信,低延迟 | 需要自己管理重连、心跳,代理配置复杂 | 实时协作、多轮交互频繁 |
| Polling | 最简单,兼容性好 | 浪费带宽,延迟高 | 不需要实时反馈的场景 |
Agent 对话这个场景,SSE 是最合适的——前端 POST 一次消息,服务端一路流式推回结果。没有双向通信的需求,WebSocket 的复杂度是多余的。
4.2 生产级的并发与压测考量
一个 Agent 推理会占用 LLM API 的连接和本地 asyncio Task。压测时要关注的指标:
- 最大并发数:取决于 LLM API 的 rate limit 和本地 CPU/内存。一般单机 50-100 并发 Agent 对话是比较安全的范围。
- 背压处理:当并发满时,新请求应该返回 429(Too Many Requests),而不是排队等。前端看到 429 可以提示用户稍后重试。
- Token 级别的速率限制:除了并发数,还要限制每用户每分钟的 Token 消耗,防止单个用户打爆预算。
4.3 会话持久化
上面的实现用了内存字典_sessions: dict。单机部署完全够用,但要做持久化的话有几种选择:
- Redis:存会话历史 JSON,设置 TTL。适合多实例部署。
- SQLite/Postgres:存结构化消息记录,方便后续做分析和评估。
- LangChain 的 BaseChatMessageHistory:LangChain 有内置的 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 的工程实践。