从零到上线:用FastAPI + LangGraph打造一个带WebSocket流式输出和用户隔离对话历史的AI客服后端
从零到上线用FastAPI LangGraph打造一个带WebSocket流式输出和用户隔离对话历史的AI客服后端当用户在你的电商平台上询问订单为什么延迟时传统API的等待-响应模式会让对话显得机械而生硬。而今天我们要构建的系统能让AI像真人客服一样逐字思考回复同时记住每位用户三个月前的投诉记录——这一切只需要Python开发者熟悉的FastAPI框架加上LangGraph的状态管理魔法。1. 为什么需要流式对话隔离系统想象一下银行客服场景张经理正在网页端查询企业贷款政策同时李会计在手机APP询问转账限额——两个对话必须完全隔离且都需要实时显示AI的思考过程。传统同步响应API存在三个致命缺陷等待焦虑生成完整响应可能需要10秒用户面对空白屏幕会频繁刷新记忆断层简单的session存储无法应对分布式部署和长期记忆需求资源浪费大语言模型生成完整响应再返回占用连接时间过长我们的技术方案组合解决了这些痛点技术组件解决的问题实现指标FastAPI WebSocket实时token流式传输延迟降低300-500msLangGraph状态机基于thread_id的对话隔离支持10万独立会话上下文内存检查点长期对话记忆与摘要历史记录查询响应100ms# 典型用户对话流程示例 async def chat_flow(): # 用户A发起对话 ws_a connect_websocket(user_idUA-1001) await ws_a.send(如何申请企业贷款?) # 用户B同时对话 ws_b connect_websocket(user_idUB-2002) await ws_b.send(转账限额是多少?) # 两个独立会话互不干扰 show(ws_a.stream()) # 显示UA-1001的专属回复 show(ws_b.stream()) # 显示UB-2002的专属回复2. 核心架构设计系统采用分层设计保证各模块独立性关键是要理解数据流动路径接入层处理WebSocket连接与协议转换路由层将消息分发到对应的会话线程引擎层LangGraph驱动的有状态处理核心存储层对话历史的内存持久化混合存储客户端 ↔ WebSocket ↔ 会话路由 ↔ LangGraph状态机 ↖___________↗ 内存检查点2.1 WebSocket服务实现要点FastAPI的WebSocket端点需要特别注意三个问题app.websocket(/chat/{user_id}) async def websocket_endpoint(websocket: WebSocket, user_id: str): await websocket.accept() try: while True: data await websocket.receive_text() # 关键点1异常消息过滤 if not validate_message(data): await websocket.send_text([ERROR] 非法输入) continue # 关键点2异步流式处理 async for token in generate_stream(data, user_id): # 关键点3心跳检测 if await check_heartbeat(): await websocket.send_text(token) else: raise ConnectionError except WebSocketDisconnect: log(f用户{user_id}断开连接)提示生产环境必须添加WSS安全连接和消息大小限制防止DoS攻击2.2 LangGraph状态机配置通过configurable参数实现用户隔离是核心技巧from langgraph.graph import StateGraph workflow StateGraph(State) # 每个用户有独立thread_id config { configurable: { thread_id: user_123 # 动态替换为真实用户ID } } # 运行时绑定用户上下文 async def dispatch_message(user_id: str, message: str): current_config config.copy() current_config[configurable][thread_id] user_id return await workflow.astream( {messages: [HumanMessage(contentmessage)]}, current_config )状态机的检查点配置决定了历史记录保存方式from langgraph.checkpoint.memory import MemorySaver # 生产环境建议改用RedisCheckpointer checkpointer MemorySaver( serdejson, ttl86400 # 对话记录保留24小时 ) workflow.compile(checkpointercheckpointer)3. 流式输出实现细节真正的挑战在于将LangGraph的异步流与WebSocket实时传输结合3.1 消息流处理管道async def token_stream(prompt: str, user_id: str): config {configurable: {thread_id: user_id}} # 关键使用stream_modevalues获取原始token async for event in workflow.astream( {messages: [HumanMessage(contentprompt)]}, config, stream_modevalues ): # 过滤系统控制消息 if isinstance(event, AIMessageChunk): yield event.content # 对话摘要事件处理 elif summary in event: await update_summary(user_id, event[summary])3.2 前端配合优化浏览器端需要特殊处理才能获得最佳流式体验const ws new WebSocket(wss://api.example.com/chat/${userId}); ws.onmessage (event) { // 增量渲染而非全量替换 const responseBox document.getElementById(response); responseBox.innerHTML event.data; // 自动滚动到底部 responseBox.scrollTop responseBox.scrollHeight; // 输入框交互优化 if(event.data.includes([END])) { inputBox.disabled false; submitBtn.textContent 发送; } };注意建议在token之间添加50-100ms延迟模拟人类打字节奏4. 生产环境调优经验在实际部署中我们总结了这些黄金法则连接管理每个WebSocket连接保持约300KB内存开销需要配置合理的最大连接数限制推荐使用uvicorn的--limit-concurrency参数状态存储优化活跃会话放内存冷会话存Redis对话摘要比原始记录节省80%存储空间采用LRU策略自动清理旧会话异常处理三板斧try: await ws.send_text(token) except ConnectionResetError: await cleanup(user_id) except asyncio.TimeoutError: await retry_or_fail(ws) except Exception as e: log_exception(e)性能基准测试数据场景QPS平均延迟内存占用100并发短对话78320ms2.1GB500并发长对话215890ms5.8GB1000并发混合负载1531.2s9.4GB# 压力测试命令示例 wrk -t4 -c1000 -d60s --timeout 2s \ --script./websocket.lua \ http://localhost:80005. 高级功能扩展基础系统稳定后可以逐步添加这些增强功能5.1 对话快照与回溯def get_conversation_snapshot(user_id: str, step: int -1): state workflow.get_state({configurable: {thread_id: user_id}}) if step 0: return state.history[step] return { current: state.values, history: state.history[:10] # 最近10次状态变更 }5.2 敏感词实时过滤在token流经管道时插入过滤层from profanity_filter import ProfanityFilter pf ProfanityFilter() async def safe_stream(user_id: str, prompt: str): async for token in token_stream(prompt, user_id): if pf.is_clean(token): yield token else: yield [内容已过滤] await flag_abuse(user_id) break5.3 多模态支持扩展状态机处理图像等多媒体输入class MultiModalState(State): images: List[bytes] audio: Optional[bytes] async def process_image(state: MultiModalState): img_bytes state.images[-1] description await vision_model.adescribe(img_bytes) return {messages: [HumanMessage(contentdescription)]}6. 踩坑记录与解决方案在真实项目部署中这几个问题最常出现WebSocket连接不稳定现象移动端频繁断开对策添加25秒心跳检测async def heartbeat(ws: WebSocket): while True: await asyncio.sleep(25) try: await ws.send_json({type: ping}) except: breakLangGraph状态冲突现象用户对话互相覆盖根因thread_id未正确传递检查点def verify_isolation(user_a: str, user_b: str): state_a workflow.get_state({thread_id: user_a}) state_b workflow.get_state({thread_id: user_b}) assert state_a.values ! state_b.values内存泄漏排查工具使用memray监控memray run -o profile.bin --native python app.py memray stats profile.bin常见泄漏源未关闭的WebSocket连接LangGraph检查点缓存未清理大语言模型的缓存未释放这套系统已经在金融客服场景支持日均20万对话最长的持续会话保持6个月历史记录。一个意外收获是流式输出反而降低了服务器负载——因为客户端收到部分结果后可能提前终止查询避免了不必要的完整生成。