ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

基于WebSocket实现Agent思考过程的实时流式推送

基于WebSocket实现Agent思考过程的实时流式推送 基于 WebSocket 实现 Agent 思考过程的实时流式推送摘要多智能体系统的执行链路往往长达数十秒甚至数分钟用户盯着空白页面等待是不可接受的体验。本文详细介绍如何通过 WebSocket 单例监控器实现 Agent 思考过程的秒级实时推送涵盖架构设计、连接管理、跨线程安全、三层兜底策略及前端集成示例。其中流式解析部分参考了 LangGraph 的 chunk 结构但整个推送模块为独立实现不依赖 LangGraph 框架。一、为什么需要实时推送在单次请求-响应模式下用户提交任务后只能等待最终结果。但 Agent 的实际执行包含多个阶段意图分析、子Agent委派、工具调用、结果汇总、文档生成。如果每个阶段都黑盒运行用户无法感知进度也无法判断系统是正在思考还是已经卡死。实时推送要解决的核心问题进度可见用户能看到当前正在调用哪个子Agent、执行哪个工具异常感知工具调用失败时能第一时间通知前端而非等到超时体验升级类似 ChatGPT 的逐字输出让人感觉系统在为我工作技术选型上WebSocket 相比 SSEServer-Sent Events和短轮询有几个优势全双工通信支持前端主动发心跳保活、原生支持 JSON 结构化消息无需解析 text/event-stream 格式、单连接复用无需重复握手。在多智能体场景中前端还需要能够发送取消任务等指令全双工能力是刚需。二、架构总览Agent 执行引擎FastAPI 服务端前端asyncio.create_tasktool_calls / contentrun_coroutine_threadsafe定向推送用户提交任务POST /api/taskWebSocket 连接ws://.../ws/{thread_id}/api/task 路由创建后台异步任务/ws/{thread_id}WebSocket 端点ConnectionManager连接池 定向推送ToolMonitor 单例事件采集与分发run_deep_agent()异步流式执行_process_stream_chunk()解析每个增量Chunk数据流向用户发起 HTTP 请求 → 服务端创建后台任务立即返回 → Agent 在后台异步执行 → 每步执行通过 Monitor 单例采集事件 → 跨线程安全投递到 ConnectionManager → WebSocket 定向推送给对应 thread_id 的前端。三、核心实现3.1 连接管理按 thread_id 隔离ConnectionManager维护一个thread_id → WebSocket的映射字典确保每个会话的消息只推送给对应的前端连接classConnectionManager:def__init__(self):self.active_connections:Dict[str,WebSocket]{}self.loopNone# 延迟绑定事件循环defget_loop(self):懒加载获取当前事件循环同时自动绑定 Monitorifself.loopisNone:try:self.loopasyncio.get_running_loop()monitor.set_websocket_manager(self)# 双向绑定exceptRuntimeError:print([Monitor] Warning: No running event loop found.)returnself.loopasyncdefconnect(self,websocket:WebSocket,thread_id:str):self.get_loop()# 首次连接时绑定事件循环awaitwebsocket.accept()self.active_connections[thread_id]websocketdefdisconnect(self,websocket:WebSocket,thread_id:str):ifthread_idinself.active_connections:delself.active_connections[thread_id]asyncdefsend_to_thread(self,message:dict,thread_id:str):定向推送只发给指定 thread_id 的客户端ifthread_idinself.active_connections:websocketself.active_connections[thread_id]awaitwebsocket.send_json(message)关键设计点延迟绑定loop不在__init__中获取事件循环而是在首次connect时懒加载。这是因为 FastAPI 在启动时事件循环尚未就绪提前获取会得到None双向绑定get_loop()中自动调用monitor.set_websocket_manager(self)确保 Monitor 和 Manager 之间的引用关系建立后续 Monitor 才能将消息投递到 Manager3.2 WebSocket 端点心跳保活app.websocket(/ws/{thread_id})asyncdefwebsocket_endpoint(websocket:WebSocket,thread_id:str):awaitmanager.connect(websocket,thread_id)try:whileTrue:dataawaitwebsocket.receive_text()# 心跳响应awaitwebsocket.send_json({type:pong,message:f服务端已收到:{data}})exceptWebSocketDisconnect:manager.disconnect(websocket,thread_id)路由参数{thread_id}作为连接标识前端在建立连接时传入。进入消息循环后持续监听前端发来的心跳包ping回复 pong 保持连接活跃。一旦客户端断开或网络异常WebSocketDisconnect异常被捕获从 Manager 中移除该连接。3.3 监控器单例 三层兜底ToolMonitor是整个实时推送系统的核心枢纽采用单例模式确保全局只有一个实例classToolMonitor:_instanceNonedef__new__(cls):ifcls._instanceisNone:cls._instancesuper(ToolMonitor,cls).__new__(cls)cls._instance.websocket_managerNonereturncls._instancedef_emit(self,event_type:str,message:str,data:Optional[Dict[str,Any]]None):payload{type:monitor_event,event:event_type,message:message,data:dataor{},timestamp:datetime.datetime.now().isoformat()}# 第1层WebSocket 定向推送ifself.websocket_manager:thread_idget_thread_context()manager_loopself.websocket_manager.get_loop()ifmanager_loopandthread_id:try:current_loopasyncio.get_running_loop()exceptRuntimeError:current_loopNoneifcurrent_loopandcurrent_loopmanager_loop:current_loop.create_task(self.websocket_manager.send_to_thread(payload,thread_id))else:asyncio.run_coroutine_threadsafe(self.websocket_manager.send_to_thread(payload,thread_id),manager_loop)# 第2层脚本模式流式输出ifbuiltinsandhasattr(builtins,runtime)\andhasattr(builtins.runtime,stream_writer):try:builtins.runtime.stream_writer(payload)exceptException:pass# 第3层控制台保底print(f\n[Monitor:{event_type}]{message})_emit方法的核心逻辑是三层输出策略WebSocket 定向推送最高优先级从ContextVar中取出当前thread_id通过ConnectionManager定向推送给对应客户端脚本模式流式输出检测builtins.runtime.stream_writer是否存在兼容命令行脚本运行场景控制台print作为兜底确保任何环境下都能看到 Agent 的执行日志每层失败不影响下一层保证消息传递的鲁棒性。3.4 跨线程安全run_coroutine_threadsafe这是整个系统最容易被忽视的细节。Agent 在asyncio.create_task创建的后台任务中运行而 WebSocket 的send_json必须在 FastAPI 的事件循环中执行。如果两个循环不同比如用了线程池直接调用send_json会报错。解决方案是判断当前循环与 Manager 的循环是否一致同一循环直接create_task调度不同循环/线程使用asyncio.run_coroutine_threadsafe将协程安全投递到 Manager 所在的事件循环ifcurrent_loopandcurrent_loopmanager_loop:current_loop.create_task(...)# 同循环直接调度else:asyncio.run_coroutine_threadsafe(# 跨线程安全投递self.websocket_manager.send_to_thread(payload,thread_id),manager_loop)3.5 事件类型与触发时机Monitor 提供了四种事件上报方法覆盖 Agent 执行的完整生命周期方法事件类型触发时机携带数据report_session_dirsession_created会话环境初始化完成工作目录路径report_tooltool_start工具函数被调用时工具名、参数report_assistantassistant_call主Agent委派子Agent时子Agent名称、描述report_task_resulttask_resultAgent 输出最终回复完整回复内容事件的触发点位于流式处理函数中它通过解析 Agent 框架输出的增量chunk来识别当前正在发生什么。这里需要说明的是本项目并未使用完整的 LangGraph 框架但流式输出的 chunk 结构与 LangGraph 的astream一致——每个 chunk 是一个以节点名称为 key 的字典value 中包含messages列表def_process_stream_chunk(chunk):解析流式输出识别关键事件并上报fornode_name,stateinchunk.items():ifnotstateormessagesnotinstate:continuemessagesstate[messages]ifisinstance(messages,list)andmessages:last_msgmessages[-1]ifisinstance(last_msg,AIMessage):# AI 决定调用工具包括委派子Agentiflast_msg.tool_calls:fortoolinlast_msg.tool_calls:iftool[name]task:monitor.report_assistant(tool[args].get(subagent_type),{desc:tool[args].get(description)})# AI 输出最终回复eliflast_msg.content:monitor.report_task_result(last_msg.content)每个 chunk 的messages列表中最后一条消息反映了当前节点的状态如果tool_calls非空说明 Agent 正在调用工具如果tool_calls为空但有content说明 Agent 已经完成了本轮思考。task工具是 Agent 框架中用于委派子Agent的内置机制我们通过判断tool[name] task来特殊处理将子Agent的名称和描述推送给前端。这里有一个容易被忽略的细节流式输出的是增量状态字典每个 chunk 都可能包含多个节点的输出。我们只关心messages列表中的最后一条消息因为它代表了当前节点的最新状态。如果遍历所有消息会导致重复推送。3.6 前端集成示例前端只需两步发起任务 建立 WebSocket 监听。由于服务端采用即接即返模式thread_id在毫秒级返回前端可以立即建立 WebSocket 连接几乎不会错过任何推送消息// 1. 发起任务const{thread_id}awaitfetch(/api/task,{method:POST,body:JSON.stringify({query:分析销售数据并生成报告})}).then(rr.json());// 2. 建立 WebSocket 监听constwsnewWebSocket(ws://localhost:8000/ws/${thread_id});ws.onmessage(event){const{type,event:evt,message,data,timestamp}JSON.parse(event.data);switch(evt){casesession_created:console.log(工作目录已创建:${data.path});break;caseassistant_call:updateUI(正在调用:${data.assistant_name});break;casetool_start:updateUI(执行工具:${data.tool_name});break;casetask_result:showFinalResult(data.result);break;}};// 3. 心跳保活每30秒发送一次pingsetInterval(()ws.send(ping),30000);四、工程实践要点4.1 单例模式的必要性ToolMonitor必须全局唯一。如果每个工具调用都创建一个新的 Monitor 实例那么websocket_manager的引用会丢失事件无法送达 WebSocket。单例保证了所有模块Agent、工具函数、API 层共享同一个 Monitor 实例引用关系始终有效。Python 中实现单例的经典方式是重写__new__方法在首次实例化时创建对象并缓存到类变量_instance后续调用直接返回缓存。这种方式比装饰器和元类更直观且不会影响类的继承关系。4.2 ContextVar 实现会话隔离Monitor 的_emit方法通过get_thread_context()获取当前thread_id这在多用户并发场景下至关重要。ContextVar是 Python 3.7 的协程安全变量每个异步任务有独立的上下文副本确保用户A的任务进度不会推送到用户B的前端。thread_id的绑定发生在run_deep_agent()入口处在finally块中通过reset_session_context清理。这个清理步骤不可省略——FastAPI 的asyncio.create_task可能会复用协程如果不清理 ContextVar下一次任务的thread_id会残留上一次的值导致消息串台。4.3 降级优先的设计哲学三层输出策略体现了不阻塞主流程的原则。即使 WebSocket 断开、脚本模式不可用Monitor 的_emit方法仍然能通过print输出日志Agent 的执行不会因为推送失败而中断。这种降级优先的设计在 Agent 这类长链路系统中尤为重要——推送是锦上添花但任务执行不能因此受阻。在实际运行中我们观察到最常见的故障场景是用户在任务执行中途关闭了浏览器标签页导致 WebSocket 断开。此时ConnectionManager的send_to_thread会因thread_id不在active_connections中而静默跳过Monitor 继续走第二层和第三层输出不影响 Agent 继续执行。等用户重新打开页面并建立新的 WebSocket 连接后由于thread_id相同后续消息仍能正常推送。五、总结本文从零搭建了一套 Agent 实时推送系统核心组件包括ConnectionManager维护thread_id → WebSocket映射实现会话级定向推送ToolMonitor 单例全局事件采集与三层分发跨线程安全投递asyncio.run_coroutine_threadsafe解决事件循环隔离问题流式 chunk 解析从增量 chunk 中识别工具调用和子Agent委派事件这套方案虽然代码量不大核心逻辑约 200 行但覆盖了实时推送系统的关键工程问题连接管理、消息路由、跨线程安全、降级策略。如果你的项目也需要让 Agent 的思考过程对用户可见这套架构可以直接复用。技术栈Python 3.10 / FastAPI / WebSocket / asyncio / ContextVar适用场景Agent 执行过程可视化、实时日志监控、多智能体调试面板
返回列表