ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

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

2026/8/9 8:22:22 拓冰建站 浏览量
基于WebSocket实现Agent思考过程的实时流式推送

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

摘要:多智能体系统的执行链路往往长达数十秒甚至数分钟,用户盯着空白页面等待是不可接受的体验。本文详细介绍如何通过 WebSocket + 单例监控器,实现 Agent 思考过程的秒级实时推送,涵盖架构设计、连接管理、跨线程安全、三层兜底策略及前端集成示例。其中流式解析部分参考了 LangGraph 的 chunk 结构,但整个推送模块为独立实现,不依赖 LangGraph 框架。


一、为什么需要实时推送

在单次请求-响应模式下,用户提交任务后只能等待最终结果。但 Agent 的实际执行包含多个阶段:意图分析、子Agent委派、工具调用、结果汇总、文档生成。如果每个阶段都黑盒运行,用户无法感知进度,也无法判断系统是"正在思考"还是"已经卡死"。

实时推送要解决的核心问题

  • 进度可见:用户能看到当前正在调用哪个子Agent、执行哪个工具
  • 异常感知:工具调用失败时能第一时间通知前端,而非等到超时
  • 体验升级:类似 ChatGPT 的逐字输出,让人感觉"系统在为我工作"

技术选型上,WebSocket 相比 SSE(Server-Sent Events)和短轮询有几个优势:全双工通信支持前端主动发心跳保活、原生支持 JSON 结构化消息无需解析 text/event-stream 格式、单连接复用无需重复握手。在多智能体场景中,前端还需要能够发送取消任务等指令,全双工能力是刚需。


二、架构总览

Agent 执行引擎

FastAPI 服务端

前端

asyncio.create_task

tool_calls / content

run_coroutine_threadsafe

定向推送

用户提交任务
POST /api/task

WebSocket 连接
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.loop=None# 延迟绑定事件循环defget_loop(self):"""懒加载获取当前事件循环,同时自动绑定 Monitor"""ifself.loopisNone:try:self.loop=asyncio.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:websocket=self.active_connections[thread_id]awaitwebsocket.send_json(message)

关键设计点:

  • 延迟绑定loop:不在__init__中获取事件循环,而是在首次connect时懒加载。这是因为 FastAPI 在启动时事件循环尚未就绪,提前获取会得到None
  • 双向绑定get_loop()中自动调用monitor.set_websocket_manager(self),确保 Monitor 和 Manager 之间的引用关系建立,后续 Monitor 才能将消息投递到 Manager

3.2 WebSocket 端点:心跳保活

@app.websocket("/ws/{thread_id}")asyncdefwebsocket_endpoint(websocket:WebSocket,thread_id:str):awaitmanager.connect(websocket,thread_id)try:whileTrue:data=awaitwebsocket.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:_instance=Nonedef__new__(cls):ifcls._instanceisNone:cls._instance=super(ToolMonitor,cls).__new__(cls)cls._instance.websocket_manager=Nonereturncls._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_id=get_thread_context()manager_loop=self.websocket_manager.get_loop()ifmanager_loopandthread_id:try:current_loop=asyncio.get_running_loop()exceptRuntimeError:current_loop=Noneifcurrent_loopandcurrent_loop==manager_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方法的核心逻辑是三层输出策略

  1. WebSocket 定向推送(最高优先级):从ContextVar中取出当前thread_id,通过ConnectionManager定向推送给对应客户端
  2. 脚本模式流式输出:检测builtins.runtime.stream_writer是否存在,兼容命令行脚本运行场景
  3. 控制台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_loop==manager_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():ifnotstateor"messages"notinstate:continuemessages=state["messages"]ifisinstance(messages,list)andmessages:last_msg=messages[-1]ifisinstance(last_msg,AIMessage):# AI 决定调用工具(包括委派子Agent)iflast_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(r=>r.json());// 2. 建立 WebSocket 监听constws=newWebSocket(`ws://localhost:8000/ws/${thread_id}`);ws.onmessage=(event)=>{const{type,event:evt,message,data,timestamp}=JSON.parse(event.data);switch(evt){case'session_created':console.log(`工作目录已创建:${data.path}`);break;case'assistant_call':updateUI(`正在调用:${data.assistant_name}`);break;case'tool_start':updateUI(`执行工具:${data.tool_name}`);break;case'task_result':showFinalResult(data.result);break;}};// 3. 心跳保活(每30秒发送一次ping)setInterval(()=>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 断开。此时ConnectionManagersend_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 执行过程可视化、实时日志监控、多智能体调试面板