数据处理管道

📎 引用文件

本文引用的文件 - base.py - registry.py - market_data.py - yfinance_loader.py - accessor.py - schema.py

目录

  1. 简介
  2. 项目结构
  3. 核心组件
  4. 架构总览
  5. 详细组件分析
  6. 依赖关系分析
  7. 性能考量
  8. 故障排除指南
  9. 结论
  10. 附录

简介

本文件为 Vibe-Trading 数据处理管道的全面技术文档,覆盖数据获取、转换、标准化与缓存的完整流程;解释数据源注册机制、适配器模式与批量处理策略;记录配置选项、错误重试机制与性能监控;并提供数据流转图、组件交互图、扩展开发指南、自定义数据源集成方法、性能调优建议以及调试与排障指南。

项目结构

数据处理管道位于 agent/backtest/loaders(数据加载器)与 agent/src/market_data.py(统一市场数据入口),并通过 src/config 提供配置访问能力。关键组织方式: - 数据加载器协议与通用工具:base.py - 数据源注册与回退链:registry.py - 统一市场数据接口与自动路由:market_data.py - 具体数据源实现示例:yfinance_loader.py - 配置访问与布尔解析:accessor.py、schema.py

graph TB A["调用方<br/>MCP/CLI/工具"] --> B["统一入口<br/>fetch_market_data()"] B --> C["数据源选择<br/>detect_source()"] C --> D["注册表与回退链<br/>resolve_loader()/FALLBACK_CHAINS"] D --> E["具体加载器实例<br/>loader.fetch()"] E --> F["标准化与校验<br/>validate_ohlc()"] F --> G["可选本地缓存<br/>loader_cache_*"] G --> H["结果聚合与截断<br/>cap_rows()"] H --> I["返回结构化结果<br/>含来源溯源(可选)"]

图表来源 - market_data.py:97-223 - registry.py:136-193 - base.py:343-439

章节来源 - market_data.py:1-229 - registry.py:1-249 - base.py:1-645

核心组件

章节来源 - base.py:27-119 - base.py:127-236 - base.py:243-439 - registry.py:23-155 - registry.py:158-249 - market_data.py:16-64 - market_data.py:97-223 - yfinance_loader.py:22-95 - yfinance_loader.py:172-200

架构总览

下图展示从调用方到数据源的端到端流程,包括自动路由、回退链、标准化、缓存与结果聚合。

sequenceDiagram participant U as "调用方" participant M as "market_data.fetch_market_data" participant R as "registry.resolve_loader" participant L as "具体加载器" participant B as "base 校验/缓存" U->>M : 传入 codes/start/end/source/interval M->>M : detect_source + 分组 M->>R : get_loader_cls_with_fallback(或按链尝试) R-->>M : 返回可用加载器类 M->>L : loader.fetch(codes, start, end, interval) L->>B : validate_date_range / validate_ohlc L->>B : cached_loader_fetch (可选) B-->>L : DataFrame 或 None L-->>M : {symbol : DataFrame} M->>M : cap_rows + JSON安全化 M-->>U : 结果字典(可含_provenance/_unresolved)

图表来源 - market_data.py:97-223 - registry.py:158-249 - base.py:343-439

详细组件分析

数据源注册与回退链(registry.py)

flowchart TD Start(["开始"]) --> CheckReg["_ensure_registered()"] CheckReg --> GetChain["获取市场回退链"] GetChain --> ForEach{"遍历候选源"} ForEach --> |存在| TryInst["构造实例并 is_available()"] TryInst --> |可用| ReturnLoader["返回加载器"] TryInst --> |不可用| Next["下一个候选"] Next --> ForEach ForEach --> |无可用| RaiseErr["抛出 NoAvailableSourceError"] ReturnLoader --> End(["结束"]) RaiseErr --> End

图表来源 - registry.py:71-115 - registry.py:136-193

章节来源 - registry.py:1-249

统一市场数据入口(market_data.py)

sequenceDiagram participant C as "调用方" participant F as "fetch_market_data" participant D as "detect_source" participant R as "get_loader_cls_with_fallback" participant L as "loader.fetch" C->>F : codes,start,end,source,interval F->>D : 对每个 code 推断首选源 F->>F : 按(src, market)分组 loop 每组 F->>R : 获取加载器类(可回退) R-->>F : 加载器类 F->>L : fetch(codes,...) alt 成功且有数据 L-->>F : {symbol : DataFrame} F->>F : cap_rows + JSON安全化 else 失败 F->>F : 记录日志并尝试下一个候选 end end F-->>C : 结果(可能含_unresolved/_provenance)

图表来源 - market_data.py:16-64 - market_data.py:97-223 - registry.py:158-249

章节来源 - market_data.py:1-229

数据标准化与校验(base.py)

flowchart TD S(["输入DataFrame"]) --> CheckCols{"包含OHLC列?"} CheckCols --> |否| RetOrig["原样返回"] CheckCols --> |是| Structural{"结构性违规?"} Structural --> |是| Strategy{"策略? drop/warn/raise"} Strategy --> |drop| DropRows["删除违规行"] Strategy --> |warn| KeepWarn["保留并记录警告"] Strategy --> |raise| RaiseErr["抛出异常"] Structural --> |否| PosCheck{"允许非正价格?"} PosCheck --> |否| NonPos{"是否存在<=0价格?"} PosCheck --> |是| ZeroOnly{"是否存在=0价格?"} NonPos --> |是| ApplyStrategy["同上"] NonPos --> |否| Pass["通过"] ZeroOnly --> |是| ApplyStrategy ZeroOnly --> |否| Pass DropRows --> End(["输出"]) KeepWarn --> End RaiseErr --> End Pass --> End

图表来源 - base.py:31-119

章节来源 - base.py:31-119

本地缓存机制(base.py)

flowchart TD Start(["进入 cached_loader_fetch"]) --> CheckEnabled{"缓存启用?"} CheckEnabled --> |否| DoFetch["直接 fetch()"] CheckEnabled --> |是| RangeOK{"end_date 已结算?"} RangeOK --> |否| DoFetch RangeOK --> |是| ReadCache["loader_cache_get(...)"] ReadCache --> Hit{"命中?"} Hit --> |是| ReturnCache["返回缓存DataFrame"] Hit --> |否| DoFetch DoFetch --> PutCache["loader_cache_put(..., frame)"] PutCache --> ReturnFetch["返回新获取的DataFrame"] ReturnCache --> End(["结束"]) ReturnFetch --> End

图表来源 - base.py:243-439 - base.py:475-595

章节来源 - base.py:243-595

重试与预算(base.py)

flowchart TD S(["调用 retry_with_budget(fn)"]) --> Attempt1["第1次尝试"] Attempt1 --> Ok{"成功?"} Ok --> |是| Return["返回结果"] Ok --> |否| Transient{"是否为瞬时常异?"} Transient --> |否| Propagate["直接传播异常"] Transient --> |是| Budget{"剩余预算>0 且未达上限?"} Budget --> |否| Timeout["抛出TimeoutError(包装原异常)"] Budget --> |是| Sleep["sleep(min(backoff,剩余))"] --> Next["下一次尝试"] Next --> AttemptN["第N次尝试"] AttemptN --> Ok

图表来源 - base.py:127-236

章节来源 - base.py:127-236

示例加载器:yfinance(yfinance_loader.py)

章节来源 - yfinance_loader.py:22-95 - yfinance_loader.py:172-200

依赖关系分析

graph LR MD["market_data.py"] --> REG["registry.py"] MD --> BASE["base.py"] REG --> BASE YF["yfinance_loader.py"] --> BASE YF --> REG BASE --> ACC["accessor.py"]

图表来源 - market_data.py:97-223 - registry.py:158-249 - base.py:243-439 - yfinance_loader.py:14-21 - accessor.py:52-76

章节来源 - market_data.py:1-229 - registry.py:1-249 - base.py:1-645 - yfinance_loader.py:1-200 - accessor.py:1-149

性能考量

[本节为通用指导,无需特定文件引用]

故障排除指南

章节来源 - registry.py:158-249 - base.py:50-119 - base.py:243-439 - base.py:127-236 - market_data.py:66-84 - market_data.py:198-223

结论

Vibe-Trading 的数据处理管道通过统一的入口、灵活的回退链、严格的标准化与健壮的重试/缓存机制,实现了跨市场、多数据源的高可用数据供给。借助内容寻址缓存与批量分组策略,系统在稳定性与性能之间取得良好平衡;通过溯源信息与截断策略,便于运维观测与容量规划。新增数据源只需遵循 DataLoaderProtocol 并自注册,即可无缝融入现有生态。

[本节为总结性内容,无需特定文件引用]

附录

管道配置选项与环境变量

章节来源 - base.py:243-281 - accessor.py:52-76 - market_data.py:97-117

扩展开发指南:自定义数据源集成

章节来源 - registry.py:62-68 - registry.py:136-155 - base.py:618-645

性能调优方法

[本节为通用指导,无需特定文件引用]

调试工具与技巧

章节来源 - market_data.py:198-223 - base.py:475-595