统一数据管理层

📎 引用文件

本文引用的文件 - registry.py - base.py - market_data.py - yfinance_loader.py - local_loader.py - ccxt_loader.py - tushare.py - eastmoney_loader.py - _symbol_utils.py - cn_adjust.py

目录

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

简介

本文件为 Vibe-Trading 的统一数据管理层提供系统化文档,聚焦于“多源数据接入 + 一致访问接口”的架构设计。内容涵盖: - 数据加载器注册机制与按市场类型的回退链 - 数据格式标准化(OHLCV)与校验规则 - 本地缓存策略(启用条件、键生成、读写流程) - 错误处理与重试预算 - 各数据源特性对比、配置方法与性能调优建议 - 数据验证、缺失值处理与异常检测 - 自定义数据源开发指南(接口规范、测试要求、部署步骤) - 查询、过滤与转换的实际使用示例

项目结构

统一数据管理层位于 agent/backtest/loaders 与 agent/src/market_data.py 中,采用“协议 + 注册表 + 回退链”的分层设计: - 协议层:定义 DataLoaderProtocol,统一 fetch/is_available 等接口 - 注册层:LOADER_REGISTRY 与 @register 装饰器完成加载器自动发现 - 路由层:FALLBACK_CHAINS 按市场类型组织回退顺序;detect_source 将符号映射到首选源 - 工具层:日期/价格校验、重试预算、本地缓存、JSON 安全化等通用能力 - 入口层:fetch_market_data 聚合分组、回退、结果裁剪与溯源信息

graph TB A["调用方<br/>MCP/工具"] --> B["market_data.fetch_market_data"] B --> C["registry.resolve_loader / get_loader_cls_with_fallback"] C --> D["具体 Loader<br/>yfinance/local/ccxt/tushare/eastmoney..."] D --> E["基础能力<br/>校验/重试/缓存"] D --> F["外部数据源<br/>交易所/公开API/本地文件"]

图表来源 - market_data.py:97-223 - registry.py:158-248 - base.py:184-236

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

核心组件

章节来源 - base.py:27-119 - base.py:184-236 - base.py:243-439 - registry.py:23-155 - registry.py:158-248 - market_data.py:16-63 - market_data.py:97-223

架构总览

下图展示从高层 API 到底层数据源的完整调用路径,包括回退链与缓存命中分支。

sequenceDiagram participant U as "调用方" participant M as "market_data.fetch_market_data" participant R as "registry.get_loader_cls_with_fallback" participant L as "具体Loader实例" participant C as "缓存(可选)" participant S as "外部数据源" U->>M : 传入 codes/start/end/source/interval M->>M : detect_source + 按市场分组 loop 每个(源,市场)组 M->>R : 获取可回退的Loader类 R-->>M : Loader类或抛出NoAvailableSourceError M->>L : 构造并调用fetch(codes,...) alt 命中缓存 L->>C : 读取Parquet C-->>L : DataFrame else 未命中缓存 L->>S : 拉取原始数据 S-->>L : 原始DataFrame L->>C : 写入Parquet(可选) end L-->>M : {symbol : DataFrame} M->>M : 标准化/校验/裁剪/JSON安全化 end M-->>U : 返回结果 + 可选_provenance

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

详细组件分析

数据加载器注册与回退机制

flowchart TD Start(["开始"]) --> Ensure["_ensure_registered()<br/>导入所有loader模块"] Ensure --> Chain{"按市场获取回退链"} Chain --> TryNext{"遍历候选源"} TryNext --> |存在| Construct["构造Loader实例"] Construct --> Avail{"is_available() ?"} Avail --> |是| Return["返回可用Loader"] Avail --> |否| Next["下一个候选"] Next --> TryNext TryNext --> |全部失败| Raise["抛出NoAvailableSourceError"]

图表来源 - registry.py:71-115 - registry.py:158-193 - registry.py:196-248

章节来源 - registry.py:1-249 - market_data.py:16-63

数据格式标准化与校验

flowchart TD In["原始DataFrame"] --> CheckCols{"是否包含OHLC列?"} CheckCols --> |否| OutEmpty["返回空/原样"] CheckCols --> |是| Structural["检查结构性不变式"] Structural --> Pos{"是否允许非正价格?"} Pos --> |否| RejectZero["拒绝<=0的价格"] Pos --> |是| AllowNeg["仅拒绝==0的价格"] RejectZero --> Invalid{"是否存在无效行?"} AllowNeg --> Invalid Invalid --> |是| Strategy{"策略: drop/warn/raise"} Strategy --> Drop["删除无效行"] Strategy --> Warn["记录警告并保留"] Strategy --> Raise["抛出异常"] Invalid --> |否| Pass["通过"]

图表来源 - base.py:31-119 - local_loader.py:83-127

章节来源 - base.py:31-119 - local_loader.py:83-127

缓存策略(本地 Parquet)

sequenceDiagram participant L as "Loader" participant K as "make_loader_cache_key" participant P as "loader_cache_path" participant R as "loader_cache_get" participant W as "loader_cache_put" L->>K : 计算内容哈希键 L->>P : 生成parquet路径 L->>R : 尝试读取缓存 alt 命中 R-->>L : 返回DataFrame else 未命中 L->>L : 执行fetch() L->>W : 写入缓存(非空且可缓存) W-->>L : 完成 L-->>L : 返回DataFrame end

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

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

错误处理与重试预算

章节来源 - base.py:184-236 - tushare.py:18-79 - ccxt_loader.py:50-57 - market_data.py:170-190

各数据源特性对比与配置要点

章节来源 - yfinance_loader.py:1-200 - local_loader.py:1-200 - ccxt_loader.py:1-200 - tushare.py:1-200 - eastmoney_loader.py:1-177 - cn_adjust.py:1-78

数据查询、过滤与转换示例

章节来源 - market_data.py:66-84 - market_data.py:97-223 - yfinance_loader.py:172-200 - eastmoney_loader.py:145-177 - local_loader.py:83-127

自定义数据源开发指南

章节来源 - base.py:31-119 - base.py:401-439 - registry.py:27-59 - registry.py:83-108 - registry.py:136-155

依赖关系分析

graph LR MD["market_data.py"] --> REG["registry.py"] MD --> BASE["base.py"] YF["yfinance_loader.py"] --> BASE LL["local_loader.py"] --> BASE CC["ccxt_loader.py"] --> BASE TS["tushare.py"] --> BASE TS --> ADJ["cn_adjust.py"] EM["eastmoney_loader.py"] --> BASE TS --> SYM["_symbol_utils.py"]

图表来源 - market_data.py:97-223 - registry.py:158-248 - base.py:184-439 - tushare.py:1-200 - cn_adjust.py:1-78 - _symbol_utils.py:1-21

章节来源 - market_data.py:97-223 - registry.py:158-248 - base.py:184-439 - tushare.py:1-200 - cn_adjust.py:1-78 - _symbol_utils.py:1-21

性能考量

故障排查指南

章节来源 - registry.py:158-193 - base.py:27-29 - base.py:50-119 - base.py:343-439 - tushare.py:18-79 - ccxt_loader.py:50-57

结论

统一数据管理层通过“协议 + 注册表 + 回退链 + 通用工具”的组合,实现了跨市场、跨供应商的一致数据访问。其关键优势在于: - 稳定的接口契约与严格的 OHLC 校验 - 健壮的回退链与受限重试预算,提升可用性 - 可选本地缓存显著降低重复请求成本 - 灵活的符号路由与本地文件重采样能力 - 可扩展的自定义数据源开发与部署流程

附录