实时数据流

📎 引用文件

本文引用的文件 - agent/src/api/swarm_routes.py - frontend/src/hooks/useSSE.ts - frontend/src/lib/apiAuth.ts - frontend/src/lib/api.ts - frontend/src/components/layout/ConnectionBanner.tsx - agent/src/api/live_routes.py - agent/src/channels/runtime.py - agent/src/market_data.py

目录

  1. 简介
  2. 项目结构
  3. 核心组件
  4. 架构总览
  5. 详细组件分析
  6. 依赖关系分析
  7. 性能与优化
  8. 故障排查指南
  9. 结论
  10. 附录:典型实时应用场景

简介

本架构文档聚焦 Vibe-Trading 的实时数据流系统,覆盖 WebSocket/SSE 连接管理、事件推送机制、前端实时通信钩子、状态同步与更新策略,以及重连、负载均衡、带宽优化与延迟控制。同时提供监控、性能分析与故障排查方法,并展示行情推送、交易执行反馈与系统状态监控等典型场景。

项目结构

graph TB subgraph "前端" FE_HOOK["useSSE 钩子"] FE_AUTH["withAuthTicket 票据"] FE_API["sseUrl 构建"] FE_UI["连接横幅"] end subgraph "后端 FastAPI" SWARM_SSE["/sessions/{id}/events<br/>SSE 事件流"] LIVE_ROUTES["/live/*<br/>授权/指令/状态/运行器"] CHANNEL_RT["channels 运行时"] end FE_HOOK --> FE_AUTH FE_HOOK --> FE_API FE_HOOK --> SWARM_SSE FE_UI --> FE_HOOK SWARM_SSE --> CHANNEL_RT LIVE_ROUTES --> SWARM_SSE

图表来源 - agent/src/api/swarm_routes.py:187-211 - frontend/src/hooks/useSSE.ts:27-214 - frontend/src/lib/apiAuth.ts:18-53 - frontend/src/lib/api.ts:181-188 - frontend/src/components/layout/ConnectionBanner.tsx:1-43 - agent/src/api/live_routes.py:631-720 - agent/src/channels/runtime.py:72-102

章节来源 - agent/src/api/swarm_routes.py:187-211 - frontend/src/hooks/useSSE.ts:27-214 - frontend/src/lib/apiAuth.ts:18-53 - frontend/src/lib/api.ts:181-188 - frontend/src/components/layout/ConnectionBanner.tsx:1-43 - agent/src/api/live_routes.py:631-720 - agent/src/channels/runtime.py:72-102

核心组件

章节来源 - agent/src/api/swarm_routes.py:187-211 - frontend/src/hooks/useSSE.ts:27-214 - agent/src/api/live_routes.py:631-720 - agent/src/channels/runtime.py:72-102

架构总览

下图展示了从前端连接到后端事件流的完整时序,包括票据交换、事件订阅、断线重连与续传。

sequenceDiagram participant FE as "前端 useSSE" participant AUTH as "withAuthTicket" participant API as "api.ts sseUrl" participant SSE as "后端 /sessions/{id}/events" participant BUS as "Session EventBus" FE->>API : 构建 SSE URL FE->>AUTH : 请求一次性票据 AUTH-->>FE : 返回 ticket FE->>SSE : EventSource(ticket) SSE-->>FE : id/event/data (增量事件) Note over FE,SSE : 断线后使用 Last-Event-ID 续传 SSE-->>FE : event : done FE->>BUS : 分发事件到业务处理器

图表来源 - frontend/src/lib/api.ts:181-188 - frontend/src/lib/apiAuth.ts:18-53 - agent/src/api/swarm_routes.py:187-211

详细组件分析

SSE 事件流(后端)

flowchart TD Start(["进入事件流"]) --> CheckConn{"连接是否断开?"} CheckConn --> |是| End(["结束流"]) CheckConn --> |否| ReadEvents["读取会话事件(after_index=idx)"] ReadEvents --> ForEachEvt{"有事件?"} ForEachEvt --> |是| EmitEvt["输出 id/event/data"] EmitEvt --> IncIdx["idx++"] --> ForEachEvt ForEachEvt --> |否| LoadRun["加载运行状态"] LoadRun --> RunExists{"运行存在?"} RunExists --> |否| DoneMissing["输出 done{status: missing}"] --> End RunExists --> |是| Reconcile["reconcile_run 写入"] Reconcile --> Terminal{"完成/失败/取消?"} Terminal --> |是| DoneTerm["输出 done{status}"] --> End Terminal --> |否| Sleep["等待 2s"] --> CheckConn

图表来源 - agent/src/api/swarm_routes.py:187-211

章节来源 - agent/src/api/swarm_routes.py:187-211

前端 SSE 钩子(useSSE)

classDiagram class UseSSE { +connect(url, handlers) +disconnect() +getStatus() +onStatusChange(cb) -attach(url, generation) -doConnect(generation) -scheduleReconnect(generation) -buildUrl(baseUrl) -trackEventId(eventId) bool }

图表来源 - frontend/src/hooks/useSSE.ts:27-214

章节来源 - frontend/src/hooks/useSSE.ts:27-214 - frontend/src/components/layout/ConnectionBanner.tsx:1-43

实时交易通道(/live/*)

sequenceDiagram participant UI as "前端 LiveRuntimePanel" participant API as "/live/halt | /live/resume" participant BUS as "Session EventBus" participant SSE as "SSE 流" UI->>API : POST /live/halt 或 /live/resume API-->>UI : 返回操作结果 API->>BUS : emit("live.halted"|"live.resumed", data) BUS-->>SSE : 推送到已订阅的会话流 SSE-->>UI : 事件到达,刷新状态

图表来源 - agent/src/api/live_routes.py:684-720 - agent/src/api/live_routes.py:722-800

章节来源 - agent/src/api/live_routes.py:631-800

IM 通道运行时

章节来源 - agent/src/channels/runtime.py:72-102

市场数据与压缩/限流

章节来源 - agent/src/market_data.py:1-229

依赖关系分析

graph LR FE_USESSE["useSSE"] --> FE_TICKET["withAuthTicket"] FE_USESSE --> FE_URL["api.sseUrl"] FE_BANNER["ConnectionBanner"] --> FE_USESSE SWARM["swarm_routes.events"] --> STORE["会话存储/运行状态"] LIVE["live_routes"] --> BUS["Session EventBus"] BUS --> SWARM CHANNELS["channels.runtime"] --> ADAPTERS["平台适配器"]

图表来源 - frontend/src/hooks/useSSE.ts:27-214 - frontend/src/lib/apiAuth.ts:18-53 - frontend/src/lib/api.ts:181-188 - frontend/src/components/layout/ConnectionBanner.tsx:1-43 - agent/src/api/swarm_routes.py:187-211 - agent/src/api/live_routes.py:631-720 - agent/src/channels/runtime.py:72-102

章节来源 - frontend/src/hooks/useSSE.ts:27-214 - frontend/src/lib/apiAuth.ts:18-53 - frontend/src/lib/api.ts:181-188 - frontend/src/components/layout/ConnectionBanner.tsx:1-43 - agent/src/api/swarm_routes.py:187-211 - agent/src/api/live_routes.py:631-720 - agent/src/channels/runtime.py:72-102

性能与优化

[本节为通用指导,不直接分析具体文件]

故障排查指南

章节来源 - frontend/src/lib/apiAuth.ts:18-53 - frontend/src/hooks/useSSE.ts:27-214 - agent/src/api/swarm_routes.py:187-211 - agent/src/api/live_routes.py:631-720 - agent/src/market_data.py:66-84

结论

Vibe-Trading 的实时数据流以 SSE 为核心,结合前端 useSSE 钩子的重连、去重与续传机制,实现了高可靠、低延迟的事件推送。实时交易通道通过 Session EventBus 将关键动作广播到已有 SSE 流,简化了前端集成。配合数据源回退链与行级采样,系统在可用性与性能之间取得良好平衡。建议在生产环境持续监控连接状态、事件吞吐与延迟指标,并结合票据机制与断点续传保障鲁棒性。

[本节为总结,不直接分析具体文件]

附录:典型实时应用场景

章节来源 - agent/src/market_data.py:97-223 - agent/src/api/live_routes.py:684-720 - agent/src/channels/runtime.py:72-102