TradingAgents-CN WebSocket 通知系统实战指南:从 SSE + Redis PubSub 到双向实时推送 TradingAgents-CN WebSocket 通知系统实战指南从 SSE Redis PubSub 到双向实时推送【免费下载链接】TradingAgents-CN基于多智能体LLM的中文金融交易框架 - TradingAgents中文增强版项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-CNTradingAgents-CN 在实时通知场景中引入了一套 WebSocket 通知系统用于替代此前的 SSE Redis PubSub 方案从根源上解决 Redis 连接泄漏问题。本文以仓库中的 docs/guides/websocket_notifications.md 为主体脉络结合 app/routers/websocket_notifications.py、app/services/websocket_manager.py、app/services/notifications_service.py 等源码完整讲解后端 WebSocket 端点、前端 Vue 3 TypeScript 集成、Nginx 反向代理配置以及从 SSE 平滑迁移的实战方案。读完本文你将掌握如何在 TradingAgents-CN 中接入实时通知流与任务进度流并理解其底层连接管理的实现原理。为什么用 WebSocket 替代 SSE Redis PubSubTradingAgents-CN 的分析任务耗时较长用户需要实时看到进度与结果通知。旧方案采用 SSE Redis PubSub每个 SSE 连接都会在 Redis 中创建独立的 PubSub 连接且不使用连接池在高并发下极易造成 Redis 连接泄漏最终拖垮整个数据同步与任务队列链路。WebSocket 方案的核心改进在于直接管理连接、不依赖 Redis PubSub。两者的对比如下特性SSE Redis PubSubWebSocket连接管理每个 SSE 连接创建独立的 PubSub 连接 ❌直接管理 WebSocket 连接 ✅Redis 连接不使用连接池容易泄漏 ❌不需要 Redis PubSub ✅双向通信单向服务器→客户端❌双向服务器↔客户端✅实时性较好 ⚠️更好 ✅连接数限制受 Redis 连接数限制 ❌只受服务器资源限制 ✅自动重连浏览器自动重连 ✅需要手动实现 ⚠️从源码看旧 SSE 实现 app/routers/sse.py 中每次task_progress_generator都会执行r.pubsub()创建独立连接并订阅task_progress:{task_id}频道尽管后续修复中补充了unsubscribe/close/reset等多级清理逻辑见 app/routers/sse.py但每个连接一条 PubSub 连接的架构缺陷依然存在。WebSocket 方案则通过进程内连接管理器直接持有连接对象彻底绕开了 Redis 这一中间层。后端 API 全景WebSocket 相关路由在 app/main.py 中被导入并在 app/main.py 处以/api前缀注册app.include_router(websocket_notifications_router.router, prefix/api, tags[websocket])因此所有端点的实际地址均以http://localhost:8000/api为前缀。1. WebSocket 通知端点ws://localhost:8000/api/ws/notifications?tokenjwt_token对应实现为 app/routers/websocket_notifications.py 中的websocket_notifications_endpoint。连接建立流程从 Query 参数取出token调用AuthService.verify_token(token)校验 JWT失败则await websocket.close(code1008, reasonUnauthorized)拒绝连接校验通过后调用manager.connect(websocket, user_id)注册连接立即下发connected类型的连接确认消息通过asyncio.create_task启动后台心跳任务每 30 秒发送一次heartbeat进入while True循环接收客户端消息主要用于保持连接活跃可处理 ping/pong客户端断开WebSocketDisconnect或异常时在finally中取消心跳任务并执行manager.disconnect清理连接。服务端下发的消息格式{ type: notification, // 消息类型: notification, heartbeat, connected data: { id: ..., title: 分析完成, content: 000001 分析已完成, type: analysis, link: /stocks/000001, source: analysis, created_at: 2025-10-23T12:00:00, status: unread } }其中data字段结构与 app/models/notification.py 中的NotificationOut模型一一对应type限定为analysis | alert | systemstatus限定为unread | read。2. WebSocket 任务进度端点ws://localhost:8000/api/ws/tasks/task_id?tokenjwt_token对应实现为 app/routers/websocket_notifications.py 中的websocket_task_progress_endpoint。连接后同样先发送connected确认消息包含task_id随后保持长连接等待任务状态流转。任务进度通过send_task_progress_via_websocket辅助函数推送消息格式如下{ type: progress, // 消息类型: progress, completed, error, heartbeat data: { task_id: ..., message: 正在分析..., step: 1, total_steps: 5, progress: 20.0, timestamp: 2025-10-23T12:00:00 } }需要说明的是当前实现中send_task_progress_via_websocket暂以manager.broadcast广播给所有连接源码注释明确提示生产环境应该只发给任务所属用户app/routers/websocket_notifications.py接入方若需要按用户隔离可从数据库查询任务归属或在progress_data中显式传递user_id。此外仓库还保留了另一套按任务维度管理连接的进度推送通道app/routers/analysis.py 中的/ws/task/{task_id}端点配合 app/services/websocket_manager.py 的WebSocketManager以task_id - Set[WebSocket]组织连接实际分析进度更新由 app/services/memory_state_manager.py 调用send_progress_update推送给订阅该任务的所有前端。3. WebSocket 连接统计GET /api/ws/stats响应示例{ total_users: 5, total_connections: 8, users: { admin: 2, user1: 1, user2: 1 } }该接口直接返回ConnectionManager.get_stats()的结果app/routers/websocket_notifications.pytotal_users为当前在线用户数total_connections为全部活跃连接数users为每个用户的连接数明细。它是运维排障与压力观测的第一手数据源。后端核心实现ConnectionManager 源码级解析WebSocket 通知系统的核心是定义在 app/routers/websocket_notifications.py 中的ConnectionManager类全局单例为模块底部的manager ConnectionManager()class ConnectionManager: def __init__(self): # user_id - Set[WebSocket] self.active_connections: Dict[str, Set[WebSocket]] {} self._lock asyncio.Lock()关键设计点以用户为维度管理连接active_connections采用user_id - Set[WebSocket]的结构天然支持一个用户多个连接如多个浏览器标签页send_personal_message会把消息发给指定用户的所有连接。锁内读、锁外写connect/disconnect在async with self._lock内修改集合而send_personal_message先在锁内拷贝连接列表再在锁外逐个send_text避免 I/O 操作阻塞其他连接的注册与注销发送失败死连接会在锁内统一清理discard后若集合为空则删除该用户条目。JSON 序列化统一处理消息统一json.dumps(message, ensure_asciiFalse)保证中文通知内容不被转义。心跳保活端点内部send_heartbeat任务每 30 秒发送一次heartbeat消息既让代理层感知连接活跃也让客户端可以据此判断连接健康状态。辅助发送函数send_notification_via_websocket(user_id, notification)与send_task_progress_via_websocket(task_id, progress_data)是给其他模块调用的统一入口见 app/routers/websocket_notifications.py。通知的完整链路在 app/services/notifications_service.py 的create_and_publish中体现通知先持久化到 MongoDBnotifications集合自动创建(user_id, created_at)与(user_id, status)索引再调用send_notification_via_websocket实时推送给用户若 WebSocket 发送失败仅记录 warning 日志并降级继续配合旧 SSE/Redis 通道实现兼容。持久化同时还内置清理策略默认保留最近 90 天、每用户最多 1000 条超出部分按时间删除最旧记录。前端集成Vue 3 TypeScript1. 创建 WebSocket Store前端采用 Pinia 管理 WebSocket 状态。以下 Store 封装了连接、断线自动重连指数退避、消息分发与桌面通知等完整逻辑// stores/websocket.ts import { ref, computed } from vue import { defineStore } from pinia import { useAuthStore } from ./auth export const useWebSocketStore defineStore(websocket, () { const ws refWebSocket | null(null) const connected ref(false) const reconnectTimer refnumber | null(null) const reconnectAttempts ref(0) const maxReconnectAttempts 5 // 连接 WebSocket function connect() { try { // 关闭现有连接 if (ws.value) { ws.value.close() ws.value null } const authStore useAuthStore() const token authStore.token || localStorage.getItem(auth-token) || const base import.meta.env.VITE_API_BASE_URL || const wsProtocol window.location.protocol https: ? wss: : ws: const wsHost base.replace(/^https?:\/\//, ).replace(/\/$/, ) const url ${wsProtocol}//${wsHost}/api/ws/notifications?token${encodeURIComponent(token)} console.log([WS] 连接到:, url) const socket new WebSocket(url) ws.value socket socket.onopen () { console.log([WS] 连接成功) connected.value true reconnectAttempts.value 0 } socket.onclose (event) { console.log([WS] 连接关闭:, event.code, event.reason) connected.value false ws.value null // 自动重连 if (reconnectAttempts.value maxReconnectAttempts) { const delay Math.min(1000 * Math.pow(2, reconnectAttempts.value), 30000) console.log([WS] ${delay}ms 后重连 (尝试 ${reconnectAttempts.value 1}/${maxReconnectAttempts})) reconnectTimer.value window.setTimeout(() { reconnectAttempts.value connect() }, delay) } else { console.error([WS] 达到最大重连次数停止重连) } } socket.onerror (error) { console.error([WS] 连接错误:, error) connected.value false } socket.onmessage (event) { try { const message JSON.parse(event.data) handleMessage(message) } catch (error) { console.error([WS] 解析消息失败:, error) } } } catch (error) { console.error([WS] 连接失败:, error) connected.value false } } // 处理消息 function handleMessage(message: any) { console.log([WS] 收到消息:, message) switch (message.type) { case connected: console.log([WS] 连接确认:, message.data) break case notification: // 处理通知 handleNotification(message.data) break case heartbeat: // 心跳消息无需处理 break default: console.warn([WS] 未知消息类型:, message.type) } } // 处理通知 function handleNotification(data: any) { // 添加到通知列表 const notificationsStore useNotificationsStore() notificationsStore.addNotification(data) // 显示桌面通知 if (Notification in window Notification.permission granted) { new Notification(data.title, { body: data.content, icon: /favicon.ico }) } } // 断开连接 function disconnect() { if (reconnectTimer.value) { clearTimeout(reconnectTimer.value) reconnectTimer.value null } if (ws.value) { ws.value.close() ws.value null } connected.value false reconnectAttempts.value 0 } // 发送消息 function send(message: any) { if (ws.value connected.value) { ws.value.send(JSON.stringify(message)) } else { console.warn([WS] 未连接无法发送消息) } } return { ws, connected, connect, disconnect, send } })几个值得注意的实现细节协议与地址推导根据window.location.protocol自动选择wss:/ws:并根据import.meta.env.VITE_API_BASE_URLVite 环境变量推导 WebSocket 主机前端部署在 HTTPS 下时不会出现混合内容被浏览器拦截的问题。重连策略指数退避Math.min(1000 * 2^n, 30000)最多尝试 5 次每次重连前会先关闭旧连接避免重复连接。消息分发通过type字段区分connected连接确认、notification业务通知、heartbeat心跳保活无需处理未知类型仅告警不中断。桌面通知handleNotification在将通知写入useNotificationsStore的同时若浏览器已授予Notification权限还会弹出系统级桌面通知icon 指向/favicon.ico适合分析完成这类需要用户留意的事件。2. 在 App.vue 中初始化在应用根组件挂载时按登录态建立连接卸载时断开script setup langts import { onMounted, onUnmounted } from vue import { useWebSocketStore } from /stores/websocket import { useAuthStore } from /stores/auth const wsStore useWebSocketStore() const authStore useAuthStore() onMounted(() { // 用户登录后连接 WebSocket if (authStore.isAuthenticated) { wsStore.connect() } }) onUnmounted(() { // 组件卸载时断开连接 wsStore.disconnect() }) /script建议实际项目中将connect()与登录动作绑定登录成功后调用、登出时调用disconnect()并清除重连定时器避免未登录状态下发起无效连接。配置环境变量# WebSocket 配置可选 WS_HEARTBEAT_INTERVAL30 # 心跳间隔秒 WS_MAX_CONNECTIONS_PER_USER3 # 每个用户最大连接数需要说明的是当前仓库版本中app/routers/websocket_notifications.py 的心跳间隔在端点内硬编码为 30 秒await asyncio.sleep(30)上述两个环境变量属于文档声明的预留配置项SSE 侧的同类参数如sse_heartbeat_interval_seconds、sse_poll_timeout_seconds则已通过 app/services/config_service.py 与 app/routers/sse.py 支持动态配置。若需要在 WebSocket 侧启用同样的可配置能力可参照 SSE 的实现模式将间隔值改为从配置中心读取。Nginx 配置当使用 Nginx 作为反向代理时/api/下必须显式开启 WebSocket 协议升级支持location /api/ { proxy_pass http://backend/api/; # WebSocket 支持必需 proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; # 超时设置重要 # WebSocket 长连接需要更长的超时时间 proxy_connect_timeout 120s; proxy_send_timeout 3600s; # 1小时 proxy_read_timeout 3600s; # 1小时 # 禁用缓存 proxy_buffering off; proxy_cache off; }关键配置说明proxy_http_version 1.1WebSocket 协议升级依赖 HTTP/1.1Nginx 默认的 HTTP/1.0 无法完成握手Upgrade和Connection头$http_upgrade会透传客户端的Upgrade: websocket请求头Connection upgrade指示 Nginx 保持升级后的长连接proxy_send_timeout和proxy_read_timeout设置为 3600s1 小时或更长如果设置太短如 120sWebSocket 连接会被意外关闭后端有心跳机制每 30 秒可以保持连接活跃proxy_buffering off禁用缓冲确保消息实时转发否则 Nginx 可能攒批发送导致前端感知延迟。仓库实际部署配置 nginx/nginx.conf 已按上述规范落地proxy_http_version 1.1、Upgrade/Connection头、proxy_buffering off; proxy_cache off;以及proxy_send_timeout 3600s; proxy_read_timeout 3600s;并额外设置了proxy_buffer_size 128k等缓冲区参数避免大响应被截断可直接作为生产参考。监控与调试查看连接统计curl http://localhost:8000/api/ws/stats响应示例{ total_users: 5, total_connections: 8, users: { admin: 2, user1: 1, user2: 1 } }观察服务端日志ConnectionManager在连接、断开、发送成功/失败时均输出结构化日志logger 名为webapi.websocket包括新连接/断开连接的用户与总连接数、心跳发送失败、死连接清理等可直接用docker logs或仓库提供的 scripts/view_logs.py 跟进实时状态。正常运行时每 30 秒应能观察到心跳发送记录若出现大量发送消息失败警告说明客户端断连未及时清理或 Nginx 超时配置过短。从 SSE 迁移到 WebSocket1. 后端无需修改通知服务会自动尝试 WebSocket失败时降级到 Redis PubSub兼容 SSE。这一降级逻辑在 app/services/notifications_service.py 中可见send_notification_via_websocket抛出的任何异常都只记 warning 日志不影响通知的持久化与后续查询同时 app/routers/sse.py 的 SSE 端点仍保留在路由表中旧客户端可以继续工作。2. 前端修改旧代码SSEconst sse new EventSource(/api/notifications/stream?token...) sse.addEventListener(notification, (event) { const data JSON.parse(event.data) // 处理通知 })新代码WebSocketconst ws new WebSocket(ws://localhost:8000/api/ws/notifications?token...) ws.onmessage (event) { const message JSON.parse(event.data) if (message.type notification) { // 处理通知 } }两者在解析data后处理通知的语义上基本一致迁移时只需把addEventListener回调改为onmessage统一入口并按type字段分流即可唯一的额外工作是为 WebSocket 补充手动重连逻辑上文 Store 已给出完整实现。若旧前端只依赖浏览器对EventSource的自动重连能力可暂时保留旧代码通过后端降级通道继续获得通知再择机切换。注意事项自动重连WebSocket 需要手动实现重连逻辑示例代码已包含建议配合指数退避与最大重试次数避免断网抖动时雪崩式重连心跳机制服务器每 30 秒发送一次心跳保持连接活跃同时可让 Nginx 长连接超时3600s始终被重置前端也可以将超过 N 秒未收到心跳视为连接失效并主动重连连接限制每个用户可以有多个连接例如多个浏览器标签页ConnectionManager以用户为维度组织连接集合推送时自动群发到该用户的全部连接兼容性旧的 SSE 客户端仍然可以工作通过 Redis PubSub 降级通道迁移过程可以前后端分阶段进行无需停机切换按用户隔离进度当前send_task_progress_via_websocket为广播实现生产环境接入时建议结合任务归属查询将进度消息收敛到任务所属用户的连接上。总结TradingAgents-CN 的 WebSocket 通知系统从架构上规避了 SSE Redis PubSub 每连接一条 Redis 连接的资源模型以进程内ConnectionManager用户维度 每 30 秒心跳 断线清理 自动降级的组合同时解决了连接泄漏、实时性和兼容性三个核心问题。对于新接入实时能力的功能模块推荐优先使用 WebSocket存量 SSE 客户端则可通过后端降级通道平稳过渡。相关的完整实现可继续研读 app/routers/websocket_notifications.py、app/services/websocket_manager.py、app/services/notifications_service.py 与 nginx/nginx.conf。【免费下载链接】TradingAgents-CN基于多智能体LLM的中文金融交易框架 - TradingAgents中文增强版项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-CN创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考