TradingAgents-CN WebSocket 通知系统实战指南:从 SSE + Redis PubSub 到双向实时推送
2026/9/13 1:39:06 网站建设 项目流程

TradingAgents-CN WebSocket 通知系统实战指南:从 SSE + Redis PubSub 到双向实时推送

【免费下载链接】TradingAgents-CN基于多智能体LLM的中文金融交易框架 - TradingAgents中文增强版项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-CN

TradingAgents-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 PubSub

TradingAgents-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?token=<jwt_token>

对应实现为 app/routers/websocket_notifications.py 中的websocket_notifications_endpoint。连接建立流程:

  1. 从 Query 参数取出token,调用AuthService.verify_token(token)校验 JWT,失败则await websocket.close(code=1008, reason="Unauthorized")拒绝连接;
  2. 校验通过后调用manager.connect(websocket, user_id)注册连接;
  3. 立即下发connected类型的连接确认消息;
  4. 通过asyncio.create_task启动后台心跳任务(每 30 秒发送一次heartbeat);
  5. 进入while True循环接收客户端消息(主要用于保持连接活跃,可处理 ping/pong);
  6. 客户端断开(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>?token=<jwt_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.py):total_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/disconnectasync with self._lock内修改集合;而send_personal_message先在锁内拷贝连接列表,再在锁外逐个send_text,避免 I/O 操作阻塞其他连接的注册与注销;发送失败(死连接)会在锁内统一清理(discard后若集合为空则删除该用户条目)。
  • JSON 序列化统一处理:消息统一json.dumps(message, ensure_ascii=False),保证中文通知内容不被转义。
  • 心跳保活:端点内部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中体现:通知先持久化到 MongoDB(notifications集合,自动创建(user_id, created_at)(user_id, status)索引),再调用send_notification_via_websocket实时推送给用户;若 WebSocket 发送失败,仅记录 warning 日志并降级继续(配合旧 SSE/Redis 通道实现兼容)。持久化同时还内置清理策略:默认保留最近 90 天、每用户最多 1000 条,超出部分按时间删除最旧记录。

前端集成(Vue 3 + TypeScript)

1. 创建 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 = ref<WebSocket | null>(null) const connected = ref(false) const reconnectTimer = ref<number | 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_URL(Vite 环境变量)推导 WebSocket 主机,前端部署在 HTTPS 下时不会出现混合内容被浏览器拦截的问题。
  • 重连策略:指数退避Math.min(1000 * 2^n, 30000),最多尝试 5 次,每次重连前会先关闭旧连接避免重复连接。
  • 消息分发:通过type字段区分connected(连接确认)、notification(业务通知)、heartbeat(心跳保活,无需处理),未知类型仅告警不中断。
  • 桌面通知handleNotification在将通知写入useNotificationsStore的同时,若浏览器已授予Notification权限,还会弹出系统级桌面通知(icon 指向/favicon.ico),适合分析完成这类需要用户留意的事件。

2. 在 App.vue 中初始化

在应用根组件挂载时按登录态建立连接,卸载时断开:

<script setup lang="ts"> 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_INTERVAL=30 # 心跳间隔(秒) WS_MAX_CONNECTIONS_PER_USER=3 # 每个用户最大连接数

需要说明的是:当前仓库版本中,app/routers/websocket_notifications.py 的心跳间隔在端点内硬编码为 30 秒(await asyncio.sleep(30)),上述两个环境变量属于文档声明的预留配置项;SSE 侧的同类参数(如sse_heartbeat_interval_secondssse_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; }

关键配置说明:

  1. proxy_http_version 1.1:WebSocket 协议升级依赖 HTTP/1.1,Nginx 默认的 HTTP/1.0 无法完成握手;
  2. UpgradeConnection$http_upgrade会透传客户端的Upgrade: websocket请求头,Connection "upgrade"指示 Nginx 保持升级后的长连接;
  3. proxy_send_timeoutproxy_read_timeout
    • 设置为 3600s(1 小时)或更长;
    • 如果设置太短(如 120s),WebSocket 连接会被意外关闭;
    • 后端有心跳机制(每 30 秒),可以保持连接活跃;
  4. proxy_buffering off:禁用缓冲,确保消息实时转发,否则 Nginx 可能攒批发送导致前端感知延迟。

仓库实际部署配置 nginx/nginx.conf 已按上述规范落地:proxy_http_version 1.1Upgrade/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 迁移到 WebSocket

1. 后端无需修改

通知服务会自动尝试 WebSocket,失败时降级到 Redis PubSub(兼容 SSE)。这一降级逻辑在 app/services/notifications_service.py 中可见:send_notification_via_websocket抛出的任何异常都只记 warning 日志,不影响通知的持久化与后续查询;同时 app/routers/sse.py 的 SSE 端点仍保留在路由表中,旧客户端可以继续工作。

2. 前端修改

旧代码(SSE)

const sse = new EventSource('/api/notifications/stream?token=...') sse.addEventListener('notification', (event) => { const data = JSON.parse(event.data) // 处理通知 })

新代码(WebSocket)

const 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的自动重连能力,可暂时保留旧代码,通过后端降级通道继续获得通知,再择机切换。

注意事项

  1. 自动重连:WebSocket 需要手动实现重连逻辑(示例代码已包含),建议配合指数退避与最大重试次数,避免断网抖动时雪崩式重连;
  2. 心跳机制:服务器每 30 秒发送一次心跳,保持连接活跃,同时可让 Nginx 长连接超时(3600s)始终被重置,前端也可以将"超过 N 秒未收到心跳"视为连接失效并主动重连;
  3. 连接限制:每个用户可以有多个连接(例如多个浏览器标签页),ConnectionManager以用户为维度组织连接集合,推送时自动群发到该用户的全部连接;
  4. 兼容性:旧的 SSE 客户端仍然可以工作(通过 Redis PubSub 降级通道),迁移过程可以前后端分阶段进行,无需停机切换;
  5. 按用户隔离进度:当前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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询