数据流设计

📎 引用文件

本文引用的文件 - agent/src/channels/bus/queue.py - agent/src/channels/bus/events.py - agent/src/channels/runtime.py - agent/backtest/loaders/base.py - agent/backtest/loaders/registry.py - agent/src/session/service.py - agent/src/memory/persistent.py - agent/src/api/swarm_routes.py - frontend/src/stores/agent.ts

目录

  1. 简介
  2. 项目结构
  3. 核心组件
  4. 架构总览
  5. 详细组件分析
  6. 依赖关系分析
  7. 性能考虑
  8. 故障恢复与监控
  9. 结论
  10. 附录:典型数据流转场景

简介

本架构文档面向 Vibe-Trading 的数据流设计,覆盖数据采集、处理、存储与传输的全链路。系统采用事件驱动的消息总线与可插拔的数据加载器注册表,结合会话服务驱动的 Agent 执行管线,提供实时 SSE 推送与前端状态缓存,形成从外部数据源到研究/回测再到用户界面的完整闭环。文档同时说明数据模型、序列化格式、版本兼容策略、数据验证规则、转换管道、缓存策略,以及监控、故障恢复与性能优化方案,并给出典型数据流转场景的图示。

项目结构

围绕数据流的关键目录与职责如下: - channels/bus:异步消息队列与事件类型定义,解耦 IM 通道与 Agent 核心。 - backtest/loaders:统一的数据加载协议、校验、重试与本地缓存;按市场维护 fallback 链。 - session/service:会话生命周期编排,创建尝试(attempt)、调度 AgentLoop、持久化结果并通过事件总线广播。 - memory/persistent:跨会话持久化记忆,支持索引、去重、语义链接与可选 FTS5 搜索。 - api/swarm_routes:SSE 事件流路由,将运行事件以事件流形式推送到前端。 - frontend/stores/agent.ts:前端状态管理,包含 SSE 增量更新、会话缓存与流式文本累积。

graph TB subgraph "采集层" L1["数据加载器<br/>backtest/loaders"] R["注册表与回退链<br/>registry.py"] end subgraph "处理层" SB["会话服务<br/>session/service.py"] AL["AgentLoop(由服务调度)"] MB["消息总线<br/>channels/bus"] end subgraph "存储层" M["持久化记忆<br/>memory/persistent.py"] ST["会话与尝试持久化"] end subgraph "传输层" SSE["SSE 事件流<br/>api/swarm_routes.py"] FE["前端状态<br/>frontend/stores/agent.ts"] end L1 --> R R --> SB SB --> AL AL --> M AL --> ST AL --> MB MB --> SSE SSE --> FE

图表来源 - agent/backtest/loaders/registry.py:136-155 - agent/src/session/service.py:158-218 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165

章节来源 - agent/backtest/loaders/registry.py:136-155 - agent/src/session/service.py:158-218 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165

核心组件

章节来源 - agent/src/channels/bus/events.py:20-55 - agent/src/channels/bus/queue.py:8-44 - agent/backtest/loaders/base.py:618-645 - agent/backtest/loaders/registry.py:158-193 - agent/src/session/service.py:53-91 - agent/src/memory/persistent.py:196-220 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165

架构总览

下图展示从数据源到用户界面的端到端数据流:数据通过加载器获取并校验,进入会话服务编排的 Agent 执行;Agent 产出结果并写入持久化存储;运行时事件经消息总线或 SSE 推送至前端,前端进行增量渲染与状态管理。

sequenceDiagram participant DS as "数据源" participant L as "加载器(Registry)" participant SS as "会话服务" participant AL as "AgentLoop" participant MS as "消息总线" participant API as "SSE路由" participant FE as "前端Store" DS->>L : "fetch(codes, start, end, interval)" L-->>SS : "DataFrame(symbols -> OHLCV)" SS->>AL : "调度执行(历史+上下文)" AL-->>MS : "发布出站事件/结果" MS-->>API : "消费出站消息" API-->>FE : "SSE事件流(id/data/event)" FE-->>FE : "appendDelta/缓存/状态更新"

图表来源 - agent/backtest/loaders/registry.py:158-193 - agent/src/session/service.py:248-344 - agent/src/channels/bus/queue.py:19-33 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165

详细组件分析

消息总线与事件模型

classDiagram class InboundMessage { +string channel +string sender_id +string chat_id +string content +datetime timestamp +string[] media +dict metadata +string session_key_override +session_key() string } class OutboundMessage { +string channel +string chat_id +string content +string reply_to +string[] media +dict metadata +list[]string~~ buttons } class MessageBus { +publish_inbound(msg) void +consume_inbound() InboundMessage +publish_outbound(msg) void +consume_outbound() OutboundMessage +inbound_size int +outbound_size int } InboundMessage <.. MessageBus : "入队" OutboundMessage <.. MessageBus : "出队"

图表来源 - agent/src/channels/bus/events.py:20-55 - agent/src/channels/bus/queue.py:8-44

章节来源 - agent/src/channels/bus/events.py:20-55 - agent/src/channels/bus/queue.py:8-44 - agent/src/channels/runtime.py:34-70

数据加载器协议与回退链

flowchart TD Start(["请求加载"]) --> CheckCache{"启用本地缓存?"} CheckCache --> |是| GetCache["读取Parquet缓存"] CheckCache --> |否| FetchChain["按市场回退链尝试"] GetCache --> Hit{"命中?"} Hit --> |是| ReturnCache["返回DataFrame"] Hit --> |否| FetchChain FetchChain --> TryNext{"下一个源可用?"} TryNext --> |是| LoadData["调用loader.fetch()"] TryNext --> |否| Error["NoAvailableSourceError"] LoadData --> Validate["validate_ohlc()"] Validate --> CachePut{"范围可缓存?"} CachePut --> |是| PutCache["写入Parquet+元数据"] CachePut --> |否| ReturnData["返回DataFrame"] PutCache --> ReturnData

图表来源 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:243-439 - agent/backtest/loaders/registry.py:136-193

章节来源 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:243-439 - agent/backtest/loaders/registry.py:136-193

会话服务与执行管线

sequenceDiagram participant Client as "客户端" participant Service as "会话服务" participant Store as "持久化存储" participant Loop as "AgentLoop" participant Bus as "事件总线" Client->>Service : "send_message(session_id, content)" Service->>Store : "append_message / create_attempt" Service->>Service : "_reserve_session()" Service->>Loop : "run(user_message, history)" Loop-->>Service : "result(status, metrics, run_dir)" Service->>Store : "update_attempt / append_reply" Service->>Bus : "emit(attempt.completed/failed/cancelled)"

图表来源 - agent/src/session/service.py:158-218 - agent/src/session/service.py:248-344

章节来源 - agent/src/session/service.py:158-218 - agent/src/session/service.py:248-344

持久化记忆与索引

flowchart TD Add["写入记忆"] --> Sanitize["清理控制字符/截断"] Sanitize --> WriteFile["写入.md + 更新索引"] WriteFile --> Links{"启用语义链接?"} Links --> |是| Discover["发现关联条目"] Links --> |否| Done Discover --> Done["完成"]

图表来源 - agent/src/memory/persistent.py:462-578 - agent/src/memory/persistent.py:358-438

章节来源 - agent/src/memory/persistent.py:462-578 - agent/src/memory/persistent.py:358-438

SSE 事件流与前端状态

sequenceDiagram participant API as "SSE路由" participant Store as "前端Store" API->>Store : "event : message/tool_call/tool_result" Store->>Store : "appendDelta / setSseStatus" API-->>Store : "event : done"

图表来源 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165 - frontend/src/stores/agent.ts:282-335

章节来源 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165 - frontend/src/stores/agent.ts:282-335

依赖关系分析

graph LR Base["base.py"] --> Registry["registry.py"] Registry --> Service["session/service.py"] Service --> Bus["channels/bus/queue.py"] Bus --> Routes["api/swarm_routes.py"] Routes --> Frontend["frontend/stores/agent.ts"]

图表来源 - agent/backtest/loaders/base.py:618-645 - agent/backtest/loaders/registry.py:158-193 - agent/src/session/service.py:158-218 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165

章节来源 - agent/backtest/loaders/base.py:618-645 - agent/backtest/loaders/registry.py:158-193 - agent/src/session/service.py:158-218 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:136-165

性能考虑

[本节为通用性能建议,不直接分析具体文件]

故障恢复与监控

章节来源 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:163-236 - agent/src/session/service.py:43-50 - agent/src/session/service.py:227-246 - agent/src/channels/bus/queue.py:35-44 - agent/src/api/swarm_routes.py:187-211 - frontend/src/stores/agent.ts:304-305

结论

Vibe-Trading 的数据流以事件驱动为核心,通过消息总线解耦渠道与 Agent,以可插拔的加载器注册表实现多市场、多源的数据接入,配合会话服务编排执行与持久化,最终通过 SSE 将实时事件推送至前端。系统在数据验证、缓存、重试、取消与重连等方面具备完善的健壮性保障,并提供丰富的监控点与性能优化手段。该设计既支持离线回测与研究,也支撑实时交互与可视化,满足复杂量化工作流的端到端需求。

附录:典型数据流转场景

[本节为概念性场景描述,不直接分析具体文件]